Compare commits

..

6 Commits

Author SHA1 Message Date
fdc310 e211d3aae6 fix(itchat): correct group/private message handling and sender resolution
- filter the bot's own messages to prevent reply loops
- use startswith('@@') for group detection
- use ActualUserName as stable GroupMember id
- fill Friend.remark from RemarkName
- strip @mention prefix with a regex
2026-08-13 11:18:09 +08:00
fdc310 42f1f00772 Merge remote-tracking branch 'origin/master' into feature/itchat-adapter
# Conflicts:
#	src/langbot/pkg/api/http/controller/groups/platform/adapters.py
#	src/langbot/pkg/plugin/connector.py
#	uv.lock
#	web/src/app/home/bots/BotDetailContent.tsx
#	web/src/app/home/components/dynamic-form/DynamicFormItemComponent.tsx
#	web/src/app/home/components/qrcode-login/QrCodeLoginDialog.tsx
2026-08-12 16:58:58 +08:00
fdc310 1ed107c9d5 feat(dynamic-form): enhance select handling with empty option support and UUID filtering 2026-07-06 13:36:58 +08:00
fdc310 0cce418956 feat(itchat): improve login session handling and error reporting 2026-07-02 14:40:13 +08:00
fdc310 d4e8ccd161 feat: enhance itchat adapter with runtime status tracking and UI integration 2026-07-02 13:56:24 +08:00
fdc310 78fb40a28a feat: add itchat-uos WeChat adapter with QR code login
- Add itchat-uos adapter supporting personal WeChat via QR code login
- Implement message/event converters for text, image, voice, sharing types
- Bridge sync itchat callbacks to async LangBot pipeline via asyncio
- Add QR login API endpoints with session management
- Add frontend QR code login dialog integration
- Fix plugin connector handler attribute check
- Use fresh Core instance per login to avoid singleton state pollution
2026-07-01 18:03:26 +08:00
324 changed files with 4187 additions and 23437 deletions
+2 -2
View File
@@ -1,5 +1,5 @@
name: 漏洞反馈
description: 【供中文用户】报错或漏洞请使用这个模板创建,不使用此模板创建的异常、漏洞相关issue将被直接关闭。由于自己操作不当/不甚了解所用技术栈引起的网络连接问题恕无法解决,请勿提 issue。容器间网络连接问题,参考文档 https://langbot.app/docs/zh/workshop/network-details
description: 【供中文用户】报错或漏洞请使用这个模板创建,不使用此模板创建的异常、漏洞相关issue将被直接关闭。由于自己操作不当/不甚了解所用技术栈引起的网络连接问题恕无法解决,请勿提 issue。容器间网络连接问题,参考文档 https://link.langbot.app/zh/docs/network
title: "[Bug]: "
labels: ["bug?"]
body:
@@ -22,7 +22,7 @@ body:
- type: textarea
attributes:
label: 异常情况
description: 完整描述异常情况,什么时候发生的、发生了什么。**请附带日志信息。**
description: 完整描述异常情况,什么时候发生的、发生了什么。**请附带日志信息。**
validations:
required: true
- type: textarea
+1 -1
View File
@@ -1,5 +1,5 @@
name: Bug report
description: Report bugs or vulnerabilities using this template. For container network connection issues, refer to the documentation https://langbot.app/docs/en/workshop/network-details
description: Report bugs or vulnerabilities using this template. For container network connection issues, refer to the documentation https://link.langbot.app/en/docs/network
title: "[Bug]: "
labels: ["bug?"]
body:
-111
View File
@@ -1,111 +0,0 @@
# Discord release announcements
This independent workflow announces new stable LangBot releases in the channel
selected by a dedicated Discord incoming webhook. It does not change the existing
release/build workflows, edit releases, run a persistent service, poll, or backfill.
Announcements run on publication, independently of artifact builds finishing.
## Setup and read-only validation
1. In the intended community **announcement channel**, create a dedicated incoming
webhook (Channel Settings → Integrations → Webhooks). Copy its URL; do not reuse
a webhook belonging to another automation.
2. In `langbot-app/LangBot` → Settings → Secrets and variables → Actions, create the
**repository secret** `DISCORD_RELEASE_WEBHOOK_URL`. Its value must be exactly
`https://discord.com/api/webhooks/<id>/<token>` — no query, trailing slash,
API-version segment, or alternate domain. Treat the entire URL as a password.
3. Once this workflow is on `master`, open Actions → **Discord Release Announcement**
→ Run workflow, choosing `master`. Alternatively:
```sh
gh workflow run discord-release.yml --repo langbot-app/LangBot --ref master
```
4. Inspect **Validate webhook (GET only, no message)**. It checks webhook type `1`
and reports `guild_id` and `channel_id`; compare both with the intended server
and channel using Discord Developer Mode → Copy ID. The secret determines the
destination; no channel ID is guessed or overridden. The URL/token is never
logged. Dispatch cannot send a test message or announce an old release, even
when run again. Missing/invalid secrets fail validation clearly; offline tests
do not need secrets.
GET validation confirms the webhook's identity, not delivery or notification
permissions. Verify those on the first genuine release. `mention_everyone=true`
confirms Discord parsed the mention; it cannot prove every member received a push
notification (member/server notification settings still apply).
## Activation and message
The workflow and `.github/discord-release/` helper **must be in the commit targeted
by each new release tag**. Merging to `master` does not enable announcements for
old tags whose commits lack these files. Manual dispatch becomes available when
the workflow is on the default branch. Only publish release tags from trusted,
reviewed commits: release workflows execute that tag's code with the secret.
Only `release` events with action `published`, `draft=false`, and
`prerelease=false` can send. Drafts and prereleases are skipped; release edits do
not trigger announcements. The helper requires the repository to be exactly
`langbot-app/LangBot`, a stable `vX.Y.Z` tag (ASCII digits, at most 64 characters),
and its exact canonical GitHub release URL. Other naming schemes fail closed.
Example message (the version and URL come from the validated event file):
```text
@everyone LangBot v4.10.11 is now available!
Release notes: https://github.com/langbot-app/LangBot/releases/tag/v4.10.11
```
The release title/body is never copied. There is one literal `@everyone`, explicit
`allowed_mentions.parse=["everyone"]`, empty user/role allowlists, and no reply
mention. TTS and notification-suppressing flags are disabled. Requests use HTTPS
only to `discord.com`, an explicit User-Agent, and no redirects or automatic
retries. After a webhook identity GET, one `POST ?wait=true` obtains a message ID;
an exact `/messages/<id>` GET verifies its ID, webhook/channel, content,
`mention_everyone=true`, and empty user/role mention arrays before success.
## Repeat guard and manual recovery
Production sending requires **`GITHUB_RUN_ATTEMPT == "1"`**. Any Actions rerun
(including “Re-run failed jobs”) refuses to POST and requires manual reconciliation,
even if the first attempt failed before sending. Read-only dispatch may be rerun.
This is a practical repeat guard, **not durable exactly-once delivery**. It cannot
prevent duplicates from a separate new run/event (for example deleting/recreating
a release), separate automation, or manual posting. It stores no durable dedupe
state and never modifies the release to mark delivery.
If a POST times out, returns an error, or readback fails, the message may already
exist. The workflow fails rather than blindly sending again. A returned message ID
is included in the safe error when available. A runner termination can also leave
an ambiguous send without that log line.
1. Inspect the announcement channel and the failed run logs. Locate the canonical
release link and, if available, the returned message ID. A failed verification
does **not** mean the message was absent.
2. If present, reconcile the existing message/mention problem manually; do not
rerun, create another release event, or send a duplicate ping.
3. If an operator has positively confirmed no message exists, fix the secret or
permission issue and use read-only dispatch to validate configuration. A
maintainer may then post the announcement manually once in Discord and record
the message link in the incident/run notes. Do not override the attempt guard
or delete/recreate a release to force recovery.
4. If absence cannot be established, pause and reconcile rather than resending.
To stop future sends, disable **Discord Release Announcement** in Actions. Rotate
or delete the dedicated Discord webhook if the URL is exposed, and update the
secret before validation. No rollback of release artifacts is involved.
## Local checks
Requires Python 3.11+ and the standard library only:
```sh
python3 -m unittest discover -s .github/discord-release -p 'test_*.py' -v
python3 -m py_compile .github/discord-release/announce.py .github/discord-release/test_announce.py
```
Tests exercise policy, CLI/event-file handling, mention payloads, hostile inputs,
HTTP failures, exact message readback, and refusal to retry. Only the HTTPS
transport is mocked for Discord tests; no live Discord requests or messages are
made. Changes to this directory or its workflow run the offline tests on push and
pull request; tests also gate release sending and read-only dispatch validation.
-162
View File
@@ -1,162 +0,0 @@
"""Announce only first-attempt stable releases; dispatch is read-only validation."""
import http.client
import json
import os
from pathlib import Path
import re
import sys
REPOSITORY = 'langbot-app/LangBot'
RELEASE_PREFIX = f'https://github.com/{REPOSITORY}/releases/tag/'
RECONCILE = (
'Do not resend or bypass the run-attempt guard; manual reconciliation is required. '
'Inspect the announcement channel and workflow logs before any manual recovery '
'(see .github/discord-release/README.md).'
)
class AnnouncementError(Exception):
"""A safe, operator-facing error containing no webhook URL or response body."""
def release_payload(event, attempt):
"""Return a bounded, mention-safe payload, or None for draft/preview releases."""
if not isinstance(event, dict) or event.get('action') != 'published':
raise AnnouncementError('Only release.published events are accepted.')
repository = event.get('repository')
if not isinstance(repository, dict) or repository.get('full_name') != REPOSITORY:
raise AnnouncementError('Unexpected release repository.')
release = event.get('release')
if not isinstance(release, dict) or any(type(release.get(key)) is not bool for key in ('draft', 'prerelease')):
raise AnnouncementError('Invalid release flags.')
if release['draft'] or release['prerelease']:
return None
if attempt != '1':
raise AnnouncementError(f'Release reruns or missing run attempts are refused. {RECONCILE}')
tag = release.get('tag_name')
if not isinstance(tag, str) or len(tag) > 64 or not re.fullmatch(r'v[0-9]+\.[0-9]+\.[0-9]+', tag):
raise AnnouncementError('Expected a stable release tag in vX.Y.Z format (at most 64 characters).')
url = RELEASE_PREFIX + tag
if release.get('html_url') != url:
raise AnnouncementError('Release URL must be the canonical LangBot release URL matching its tag.')
return {
'content': f'@everyone LangBot {tag} is now available!\nRelease notes: {url}',
'allowed_mentions': {'parse': ['everyone'], 'users': [], 'roles': [], 'replied_user': False},
'tts': False,
'flags': 0,
}
def is_snowflake(value):
return isinstance(value, str) and re.fullmatch(r'[0-9]{1,20}', value) is not None
class DiscordWebhook:
def __init__(self, url):
if not url:
raise AnnouncementError('DISCORD_RELEASE_WEBHOOK_URL is missing. Set the repository Actions secret.')
match = re.fullmatch(r'https://discord\.com(/api/webhooks/([0-9]{1,20})/[A-Za-z0-9_-]+)', url)
if not match:
raise AnnouncementError('Invalid webhook URL; expected https://discord.com/api/webhooks/<id>/<token>.')
self.path, self.id = match.groups()
def _request(self, method, suffix='', payload=None):
# Direct HTTPS, default certificate verification, no proxies or redirect/retry machinery.
connection = http.client.HTTPSConnection('discord.com', timeout=20)
try:
body = json.dumps(payload).encode('utf-8') if payload is not None else None
connection.request(
method,
self.path + suffix,
body=body,
headers={'Content-Type': 'application/json', 'User-Agent': 'LangBot-Release-Announcements/1.0'},
)
response = connection.getresponse()
if response.status != 200:
raise AnnouncementError(f'Discord {method} returned HTTP {response.status}; no retry was attempted.')
raw = response.read(1_048_577)
if len(raw) > 1_048_576:
raise AnnouncementError('Discord response exceeded the size limit.')
return json.loads(raw)
except (OSError, http.client.HTTPException, ValueError, UnicodeError):
# Exceptions and bodies can contain the token; never print them or chain them.
raise AnnouncementError(
f'Discord {method} failed or returned invalid JSON; no retry was attempted.'
) from None
finally:
connection.close()
def validate(self):
"""GET only: verify an incoming webhook and return safe identifying fields."""
webhook = self._request('GET')
if (
not isinstance(webhook, dict)
or type(webhook.get('type')) is not int
or webhook['type'] != 1
or webhook.get('id') != self.id
or not is_snowflake(webhook.get('guild_id'))
or not is_snowflake(webhook.get('channel_id'))
):
raise AnnouncementError('Expected an incoming (type 1) webhook with matching ID and guild/channel IDs.')
return {key: webhook[key] for key in ('id', 'type', 'guild_id', 'channel_id')}
def send(self, payload):
"""One POST, followed by exact message GET; never automatically retry a send."""
webhook = self.validate()
message_id = None
try:
sent = self._request('POST', '?wait=true', payload)
if not isinstance(sent, dict) or not is_snowflake(sent.get('id')):
raise AnnouncementError('Discord did not return a valid message ID.')
message_id = sent['id']
saved = self._request('GET', f'/messages/{message_id}')
if (
not isinstance(saved, dict)
or saved.get('id') != message_id
or saved.get('webhook_id') != self.id
or saved.get('channel_id') != webhook['channel_id']
or saved.get('content') != payload['content']
or saved.get('mention_everyone') is not True
or saved.get('mentions') != []
or saved.get('mention_roles') != []
):
raise AnnouncementError('Discord message readback did not match content, identity, or mentions.')
except AnnouncementError as error:
reference = f' Returned message ID: {message_id}.' if message_id else ''
raise AnnouncementError(f'Delivery not confirmed. {error}{reference} {RECONCILE}') from None
return message_id
def main(env=None):
env = os.environ if env is None else env
try:
if env.get('GITHUB_REPOSITORY') != REPOSITORY:
raise AnnouncementError('This workflow is restricted to langbot-app/LangBot.')
name = env.get('GITHUB_EVENT_NAME')
if name == 'workflow_dispatch':
webhook = DiscordWebhook(env.get('DISCORD_RELEASE_WEBHOOK_URL')).validate()
print(
f'Validated incoming webhook: guild_id={webhook["guild_id"]} channel_id={webhook["channel_id"]}. No message sent.'
)
return 0
if name != 'release':
raise AnnouncementError('Only release and workflow_dispatch events are accepted by this helper.')
try:
event = json.loads(Path(env.get('GITHUB_EVENT_PATH', '')).read_text(encoding='utf-8'))
except (OSError, ValueError, UnicodeError):
raise AnnouncementError('Cannot read a valid JSON release event from GITHUB_EVENT_PATH.') from None
payload = release_payload(event, env.get('GITHUB_RUN_ATTEMPT'))
if payload is None:
print('Skipped draft or prerelease; no message sent.')
return 0
message_id = DiscordWebhook(env.get('DISCORD_RELEASE_WEBHOOK_URL')).send(payload)
print(f'Announcement verified by exact message readback: message_id={message_id}.')
return 0
except AnnouncementError as error:
print(f'Error: {error}', file=sys.stderr)
return 1
if __name__ == '__main__':
sys.exit(main())
-427
View File
@@ -1,427 +0,0 @@
"""Offline contract tests; no Discord credentials or network required."""
import contextlib
import io
import json
import os
from pathlib import Path
import subprocess
import sys
import tempfile
import unittest
from unittest.mock import MagicMock, patch
try:
import announce
except ModuleNotFoundError:
announce = None
WEBHOOK = 'https://discord.com/api/webhooks/123456789012345678/fixture_token-ONLY'
WEBHOOK_ID = '123456789012345678'
GUILD_ID = '234567890123456789'
CHANNEL_ID = '345678901234567890'
MESSAGE_ID = '456789012345678901'
REPO = 'langbot-app/LangBot'
URL = f'https://github.com/{REPO}/releases/tag/v4.10.11'
CONTENT = f'@everyone LangBot v4.10.11 is now available!\nRelease notes: {URL}'
def event():
return {
'action': 'published',
'repository': {'full_name': REPO},
'release': {
'draft': False,
'prerelease': False,
'tag_name': 'v4.10.11',
'html_url': URL,
'name': 'Hostile @everyone <@123> $(touch /tmp/unsafe)',
'body': '@everyone @here <@123> <@&456> `hostile`',
},
}
def metadata():
return {'id': WEBHOOK_ID, 'type': 1, 'guild_id': GUILD_ID, 'channel_id': CHANNEL_ID}
def message():
return {
'id': MESSAGE_ID,
'webhook_id': WEBHOOK_ID,
'channel_id': CHANNEL_ID,
'content': CONTENT,
'mention_everyone': True,
'mentions': [],
'mention_roles': [],
}
class BaseTest(unittest.TestCase):
def setUp(self):
self.assertIsNotNone(announce, 'The release announcement helper must exist')
class PolicyTests(BaseTest):
def test_payload_has_one_literal_everyone_and_no_untrusted_body(self):
payload = announce.release_payload(event(), '1')
self.assertEqual(payload['content'], CONTENT)
self.assertEqual(json.dumps(payload).count('@everyone'), 1)
self.assertEqual(
payload['allowed_mentions'],
{
'parse': ['everyone'],
'users': [],
'roles': [],
'replied_user': False,
},
)
self.assertIs(payload['tts'], False)
self.assertEqual(payload['flags'], 0)
def test_drafts_and_prereleases_are_skipped(self):
for flag in ('draft', 'prerelease'):
with self.subTest(flag=flag):
value = event()
value['release'][flag] = True
self.assertIsNone(announce.release_payload(value, '1'))
def test_only_published_action_is_accepted(self):
for action in ('edited', 'created', 'released', 'deleted', '', None):
with self.subTest(action=action):
value = event()
value['action'] = action
with self.assertRaises(announce.AnnouncementError):
announce.release_payload(value, '1')
def test_reruns_and_missing_attempt_refuse_manual_reconciliation(self):
for attempt in ('2', '3', '', None, '01', '0', '1\n'):
with self.subTest(attempt=attempt):
with self.assertRaisesRegex(announce.AnnouncementError, 'manual reconciliation'):
announce.release_payload(event(), attempt)
def test_repository_must_match_exactly(self):
for repo in ('evil/LangBot', 'langbot-app/langbot', None):
value = event()
value['repository']['full_name'] = repo
with self.assertRaises(announce.AnnouncementError):
announce.release_payload(value, '1')
def test_hostile_and_noncanonical_tags_are_rejected(self):
for tag in (
'v1.2.3 @everyone',
'v1.2.3\n',
'v1.2.3/../../x',
'v1.2.3?x=y',
'$(id)',
'v1.2.3-rc.1',
'v.2.3',
'v1.2.3%0a',
'<@123>',
'v1.2.' + '3' * 100,
'',
None,
123,
):
with self.subTest(tag=tag):
value = event()
value['release']['tag_name'] = tag
value['release']['html_url'] = f'https://github.com/{REPO}/releases/tag/{tag}'
with self.assertRaises(announce.AnnouncementError):
announce.release_payload(value, '1')
def test_release_url_must_be_canonical_and_match_tag(self):
for url in (
'https://evil.example/tag/v4.10.11',
URL + '?x=y',
URL + '#anchor',
URL + '/',
URL.replace('v4.10.11', 'v4.10.12'),
URL.replace('github.com', 'github.com@evil.example'),
URL.replace('https:', 'http:'),
URL + '\n',
None,
):
with self.subTest(url=url):
value = event()
value['release']['html_url'] = url
with self.assertRaises(announce.AnnouncementError):
announce.release_payload(value, '1')
def test_malformed_events_fail_closed(self):
for value in (None, [], {}, {'release': []}, {'repository': None}):
with self.subTest(value=value):
with self.assertRaises(announce.AnnouncementError):
announce.release_payload(value, '1')
for flag in ('draft', 'prerelease'):
for bad in (None, 'false', 0, 1):
value = event()
value['release'][flag] = bad
with self.assertRaises(announce.AnnouncementError):
announce.release_payload(value, '1')
class DiscordTests(BaseTest):
def setUp(self):
super().setUp()
self.patch = patch('announce.http.client.HTTPSConnection')
self.connection_class = self.patch.start()
self.addCleanup(self.patch.stop)
self.connection = self.connection_class.return_value
def respond(self, *values):
responses = []
for value in values:
response = MagicMock()
response.status = 200
response.read.return_value = json.dumps(value).encode()
responses.append(response)
self.connection.getresponse.side_effect = responses
def methods(self):
return [call.args[0] for call in self.connection.request.call_args_list]
def test_webhook_validation_is_get_only_and_reports_ids(self):
self.respond(metadata())
result = announce.DiscordWebhook(WEBHOOK).validate()
self.assertEqual(result, metadata())
self.assertEqual(self.methods(), ['GET'])
self.assertEqual(
self.connection.request.call_args.args[:2], ('GET', f'/api/webhooks/{WEBHOOK_ID}/fixture_token-ONLY')
)
self.connection_class.assert_called_with('discord.com', timeout=20)
self.connection.close.assert_called_once()
def test_invalid_webhook_urls_are_rejected_before_network(self):
for url in (
'',
None,
WEBHOOK + '/',
WEBHOOK + '?wait=true',
WEBHOOK + '#x',
WEBHOOK + '\n',
' ' + WEBHOOK,
WEBHOOK.replace('https:', 'http:'),
WEBHOOK.replace('discord.com', 'discord.com.evil.example'),
WEBHOOK.replace('discord.com', 'discord.com@evil.example'),
WEBHOOK.replace('discord.com', 'discord.com:443'),
WEBHOOK.replace('/api/', '/api/v10/'),
WEBHOOK.replace(WEBHOOK_ID, 'abc'),
WEBHOOK + '/../../x',
WEBHOOK.replace('fixture_token-ONLY', 'a%2Fb'),
):
with self.subTest(url=url):
with self.assertRaises(announce.AnnouncementError):
announce.DiscordWebhook(url)
self.connection_class.assert_not_called()
def test_webhook_metadata_requires_incoming_type_and_ids(self):
invalid = [
None,
[],
{},
dict(metadata(), type=2),
dict(metadata(), type=True),
dict(metadata(), id='999'),
dict(metadata(), channel_id=None),
dict(metadata(), guild_id='::error::hostile'),
]
for value in invalid:
with self.subTest(value=value):
self.respond(value)
with self.assertRaises(announce.AnnouncementError):
announce.DiscordWebhook(WEBHOOK).validate()
self.assertNotIn('POST', self.methods())
def test_send_waits_and_reads_back_exact_returned_message(self):
self.respond(metadata(), message(), message())
result = announce.DiscordWebhook(WEBHOOK).send(announce.release_payload(event(), '1'))
self.assertEqual(result, MESSAGE_ID)
self.assertEqual(self.methods(), ['GET', 'POST', 'GET'])
calls = self.connection.request.call_args_list
self.assertEqual(calls[1].args[:2], ('POST', f'/api/webhooks/{WEBHOOK_ID}/fixture_token-ONLY?wait=true'))
self.assertEqual(json.loads(calls[1].kwargs['body']), announce.release_payload(event(), '1'))
self.assertEqual(
calls[2].args[:2], ('GET', f'/api/webhooks/{WEBHOOK_ID}/fixture_token-ONLY/messages/{MESSAGE_ID}')
)
def test_readback_must_match_content_mentions_and_identity(self):
for field, bad in (
('content', 'wrong'),
('mention_everyone', False),
('mention_everyone', 1),
('mentions', [{'id': '123'}]),
('mention_roles', ['123']),
('id', '999'),
('channel_id', '999'),
('webhook_id', '999'),
):
with self.subTest(field=field, bad=bad):
self.connection.reset_mock()
self.respond(metadata(), message(), dict(message(), **{field: bad}))
with self.assertRaisesRegex(announce.AnnouncementError, 'manual reconciliation'):
announce.DiscordWebhook(WEBHOOK).send(announce.release_payload(event(), '1'))
self.assertEqual(self.methods().count('POST'), 1)
def test_missing_readback_fields_fail_closed(self):
for field in message():
value = message()
del value[field]
self.respond(metadata(), message(), value)
with self.assertRaises(announce.AnnouncementError):
announce.DiscordWebhook(WEBHOOK).send(announce.release_payload(event(), '1'))
def test_unsafe_post_message_id_never_becomes_get_path(self):
for value in (None, {}, dict(message(), id='../evil'), dict(message(), id='123?x=y')):
self.connection.reset_mock()
self.respond(metadata(), value)
with self.assertRaisesRegex(announce.AnnouncementError, 'manual reconciliation'):
announce.DiscordWebhook(WEBHOOK).send(announce.release_payload(event(), '1'))
self.assertEqual(self.methods(), ['GET', 'POST'])
def test_post_failure_never_retries_and_never_logs_secret(self):
for status in (301, 302, 307, 308, 400, 401, 403, 429, 500, 204):
with self.subTest(status=status):
self.connection.reset_mock()
self.respond(metadata(), message())
responses = list(self.connection.getresponse.side_effect)
responses[1].status = status
self.connection.getresponse.side_effect = responses
with self.assertRaisesRegex(announce.AnnouncementError, 'manual reconciliation') as caught:
announce.DiscordWebhook(WEBHOOK).send(announce.release_payload(event(), '1'))
self.assertNotIn('fixture_token', str(caught.exception))
self.assertEqual(self.methods(), ['GET', 'POST'])
def test_ambiguous_timeout_never_retries_or_echoes_exception(self):
self.respond(metadata())
first = next(self.connection.getresponse.side_effect)
self.connection.getresponse.side_effect = [first, TimeoutError(WEBHOOK)]
with self.assertRaisesRegex(announce.AnnouncementError, 'manual reconciliation') as caught:
announce.DiscordWebhook(WEBHOOK).send(announce.release_payload(event(), '1'))
self.assertNotIn('fixture_token', str(caught.exception))
self.assertEqual(self.methods(), ['GET', 'POST'])
def test_malformed_json_response_is_sanitized(self):
self.respond(metadata())
response = next(self.connection.getresponse.side_effect)
response.read.return_value = WEBHOOK.encode()
self.connection.getresponse.side_effect = [response]
with self.assertRaises(announce.AnnouncementError) as caught:
announce.DiscordWebhook(WEBHOOK).validate()
self.assertNotIn('fixture_token', str(caught.exception))
def test_get_redirect_is_not_followed(self):
self.respond(metadata())
response = next(self.connection.getresponse.side_effect)
response.status = 302
response.getheader.return_value = 'https://evil.example/'
self.connection.getresponse.side_effect = [response]
with self.assertRaises(announce.AnnouncementError):
announce.DiscordWebhook(WEBHOOK).validate()
self.assertEqual(self.methods(), ['GET'])
class EntrypointTests(BaseTest):
def run_main(self, data=None, **overrides):
with tempfile.TemporaryDirectory() as directory:
path = Path(directory) / 'event.json'
path.write_text(json.dumps(event() if data is None else data))
env = {
'GITHUB_EVENT_NAME': 'release',
'GITHUB_EVENT_PATH': str(path),
'GITHUB_REPOSITORY': REPO,
'GITHUB_RUN_ATTEMPT': '1',
'DISCORD_RELEASE_WEBHOOK_URL': WEBHOOK,
}
env.update(overrides)
output = io.StringIO()
with contextlib.redirect_stdout(output), contextlib.redirect_stderr(output):
result = announce.main(env)
return result, output.getvalue()
def test_dispatch_only_validates_even_if_event_contains_release(self):
with patch('announce.DiscordWebhook') as client:
client.return_value.validate.return_value = metadata()
result, output = self.run_main(GITHUB_EVENT_NAME='workflow_dispatch')
self.assertEqual(result, 0)
client.return_value.validate.assert_called_once()
client.return_value.send.assert_not_called()
self.assertIn(GUILD_ID, output)
self.assertIn(CHANNEL_ID, output)
self.assertNotIn('fixture_token', output)
def test_production_release_sends_once(self):
with patch('announce.DiscordWebhook') as client:
client.return_value.send.return_value = MESSAGE_ID
result, output = self.run_main()
self.assertEqual(result, 0)
client.return_value.send.assert_called_once_with(announce.release_payload(event(), '1'))
self.assertIn(MESSAGE_ID, output)
def test_skipped_releases_need_no_secret_or_network(self):
for flag in ('draft', 'prerelease'):
value = event()
value['release'][flag] = True
with patch('announce.DiscordWebhook') as client:
result, _ = self.run_main(value, DISCORD_RELEASE_WEBHOOK_URL='')
self.assertEqual(result, 0)
client.assert_not_called()
def test_rerun_never_constructs_client(self):
with patch('announce.DiscordWebhook') as client:
result, output = self.run_main(GITHUB_RUN_ATTEMPT='2')
self.assertEqual(result, 1)
self.assertIn('manual reconciliation', output)
client.assert_not_called()
def test_unexpected_event_or_repository_cannot_send(self):
for overrides in (
{'GITHUB_EVENT_NAME': 'push'},
{'GITHUB_EVENT_NAME': 'pull_request'},
{'GITHUB_REPOSITORY': 'evil/LangBot'},
):
with patch('announce.DiscordWebhook') as client:
result, _ = self.run_main(**overrides)
self.assertEqual(result, 1)
client.assert_not_called()
def test_missing_secret_fails_clearly_for_send_and_validation(self):
for name in ('release', 'workflow_dispatch'):
result, output = self.run_main(GITHUB_EVENT_NAME=name, DISCORD_RELEASE_WEBHOOK_URL='')
self.assertEqual(result, 1)
self.assertIn('DISCORD_RELEASE_WEBHOOK_URL is missing', output)
def test_cli_reads_event_file_and_redacts_invalid_input(self):
with tempfile.TemporaryDirectory() as directory:
path = Path(directory) / 'event.json'
value = event()
value['release']['tag_name'] = '::error::hostile @everyone'
path.write_text(json.dumps(value))
env = dict(
os.environ,
GITHUB_EVENT_NAME='release',
GITHUB_EVENT_PATH=str(path),
GITHUB_REPOSITORY=REPO,
GITHUB_RUN_ATTEMPT='1',
DISCORD_RELEASE_WEBHOOK_URL=WEBHOOK,
)
result = subprocess.run(
[sys.executable, str(Path(__file__).with_name('announce.py'))],
env=env,
text=True,
capture_output=True,
check=False,
)
self.assertEqual(result.returncode, 1)
self.assertNotIn('hostile', result.stderr)
self.assertNotIn('fixture_token', result.stderr)
self.assertNotIn('Traceback', result.stderr)
def test_unreadable_event_fails_safely(self):
result, output = self.run_main(GITHUB_EVENT_PATH='/nonexistent/event.json')
self.assertEqual(result, 1)
self.assertNotIn('Traceback', output)
if __name__ == '__main__':
unittest.main()
+13 -32
View File
@@ -7,42 +7,23 @@ on:
jobs:
build-dev-image:
runs-on: ubuntu-latest
# 如果是tag则跳过
if: ${{ !startsWith(github.ref, 'refs/tags/') }}
permissions:
contents: read
steps:
- name: Checkout
uses: actions/checkout@v4
uses: actions/checkout@v2
with:
persist-credentials: false
- name: Set up Docker Buildx
uses: docker/setup-buildx-action@v3
- name: Generate image metadata
id: image
shell: bash
- name: Generate Tag
id: generate_tag
run: |
set -euo pipefail
branch_tag="${GITHUB_REF#refs/heads/}"
branch_tag="${branch_tag//\//-}"
echo "branch_tag=${branch_tag}" >> "$GITHUB_OUTPUT"
echo "sha_tag=sha-${GITHUB_SHA}" >> "$GITHUB_OUTPUT"
- name: Login to Docker Hub
uses: docker/login-action@v3
with:
username: ${{ secrets.DOCKER_USERNAME }}
password: ${{ secrets.DOCKER_PASSWORD }}
- name: Build and push immutable Core image
uses: docker/build-push-action@v6
with:
context: .
push: true
tags: |
rockchin/langbot:${{ steps.image.outputs.branch_tag }}
rockchin/langbot:${{ steps.image.outputs.sha_tag }}
labels: |
org.opencontainers.image.revision=${{ github.sha }}
org.opencontainers.image.source=${{ github.server_url }}/${{ github.repository }}
# 获取分支名称,把/替换为-
echo ${{ github.ref }} | sed 's/refs\/heads\///g' | sed 's/\//-/g'
echo ::set-output name=tag::$(echo ${{ github.ref }} | sed 's/refs\/heads\///g' | sed 's/\//-/g')
- name: Login to Registry
run: docker login --username=${{ secrets.DOCKER_USERNAME }} --password ${{ secrets.DOCKER_PASSWORD }}
- name: Build Docker Image
run: |
docker buildx create --name mybuilder --use
docker build -t rockchin/langbot:${{ steps.generate_tag.outputs.tag }} . --push
-78
View File
@@ -1,78 +0,0 @@
name: Build fnOS FPK
on:
workflow_dispatch:
## 发布release的时候会自动构建
release:
types: [published]
permissions:
contents: write
jobs:
build-fnos-fpk:
runs-on: ubuntu-latest
steps:
- name: Checkout
uses: actions/checkout@v2
with:
persist-credentials: false
- name: Check version
id: check_version
run: |
echo $GITHUB_REF
# 如果是tag,则去掉refs/tags/前缀(与其他 release workflow 一致,版本号取 tag 名)
if [[ $GITHUB_REF == refs/tags/* ]]; then
echo "It's a tag"
echo "version=$(echo $GITHUB_REF | awk -F '/' '{print $3}')" >> $GITHUB_OUTPUT
else
# 手动触发(workflow_dispatch):读不到 tag,使用 manifest 内维护的版本
echo "It's not a tag"
echo "version=$(grep '^version=' packaging/fnos/manifest | cut -d= -f2)" >> $GITHUB_OUTPUT
fi
- name: Setup Node
uses: actions/setup-node@v2
with:
node-version: '22'
- name: Setup Python
uses: actions/setup-python@v5
with:
python-version: '3.12'
- name: Install build tools
run: |
pip install pillow
# fnpack:飞牛官方打包 CLI(静态二进制)
curl -fsSL -o /usr/local/bin/fnpack \
https://static2.fnnas.com/fnpack/fnpack-1.2.3-linux-amd64
chmod +x /usr/local/bin/fnpack
- name: Build FPK
env:
FPK_VERSION: ${{ steps.check_version.outputs.version }}
run: |
bash packaging/fnos/build.sh
test -f packaging/fnos/langbot.fpk
- name: Upload Artifact
uses: actions/upload-artifact@v4
with:
name: langbot-${{ steps.check_version.outputs.version }}-fnos
path: packaging/fnos/langbot.fpk
- name: Upload To Release
# 仅 release 触发时执行;手动/workflow_dispatch 触发时没有 release
# 且 github.event.release.tag_name 为空(否则 gh release upload 缺参数报错)
if: github.event_name == 'release'
env:
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
run: |
# tag 可能已带 -fnos 后缀(如 v4.10.9-fnos),产物名统一为
# langbot-<基础版本>-fnos.fpk,避免出现 -fnos-fnos
VER="${{ steps.check_version.outputs.version }}"
BASE="${VER#v}"; BASE="${BASE%-fnos}"
cp packaging/fnos/langbot.fpk "langbot-${BASE}-fnos.fpk"
gh release upload ${{ github.event.release.tag_name }} "langbot-${BASE}-fnos.fpk"
+59
View File
@@ -0,0 +1,59 @@
name: Build and deploy production
on:
push:
branches: [deploy/prod]
workflow_dispatch:
permissions:
contents: read
concurrency:
group: langbot-production
cancel-in-progress: false
env:
CORE_IMAGE: ${{ secrets.DOCKER_USERNAME }}/langbot
CLOUD_IMAGE: ${{ secrets.DOCKER_USERNAME }}/langbot-cloud-core
SPACE_REF: 58253c53933f95d81b035fbe2efedb55b6c1a82b
jobs:
build-and-deploy:
runs-on: ubuntu-latest
environment: production
steps:
- uses: actions/checkout@v4
- uses: docker/setup-buildx-action@v3
- uses: docker/login-action@v3
with:
username: ${{ secrets.DOCKER_USERNAME }}
password: ${{ secrets.DOCKER_PASSWORD }}
- name: Build exact Core image
uses: docker/build-push-action@v6
with:
context: .
push: true
tags: |
${{ env.CORE_IMAGE }}:prod-${{ github.sha }}
${{ env.CORE_IMAGE }}:deploy-prod
cache-from: type=gha,scope=core-prod
cache-to: type=gha,mode=max,scope=core-prod
- name: Checkout production Cloud adapter
uses: actions/checkout@v4
with:
repository: langbot-app/langbot-space
ref: ${{ env.SPACE_REF }}
token: ${{ secrets.CLA_PAT }}
path: .space
- name: Build exact Cloud Core image
uses: docker/build-push-action@v6
with:
context: .space
file: .space/Dockerfile.cloud
push: true
build-args: LANGBOT_CORE_IMAGE=${{ env.CORE_IMAGE }}:prod-${{ github.sha }}
tags: |
${{ env.CLOUD_IMAGE }}:prod-${{ github.sha }}
${{ env.CLOUD_IMAGE }}:deploy-prod
cache-from: type=gha,scope=cloud-core-prod
cache-to: type=gha,mode=max,scope=cloud-core-prod
-64
View File
@@ -1,64 +0,0 @@
name: Discord Release Announcement
on:
release:
types: [published]
workflow_dispatch:
push:
paths:
- '.github/workflows/discord-release.yml'
- '.github/discord-release/**'
pull_request:
paths:
- '.github/workflows/discord-release.yml'
- '.github/discord-release/**'
permissions:
contents: read
jobs:
tests:
name: Offline announcement tests
runs-on: ubuntu-24.04
timeout-minutes: 5
steps:
- uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 # v4
with:
persist-credentials: false
- name: Test helper without secrets or network
run: python3 -m unittest discover -s .github/discord-release -p 'test_*.py' -v
validate:
name: Validate webhook (GET only, no message)
if: github.repository == 'langbot-app/LangBot' && github.event_name == 'workflow_dispatch'
needs: tests
runs-on: ubuntu-24.04
timeout-minutes: 5
steps:
- uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 # v4
with:
persist-credentials: false
- name: Validate incoming webhook and report guild/channel IDs
env:
DISCORD_RELEASE_WEBHOOK_URL: ${{ secrets.DISCORD_RELEASE_WEBHOOK_URL }}
run: python3 .github/discord-release/announce.py
announce:
name: Announce published stable release
if: >-
github.repository == 'langbot-app/LangBot' &&
github.event_name == 'release' && github.event.action == 'published' &&
github.event.release.draft == false && github.event.release.prerelease == false
needs: tests
runs-on: ubuntu-24.04
timeout-minutes: 5
steps:
- uses: actions/checkout@11d5960a326750d5838078e36cf38b85af677262 # v4
with:
persist-credentials: false
# The helper refuses GITHUB_RUN_ATTEMPT != 1 with recovery guidance.
# Never interpolate release data into a shell command.
- name: Send once and verify the exact Discord message
env:
DISCORD_RELEASE_WEBHOOK_URL: ${{ secrets.DISCORD_RELEASE_WEBHOOK_URL }}
run: python3 .github/discord-release/announce.py
+4 -35
View File
@@ -2,11 +2,6 @@ name: Build and Publish to PyPI
on:
workflow_dispatch:
inputs:
source_ref:
description: 'Existing release tag to publish (for example v4.10.11)'
required: true
type: string
release:
types: [published]
@@ -16,39 +11,13 @@ jobs:
permissions:
contents: read
id-token: write # Required for trusted publishing to PyPI
steps:
- name: Checkout code
uses: actions/checkout@v4
with:
ref: ${{ inputs.source_ref || github.sha }}
fetch-depth: 0
persist-credentials: false
- name: Validate release source and version
env:
RELEASE_TAG: ${{ inputs.source_ref || github.event.release.tag_name }}
run: |
python3 - <<'PY'
import os
import re
import subprocess
import tomllib
from pathlib import Path
tag = os.environ['RELEASE_TAG']
if not re.fullmatch(r'v[0-9]+\.[0-9]+\.[0-9]+', tag):
raise SystemExit('source_ref must be an existing release tag: vX.Y.Z')
def revision(ref):
return subprocess.check_output(['git', 'rev-parse', '--verify', ref], text=True).strip()
if revision('HEAD') != revision(f'refs/tags/{tag}^{{}}'):
raise SystemExit('Checked-out commit does not match the release tag')
version = tomllib.loads(Path('pyproject.toml').read_text())['project']['version']
if version != tag[1:]:
raise SystemExit(f'Package version {version} does not match tag {tag}')
print(f'Validated {tag} at {revision("HEAD")} (package {version})')
PY
- name: Set up Node.js
uses: actions/setup-node@v4
with:
@@ -57,9 +26,9 @@ jobs:
- name: Build frontend
run: |
cd web
# Match the archive/Docker npm path; npm ci rejects older tags' stale npm lockfiles.
npm install --include=optional
npm run build
npm install -g pnpm
pnpm install
pnpm build
mkdir -p ../src/langbot/web/dist
cp -r dist ../src/langbot/web/
-6
View File
@@ -10,16 +10,12 @@ on:
- 'src/langbot/pkg/persistence/**'
- 'src/langbot/pkg/entity/persistence/**'
- 'tests/integration/persistence/**'
- 'tests/unit_tests/api/service/test_monitoring_sessions.py'
- '.github/workflows/test-migrations.yml'
pull_request:
types: [opened, synchronize, reopened, ready_for_review]
paths:
- 'src/langbot/pkg/persistence/**'
- 'src/langbot/pkg/entity/persistence/**'
- 'tests/integration/persistence/**'
- 'tests/unit_tests/api/service/test_monitoring_sessions.py'
- '.github/workflows/test-migrations.yml'
jobs:
test-migrations-sqlite:
@@ -84,8 +80,6 @@ jobs:
run: >-
uv run pytest
tests/integration/persistence/test_migrations_postgres.py
tests/integration/persistence/test_monitoring_postgres.py
tests/unit_tests/api/service/test_monitoring_sessions.py::test_postgres_upgrade_rls_and_concurrent_bot_counts
tests/integration/persistence/test_pgvector_postgres.py
tests/integration/persistence/test_release_migration_postgres.py
tests/integration/persistence/test_plugin_identity_migration.py
-12
View File
@@ -57,15 +57,3 @@ testsdk/
# Next.js build cache (legacy)
web/.next/
web/.pnpm-home
.tmp
Caddyfile
# fnOS packaging build artifacts (packaging/fnos/build.sh)
packaging/fnos/app/langbot/
packaging/fnos/app/bin/
packaging/fnos/ICON.PNG
packaging/fnos/ICON_256.PNG
packaging/fnos/app/ui/images/
packaging/fnos/app/desktop/images/
packaging/fnos/*.fpk
+2 -2
View File
@@ -43,8 +43,8 @@ Run the narrowest useful test first, then broader checks when confidence is need
## Where to Look
- Architecture map: `ARCHITECTURE.md`.
- Dev environment guide: https://langbot.app/docs/zh/develop/dev-config.
- Plugin runtime / CLI / SDK debugging: https://langbot.app/docs/zh/develop/plugin-runtime.
- Dev environment guide: https://docs.langbot.app/zh/develop/dev-config.
- Plugin runtime / CLI / SDK debugging: https://docs.langbot.app/zh/develop/plugin-runtime.
- API-key auth: `docs/API_KEY_AUTH.md`.
- Box deep-dive notes: `docs/review/box-architecture.md` and related files.
- In-repo skills: `skills/` is the single source of truth for LangBot agent skills.
+2 -2
View File
@@ -1,4 +1,4 @@
FROM --platform=$BUILDPLATFORM node:22-alpine AS node
FROM node:22-alpine AS node
WORKDIR /app
@@ -62,7 +62,7 @@ RUN apt-get update \
&& apt-get install -y --no-install-recommends nodejs \
&& rm -f /tmp/nodesource_setup.sh \
&& python -m pip install --no-cache-dir uv \
&& uv sync --extra seekdb \
&& uv sync \
&& apt-get purge -y --auto-remove curl git gnupg \
&& rm -rf /var/lib/apt/lists/* \
&& touch /.dockerenv
+17 -7
View File
@@ -19,9 +19,9 @@ English / [简体中文](README_CN.md) / [繁體中文](README_TW.md) / [日本
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Website</a>
<a href="https://langbot.app/docs/en/insight/features">Features</a>
<a href="https://langbot.app/docs/en/insight/guide">Docs</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Features</a>
<a href="https://link.langbot.app/en/docs/guide">Docs</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app/cloud">Cloud</a>
<a href="https://space.langbot.app">Plugin Market</a>
<a href="https://langbot.featurebase.app/roadmap">Roadmap</a>
@@ -49,7 +49,7 @@ LangBot is an **open-source, production-grade platform** for building AI-powered
- **Web Management Panel** — Configure, manage, and monitor your bots through an intuitive browser interface. No YAML editing required.
- **Multi-Pipeline Architecture** — Different bots for different scenarios, with comprehensive monitoring and exception handling.
[→ Learn more about all features](https://langbot.app/docs/en/insight/features)
[→ Learn more about all features](https://link.langbot.app/en/docs/features)
📍 Practical guides: [deploy a multi-platform AI bot in 5 minutes](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [connect DeepSeek to WeChat, Discord, and Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [run a Dify Agent in Discord, Telegram, and Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/), and [build an n8n-powered chatbot](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -83,13 +83,23 @@ cd LangBot/docker
docker compose --profile all up -d
```
### One-Click Cloud Deploy
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**More options:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Manual](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**More options:** [Docker](https://link.langbot.app/en/docs/docker) · [Manual](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
## Live Demo
**Try it now:** https://demo.langbot.dev/
- Email: `demo@langbot.app`
- Password: `langbot123456`
_Note: Public demo environment. Do not enter sensitive information._
---
@@ -140,7 +150,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | Gateway | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Gateway | ✅ |
[→ View all integrations](https://langbot.app/docs/en/insight/features)
[→ View all integrations](https://link.langbot.app/en/docs/features)
---
+16 -7
View File
@@ -21,9 +21,9 @@
[![star](https://gitcode.com/RockChinQ/LangBot/star/badge.svg)](https://gitcode.com/RockChinQ/LangBot)
<a href="https://langbot.app">官网</a>
<a href="https://langbot.app/docs/zh/insight/features">特性</a>
<a href="https://langbot.app/docs/zh/insight/guide">文档</a>
<a href="https://langbot.app/docs/zh/tags/readme">API</a>
<a href="https://link.langbot.app/zh/docs/features">特性</a>
<a href="https://link.langbot.app/zh/docs/guide">文档</a>
<a href="https://link.langbot.app/zh/docs/api">API</a>
<a href="https://space.langbot.app/cloud">Cloud</a>
<a href="https://space.langbot.app">扩展市场</a>
<a href="https://langbot.featurebase.app/roadmap">路线图</a>
@@ -49,7 +49,7 @@ LangBot 是一个**开源的生产级平台**,用于构建 AI 驱动的即时
- **Web 管理面板** — 通过浏览器直观地配置、管理和监控机器人,无需手动编辑配置文件。
- **多流水线架构** — 不同机器人用于不同场景,具备全面的监控和异常处理能力。
[→ 了解更多功能特性](https://langbot.app/docs/zh/insight/features)
[→ 了解更多功能特性](https://link.langbot.app/zh/docs/features)
📍 实践指南:[5 分钟部署多平台 AI 机器人](https://langbot.app/zh/blog/deploy-ai-bot-in-5-minutes/)、[将 DeepSeek 接入微信、企业微信与 Discord](https://langbot.app/zh/blog/connect-deepseek-to-wechat/)、[让 Dify Agent 跑在 Discord、Telegram 和 Slack 上](https://langbot.app/zh/blog/dify-agent-discord-telegram-slack/),以及[用 n8n 构建多平台 AI 聊天机器人](https://langbot.app/zh/blog/n8n-multi-platform-ai-chatbot/)。
@@ -83,13 +83,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### 一键云部署
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/zh-CN/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**更多方式:** [Docker](https://langbot.app/docs/zh/deploy/langbot/docker) · [手动部署](https://langbot.app/docs/zh/deploy/langbot/manual) · [宝塔面板](https://langbot.app/docs/zh/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/zh/deploy/langbot/kubernetes)
**更多方式:** [Docker](https://link.langbot.app/zh/docs/docker) · [手动部署](https://link.langbot.app/zh/docs/manual-deploy) · [宝塔面板](https://link.langbot.app/zh/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/zh/deploy/langbot/kubernetes)
---
## 在线演示
**立即体验:** https://demo.langbot.dev/
- 邮箱:`demo@langbot.app`
- 密码:`langbot123456`
*注意:公开演示环境,请不要在其中填入任何敏感信息。*
---
@@ -142,7 +151,7 @@ docker compose --profile all up -d
| [百宝箱Tbox](https://www.tbox.cn/open) | 智能体平台 | ✅ |
| [七牛云Qiniu](https://www.qiniu.com/ai/agent) | 聚合平台 | ✅ |
[→ 查看完整集成列表](https://langbot.app/docs/zh/insight/features)
[→ 查看完整集成列表](https://link.langbot.app/zh/docs/features)
### TTS(语音合成)
+16 -7
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Inicio</a>
<a href="https://langbot.app/docs/en/insight/features">Características</a>
<a href="https://langbot.app/docs/en/insight/guide">Documentación</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Características</a>
<a href="https://link.langbot.app/en/docs/guide">Documentación</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Mercado de Plugins</a>
<a href="https://langbot.featurebase.app/roadmap">Hoja de Ruta</a>
@@ -48,7 +48,7 @@ LangBot es una **plataforma de código abierto y grado de producción** para con
- **Panel de Gestión Web** — Configure, gestione y monitoree sus bots a través de una interfaz de navegador intuitiva. Sin necesidad de editar YAML.
- **Arquitectura Multi-Pipeline** — Diferentes bots para diferentes escenarios, con monitoreo completo y manejo de excepciones.
[→ Conocer más sobre todas las funcionalidades](https://langbot.app/docs/en/insight/features)
[→ Conocer más sobre todas las funcionalidades](https://link.langbot.app/en/docs/features)
📍 Guías prácticas: [desplegar un bot de IA multiplataforma en 5 minutos](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [conectar DeepSeek a WeChat, Discord y Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [ejecutar un Dify Agent en Discord, Telegram y Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) y [crear un chatbot con n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -82,13 +82,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### Despliegue en la Nube con un Clic
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Más opciones:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Manual](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Más opciones:** [Docker](https://link.langbot.app/en/docs/docker) · [Manual](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
## Demo en Vivo
**Pruébelo ahora:** https://demo.langbot.dev/
- Correo electrónico: `demo@langbot.app`
- Contraseña: `langbot123456`
*Nota: Entorno de demostración público. No ingrese información confidencial.*
---
@@ -139,7 +148,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | Pasarela | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Pasarela | ✅ |
[→ Ver todas las integraciones](https://langbot.app/docs/en/insight/features)
[→ Ver todas las integraciones](https://link.langbot.app/en/docs/features)
---
+16 -7
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Accueil</a>
<a href="https://langbot.app/docs/en/insight/features">Fonctionnalités</a>
<a href="https://langbot.app/docs/en/insight/guide">Documentation</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Fonctionnalités</a>
<a href="https://link.langbot.app/en/docs/guide">Documentation</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Marché des Plugins</a>
<a href="https://langbot.featurebase.app/roadmap">Feuille de Route</a>
@@ -48,7 +48,7 @@ LangBot est une **plateforme open-source de niveau production** pour créer des
- **Panneau de Gestion Web** — Configurez, gérez et surveillez vos bots via une interface navigateur intuitive. Aucune édition de YAML requise.
- **Architecture Multi-Pipeline** — Différents bots pour différents scénarios, avec surveillance complète et gestion des exceptions.
[→ En savoir plus sur toutes les fonctionnalités](https://langbot.app/docs/en/insight/features)
[→ En savoir plus sur toutes les fonctionnalités](https://link.langbot.app/en/docs/features)
📍 Guides pratiques : [déployer un bot IA multiplateforme en 5 minutes](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [connecter DeepSeek à WeChat, Discord et Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [exécuter un Dify Agent dans Discord, Telegram et Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) et [créer un chatbot avec n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -82,13 +82,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### Déploiement Cloud en un Clic
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Plus d'options :** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Manuel](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Plus d'options :** [Docker](https://link.langbot.app/en/docs/docker) · [Manuel](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
## Démo en Ligne
**Essayez maintenant :** https://demo.langbot.dev/
- Email : `demo@langbot.app`
- Mot de passe : `langbot123456`
*Note : Environnement de démonstration public. Ne saisissez pas d'informations sensibles.*
---
@@ -139,7 +148,7 @@ docker compose --profile all up -d
| [ShengSuanYun](https://www.shengsuanyun.com/?from=CH_KYIPP758) | Plateforme GPU | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Passerelle | ✅ |
[→ Voir toutes les intégrations](https://langbot.app/docs/en/insight/features)
[→ Voir toutes les intégrations](https://link.langbot.app/en/docs/features)
---
+16 -7
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">ホーム</a>
<a href="https://langbot.app/docs/ja/insight/features">機能</a>
<a href="https://langbot.app/docs/ja/insight/guide">ドキュメント</a>
<a href="https://langbot.app/docs/ja/tags/readme">API</a>
<a href="https://link.langbot.app/ja/docs/features">機能</a>
<a href="https://link.langbot.app/ja/docs/guide">ドキュメント</a>
<a href="https://link.langbot.app/ja/docs/api">API</a>
<a href="https://space.langbot.app">プラグインマーケット</a>
<a href="https://langbot.featurebase.app/roadmap">ロードマップ</a>
@@ -48,7 +48,7 @@ LangBot は、AI搭載のインスタントメッセージングボットを構
- **Web管理パネル** — 直感的なブラウザインターフェースからボットの設定、管理、監視が可能。YAML編集は不要。
- **マルチパイプラインアーキテクチャ** — 異なるシナリオに異なるボットを配置し、包括的な監視と例外処理を実現。
[→ すべての機能について詳しく見る](https://langbot.app/docs/ja/insight/features)
[→ すべての機能について詳しく見る](https://link.langbot.app/ja/docs/features)
📍 実践ガイド: [5分でマルチプラットフォームAIボットをデプロイ](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/)、[DeepSeekをWeChat・Discord・Telegramに接続](https://langbot.app/en/blog/connect-deepseek-to-wechat/)、[Dify AgentをDiscord・Telegram・Slackで動かす](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/)、[n8n連携チャットボットを構築](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/)。
@@ -82,13 +82,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### ワンクリッククラウドデプロイ
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**その他:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [手動デプロイ](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**その他:** [Docker](https://link.langbot.app/en/docs/docker) · [手動デプロイ](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
## ライブデモ
**今すぐ試す:** https://demo.langbot.dev/
- メール: `demo@langbot.app`
- パスワード: `langbot123456`
*注意: 公開デモ環境です。機密情報を入力しないでください。*
---
@@ -139,7 +148,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | ゲートウェイ | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | ゲートウェイ | ✅ |
[→ すべての統合を表示](https://langbot.app/docs/en/insight/features)
[→ すべての統合を表示](https://link.langbot.app/en/docs/features)
---
+16 -7
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">홈</a>
<a href="https://langbot.app/docs/en/insight/features">기능</a>
<a href="https://langbot.app/docs/en/insight/guide">문서</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">기능</a>
<a href="https://link.langbot.app/en/docs/guide">문서</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">플러그인 마켓</a>
<a href="https://langbot.featurebase.app/roadmap">로드맵</a>
@@ -48,7 +48,7 @@ LangBot은 AI 기반 인스턴트 메시징 봇을 구축하기 위한 **오픈
- **웹 관리 패널** — 직관적인 브라우저 인터페이스로 봇을 구성, 관리 및 모니터링. YAML 편집 불필요.
- **멀티 파이프라인 아키텍처** — 다양한 시나리오에 맞는 다양한 봇 구성, 종합 모니터링 및 예외 처리.
[→ 모든 기능 자세히 보기](https://langbot.app/docs/en/insight/features)
[→ 모든 기능 자세히 보기](https://link.langbot.app/en/docs/features)
📍 실전 가이드: [5분 만에 멀티 플랫폼 AI 봇 배포하기](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [DeepSeek를 WeChat, Discord, Telegram에 연결하기](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [Dify Agent를 Discord, Telegram, Slack에서 실행하기](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/), [n8n 기반 챗봇 만들기](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -82,13 +82,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### 원클릭 클라우드 배포
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**더 많은 옵션:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [수동 배포](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**더 많은 옵션:** [Docker](https://link.langbot.app/en/docs/docker) · [수동 배포](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
## 라이브 데모
**지금 체험:** https://demo.langbot.dev/
- 이메일: `demo@langbot.app`
- 비밀번호: `langbot123456`
*참고: 공개 데모 환경입니다. 민감한 정보를 입력하지 마세요.*
---
@@ -139,7 +148,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | 게이트웨이 | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | 게이트웨이 | ✅ |
[→ 모든 통합 보기](https://langbot.app/docs/en/insight/features)
[→ 모든 통합 보기](https://link.langbot.app/en/docs/features)
---
+16 -7
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Главная</a>
<a href="https://langbot.app/docs/en/insight/features">Возможности</a>
<a href="https://langbot.app/docs/en/insight/guide">Документация</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Возможности</a>
<a href="https://link.langbot.app/en/docs/guide">Документация</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Магазин плагинов</a>
<a href="https://langbot.featurebase.app/roadmap">Дорожная карта</a>
@@ -48,7 +48,7 @@ LangBot — это **платформа с открытым исходным к
- **Веб-панель управления** — Настраивайте, управляйте и мониторьте ваших ботов через интуитивный браузерный интерфейс. Ручное редактирование YAML не требуется.
- **Мультиконвейерная архитектура** — Разные боты для разных сценариев с комплексным мониторингом и обработкой исключений.
[→ Подробнее обо всех возможностях](https://langbot.app/docs/en/insight/features)
[→ Подробнее обо всех возможностях](https://link.langbot.app/en/docs/features)
📍 Практические руководства: [развернуть мультиплатформенного ИИ-бота за 5 минут](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [подключить DeepSeek к WeChat, Discord и Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [запустить Dify Agent в Discord, Telegram и Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) и [создать чат-бота на n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -82,13 +82,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### Облачное развертывание одним кликом
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Другие варианты:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Ручная установка](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Другие варианты:** [Docker](https://link.langbot.app/en/docs/docker) · [Ручная установка](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
## Демо
**Попробуйте прямо сейчас:** https://demo.langbot.dev/
- Email: `demo@langbot.app`
- Пароль: `langbot123456`
*Примечание: Публичная демо-среда. Не вводите конфиденциальную информацию.*
---
@@ -139,7 +148,7 @@ docker compose --profile all up -d
| [ShengSuanYun](https://www.shengsuanyun.com/?from=CH_KYIPP758) | Платформа GPU | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Шлюз | ✅ |
[→ Смотреть все интеграции](https://langbot.app/docs/en/insight/features)
[→ Смотреть все интеграции](https://link.langbot.app/en/docs/features)
---
+16 -7
View File
@@ -21,9 +21,9 @@
[![star](https://gitcode.com/RockChinQ/LangBot/star/badge.svg)](https://gitcode.com/RockChinQ/LangBot)
<a href="https://langbot.app">官網</a>
<a href="https://langbot.app/docs/zh/insight/features">特性</a>
<a href="https://langbot.app/docs/zh/insight/guide">文件</a>
<a href="https://langbot.app/docs/zh/tags/readme">API</a>
<a href="https://link.langbot.app/zh/docs/features">特性</a>
<a href="https://link.langbot.app/zh/docs/guide">文件</a>
<a href="https://link.langbot.app/zh/docs/api">API</a>
<a href="https://space.langbot.app">外掛市場</a>
<a href="https://langbot.featurebase.app/roadmap">路線圖</a>
@@ -50,7 +50,7 @@ LangBot 是一個**開源的生產級平台**,用於建構 AI 驅動的即時
- **Web 管理面板** — 透過瀏覽器直觀地配置、管理和監控機器人,無需手動編輯設定檔。
- **多流水線架構** — 不同機器人用於不同場景,具備全面的監控和異常處理能力。
[→ 了解更多功能特性](https://langbot.app/docs/zh/insight/features)
[→ 了解更多功能特性](https://link.langbot.app/zh/docs/features)
📍 實踐指南:[5 分鐘部署多平台 AI 機器人](https://langbot.app/zh/blog/deploy-ai-bot-in-5-minutes/)、[將 DeepSeek 接入微信、企業微信與 Discord](https://langbot.app/zh/blog/connect-deepseek-to-wechat/)、[讓 Dify Agent 跑在 Discord、Telegram 和 Slack 上](https://langbot.app/zh/blog/dify-agent-discord-telegram-slack/),以及[用 n8n 建構多平台 AI 聊天機器人](https://langbot.app/zh/blog/n8n-multi-platform-ai-chatbot/)。
@@ -84,13 +84,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### 一鍵雲端部署
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/zh-CN/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**更多方式:** [Docker](https://langbot.app/docs/zh/deploy/langbot/docker) · [手動部署](https://langbot.app/docs/zh/deploy/langbot/manual) · [寶塔面板](https://langbot.app/docs/zh/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/zh/deploy/langbot/kubernetes)
**更多方式:** [Docker](https://link.langbot.app/zh/docs/docker) · [手動部署](https://link.langbot.app/zh/docs/manual-deploy) · [寶塔面板](https://link.langbot.app/zh/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/zh/deploy/langbot/kubernetes)
---
## 線上演示
**立即體驗:** https://demo.langbot.dev/
- 信箱:`demo@langbot.app`
- 密碼:`langbot123456`
*注意:公開演示環境,請不要在其中填入任何敏感資訊。*
---
@@ -155,7 +164,7 @@ docker compose --profile all up -d
|-----------|------|
| 阿里雲百煉 | [外掛](https://github.com/Thetail001/LangBot_BailianTextToImagePlugin) |
[→ 查看完整整合列表](https://langbot.app/docs/zh/insight/features)
[→ 查看完整整合列表](https://link.langbot.app/zh/docs/features)
---
+16 -7
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Trang chủ</a>
<a href="https://langbot.app/docs/en/insight/features">Tính năng</a>
<a href="https://langbot.app/docs/en/insight/guide">Tài liệu</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Tính năng</a>
<a href="https://link.langbot.app/en/docs/guide">Tài liệu</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Chợ Plugin</a>
<a href="https://langbot.featurebase.app/roadmap">Lộ trình</a>
@@ -48,7 +48,7 @@ LangBot là một **nền tảng mã nguồn mở, cấp sản xuất** để x
- **Bảng quản lý Web** — Cấu hình, quản lý và giám sát bot thông qua giao diện trình duyệt trực quan. Không cần chỉnh sửa YAML.
- **Kiến trúc đa Pipeline** — Các bot khác nhau cho các kịch bản khác nhau, với giám sát toàn diện và xử lý ngoại lệ.
[→ Tìm hiểu thêm về tất cả tính năng](https://langbot.app/docs/en/insight/features)
[→ Tìm hiểu thêm về tất cả tính năng](https://link.langbot.app/en/docs/features)
📍 Hướng dẫn thực hành: [triển khai bot AI đa nền tảng trong 5 phút](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [kết nối DeepSeek với WeChat, Discord và Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [chạy Dify Agent trên Discord, Telegram và Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) và [xây dựng chatbot với n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -82,13 +82,22 @@ cd LangBot/docker
docker compose --profile all up -d
```
### Triển khai đám mây một cú nhấp
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Thêm tùy chọn:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Thủ công](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Thêm tùy chọn:** [Docker](https://link.langbot.app/en/docs/docker) · [Thủ công](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
## Demo trực tuyến
**Thử ngay:** https://demo.langbot.dev/
- Email: `demo@langbot.app`
- Mật khẩu: `langbot123456`
*Lưu ý: Môi trường demo công khai. Không nhập thông tin nhạy cảm.*
---
@@ -139,7 +148,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | Cổng | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Cổng | ✅ |
[→ Xem tất cả tích hợp](https://langbot.app/docs/en/insight/features)
[→ Xem tất cả tích hợp](https://link.langbot.app/en/docs/features)
---
+97
View File
@@ -0,0 +1,97 @@
#!/usr/bin/env bash
set -Eeuo pipefail
cd /opt/langbot-cloud-prod
TAG=${1:?usage: deploy.sh prod-<40-char-sha>}
[[ "$TAG" =~ ^prod-[0-9a-f]{40}$ ]] || { echo 'invalid immutable image tag' >&2; exit 2; }
[[ -s .env ]] || { echo '/opt/langbot-cloud-prod/.env is missing' >&2; exit 3; }
rendered_compose=$(docker compose config)
grep -Fq 'LANGBOT_SPACE_CONTROL_PLANE_URL: https://space.langbot.app' <<<"$rendered_compose" || {
echo 'Cloud control-plane URL must be https://space.langbot.app' >&2
exit 4
}
grep -Fq 'SPACE__URL: https://space.langbot.app' <<<"$rendered_compose" || {
echo 'Cloud user-facing Space URL must be https://space.langbot.app' >&2
exit 5
}
grep -Eq 'LANGBOT_TELEMETRY_INGEST_TOKEN: .+' <<<"$rendered_compose" || {
echo 'Cloud telemetry ingest token must be configured' >&2
exit 6
}
update_env() {
local key=$1 value=$2
python3 - "$key" "$value" <<'PY'
from pathlib import Path
import os
import sys
path = Path('.env')
key, value = sys.argv[1:]
lines = path.read_text().splitlines()
updated = False
for index, line in enumerate(lines):
if line.startswith(f'{key}='):
lines[index] = f'{key}={value}'
updated = True
break
if not updated:
lines.append(f'{key}={value}')
temporary = Path('.env.tmp')
temporary.write_text('\n'.join(lines) + '\n')
os.chmod(temporary, 0o600)
temporary.replace(path)
PY
}
update_env LANGBOT_IMAGE_TAG "$TAG"
set -a
. ./.env
set +a
: "${CLOUD_V2_CONTROL_PLANE_TOKEN:?CLOUD_V2_CONTROL_PLANE_TOKEN is required}"
for attempt in 1 2 3 4 5; do
if docker compose pull postgres redis migrate plugin-runtime core; then
break
fi
if [ "$attempt" -eq 5 ]; then
echo "docker compose pull failed after $attempt attempts" >&2
exit 1
fi
delay=$((attempt * 10))
echo "docker compose pull failed (attempt $attempt/5); retrying in ${delay}s" >&2
sleep "$delay"
done
docker compose up -d postgres redis
for _ in $(seq 1 60); do
if docker compose exec -T postgres pg_isready -U langbot_operator -d langbot >/dev/null 2>&1; then break; fi
sleep 2
done
docker compose exec -T postgres pg_isready -U langbot_operator -d langbot >/dev/null
docker compose exec -T postgres psql -v ON_ERROR_STOP=1 -U langbot_operator -d langbot \
-v runtime_password="$POSTGRES_RUNTIME_PASSWORD" <<'SQL'
SELECT format('CREATE ROLE langbot_runtime LOGIN PASSWORD %L', :'runtime_password')
WHERE NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'langbot_runtime')\gexec
ALTER ROLE langbot_runtime PASSWORD :'runtime_password';
GRANT CONNECT ON DATABASE langbot TO langbot_runtime;
REVOKE CREATE ON SCHEMA public FROM PUBLIC, langbot_runtime;
REVOKE ALL PRIVILEGES ON ALL TABLES IN SCHEMA public FROM langbot_runtime;
REVOKE ALL PRIVILEGES ON ALL SEQUENCES IN SCHEMA public FROM langbot_runtime;
ALTER DEFAULT PRIVILEGES FOR ROLE langbot_operator IN SCHEMA public REVOKE ALL ON TABLES FROM langbot_runtime;
ALTER DEFAULT PRIVILEGES FOR ROLE langbot_operator IN SCHEMA public REVOKE ALL ON SEQUENCES FROM langbot_runtime;
GRANT USAGE ON SCHEMA public TO langbot_runtime;
SQL
docker compose --profile tools run --rm migrate
docker compose up -d --remove-orphans plugin-runtime core
for _ in $(seq 1 90); do
if docker compose exec -T core python -c 'import urllib.request; urllib.request.urlopen("http://127.0.0.1:5300/healthz", timeout=3)' >/dev/null 2>&1; then
docker compose ps
exit 0
fi
sleep 2
done
docker compose logs --tail=200 core plugin-runtime >&2
exit 1
+162
View File
@@ -0,0 +1,162 @@
services:
postgres:
image: pgvector/pgvector:pg17
container_name: langbot-cloud-postgres
restart: unless-stopped
environment:
POSTGRES_DB: langbot
POSTGRES_USER: langbot_operator
POSTGRES_PASSWORD: ${POSTGRES_OPERATOR_PASSWORD}
volumes:
- postgres-data:/var/lib/postgresql/data
healthcheck:
test: [CMD-SHELL, "pg_isready -U langbot_operator -d langbot"]
interval: 5s
timeout: 5s
retries: 30
networks: [internal]
redis:
image: redis:7.4-alpine
container_name: langbot-cloud-redis
restart: unless-stopped
command: [redis-server, --appendonly, "yes", --requirepass, "${REDIS_PASSWORD}"]
volumes:
- redis-data:/data
healthcheck:
test: [CMD-SHELL, "redis-cli -a \"$${REDIS_PASSWORD}\" ping | grep PONG"]
interval: 5s
timeout: 5s
retries: 20
environment:
REDIS_PASSWORD: ${REDIS_PASSWORD}
networks: [internal]
migrate:
image: rockchin/langbot-cloud-core:${LANGBOT_IMAGE_TAG}
profiles: [tools]
command: [uv, run, langbot, migrate, --cloud]
environment: &core-env
TZ: Asia/Shanghai
SYSTEM__INSTANCE_ID: ${CLOUD_V2_INSTANCE_UUID}
SYSTEM__EDITION: cloud
SYSTEM__RECOVERY_KEY: ${SYSTEM_RECOVERY_KEY}
SYSTEM__JWT__SECRET: ${JWT_SECRET}
SYSTEM__LIMITATION__MAX_BOTS: "2"
SYSTEM__LIMITATION__MAX_PIPELINES: "3"
SYSTEM__LIMITATION__MAX_EXTENSIONS: "3"
SYSTEM__LIMITATION__MAX_KNOWLEDGE_BASES: "2"
API__WEBHOOK_PREFIX: https://cloud.langbot.app
API__WEBUI_URL: https://cloud.langbot.app
WORKSPACE__INVITATIONS__PUBLIC_WEB_URL: https://cloud.langbot.app
DATABASE__USE: postgresql
DATABASE__POSTGRESQL__URL: postgresql+asyncpg://langbot_runtime:${POSTGRES_RUNTIME_PASSWORD}@postgres:5432/langbot
DATABASE__CLOUD_MIGRATION__OPERATOR_DSN_ENV: LANGBOT_CLOUD_MIGRATION_DSN
LANGBOT_CLOUD_MIGRATION_DSN: postgresql://langbot_operator:${POSTGRES_OPERATOR_PASSWORD}@postgres:5432/langbot
VDB__USE: pgvector
VDB__PGVECTOR__USE_BUSINESS_DATABASE: "true"
VDB__PGVECTOR__ALLOWED_DIMENSIONS: "384,512,768,1024,1536"
PLUGIN__ENABLE: "true"
PLUGIN__RUNTIME_WS_URL: ws://plugin-runtime:5400/control/ws
PLUGIN__DISPLAY_PLUGIN_DEBUG_URL: wss://cloud.langbot.app/plugin/debug/ws
PLUGIN__WORKER__MAX_CPUS: "0.25"
PLUGIN__WORKER__MAX_MEMORY_MB: "256"
PLUGIN__WORKER__MAX_PIDS: "128"
PLUGIN__WORKER__MAX_WORKERS: "16"
PLUGIN__WORKER__MAX_TOTAL_CPUS: "4.0"
PLUGIN__WORKER__MAX_TOTAL_MEMORY_MB: "4096"
PLUGIN__WORKER__REQUIRE_HARD_LIMITS: "true"
LANGBOT_PLUGIN_RUNTIME_CONTROL_TOKEN: ${PLUGIN_RUNTIME_CONTROL_TOKEN}
# Cloud v2 currently grants no managed Box capability. Keep the shared
# runtime deployed but disable Core integration until a hard-quota-capable
# backend can satisfy the fail-closed Cloud readiness contract.
BOX__ENABLED: "false"
BOX__BACKEND: nsjail
BOX__RUNTIME__ENDPOINT: ws://box:5410
BOX__ADMISSION__REQUIRED: "true"
BOX__ADMISSION__LOGICAL_SESSION_ID: global
BOX__ADMISSION__REQUIRED_BACKEND: nsjail
BOX__ADMISSION__MAX_SESSIONS: "1"
BOX__ADMISSION__MAX_MANAGED_PROCESSES: "0"
BOX__ADMISSION__CPUS: "0.25"
BOX__ADMISSION__MEMORY_MB: "256"
BOX__ADMISSION__WORKSPACE_QUOTA_MB: "256"
BOX__LOCAL__HOST_ROOT: /app/data/box
BOX__LOCAL__DEFAULT_WORKSPACE: /app/data/box
BOX__LOCAL__ALLOWED_MOUNT_ROOTS: /app/data/box
LANGBOT_BOX_CONTROL_TOKEN: ${BOX_CONTROL_TOKEN}
MCP__STDIO__ENABLED: "false"
LANGBOT_SPACE_CONTROL_PLANE_URL: https://space.langbot.app
LANGBOT_SPACE_CONTROL_PLANE_TOKEN: ${CLOUD_V2_CONTROL_PLANE_TOKEN}
LANGBOT_TELEMETRY_INGEST_TOKEN: ${CLOUD_V2_CONTROL_PLANE_TOKEN}
LANGBOT_SPACE_CONTROL_PLANE_PUBLIC_KEY: ${CLOUD_V2_MANIFEST_PUBLIC_KEY}
LANGBOT_SPACE_CONTROL_PLANE_KEY_ID: ${CLOUD_V2_MANIFEST_KEY_ID}
SPACE__URL: https://space.langbot.app
depends_on:
postgres: {condition: service_healthy}
networks: [internal]
plugin-runtime:
image: rockchin/langbot:${LANGBOT_IMAGE_TAG}
container_name: langbot-cloud-plugin-runtime
restart: unless-stopped
command: [uv, run, python, -m, langbot_plugin.cli.__init__, rt]
environment:
LANGBOT_PLUGIN_RUNTIME_CONTROL_TOKEN: ${PLUGIN_RUNTIME_CONTROL_TOKEN}
volumes:
- plugin-data:/app/data
- /sys/fs/cgroup:/sys/fs/cgroup:rw
cgroup: host
privileged: true
expose: ["5400"]
networks: [internal]
box:
image: rockchin/langbot:${LANGBOT_IMAGE_TAG}
container_name: langbot-cloud-box
restart: unless-stopped
command: [uv, run, lbp, box, --host, 0.0.0.0, --ws-control-port, "5410"]
environment:
LANGBOT_BOX_CONTROL_TOKEN: ${BOX_CONTROL_TOKEN}
LANGBOT_BOX_ROOT: /app/data/box
volumes:
- box-data:/app/data/box
- /sys/fs/cgroup:/sys/fs/cgroup:rw
cgroup: host
privileged: true
expose: ["5410"]
networks: [internal]
core:
image: rockchin/langbot-cloud-core:${LANGBOT_IMAGE_TAG}
container_name: langbot-cloud-core
restart: unless-stopped
environment: *core-env
volumes:
- core-data:/app/data
- box-data:/app/data/box
depends_on:
postgres: {condition: service_healthy}
redis: {condition: service_healthy}
plugin-runtime: {condition: service_started}
box: {condition: service_started}
expose: ["5300"]
healthcheck:
test: [CMD-SHELL, "python -c 'import urllib.request; urllib.request.urlopen(\"http://127.0.0.1:5300/healthz\", timeout=3)'" ]
interval: 10s
timeout: 5s
retries: 30
start_period: 30s
networks: [internal, shared-network]
networks:
internal:
shared-network:
external: true
volumes:
postgres-data:
redis-data:
plugin-data:
box-data:
core-data:
+8 -8
View File
@@ -1,5 +1,5 @@
# Docker Compose configuration for LangBot
# For Kubernetes deployment, see kubernetes.yaml and the deployment guide at https://langbot.app/docs
# For Kubernetes deployment, see kubernetes.yaml and the deployment guide at https://docs.langbot.app
version: "3"
services:
@@ -47,10 +47,11 @@ services:
restart: on-failure
environment:
- TZ=Asia/Shanghai
# Optional shared control-plane secret used to authenticate both the RPC
# socket and managed-process relay. Leave unset on both OSS services, or
# generate one with ``openssl rand -hex 32`` and set the same value on
# both ends. Strongly recommended when the deployment is Internet-accessible.
# Shared control-plane secret used to authenticate both the RPC socket
# and managed-process relay. Generate once (for example with
# ``openssl rand -hex 32``) and export it before enabling this profile.
# An empty value is accepted by Compose so Box can remain optional, but
# the Box runtime itself fails closed when the profile is started.
- LANGBOT_BOX_CONTROL_TOKEN=${LANGBOT_BOX_CONTROL_TOKEN:-}
# Box has its own process-wide blocking-work budget.
- LANGBOT_BLOCKING_EXECUTOR_MAX_WORKERS=${LANGBOT_BLOCKING_EXECUTOR_MAX_WORKERS:-8}
@@ -78,9 +79,8 @@ services:
- TZ=Asia/Shanghai
# Optional. Leave unset on both OSS services, or match plugin Runtime.
- LANGBOT_PLUGIN_RUNTIME_CONTROL_TOKEN=${LANGBOT_PLUGIN_RUNTIME_CONTROL_TOKEN:-}
# When set, this must match langbot_box. If both ends leave it unset,
# OSS permits the connection without token authentication. The token is
# sent only in WebSocket handshake headers, never in URLs or payloads.
# Must match the value supplied to langbot_box. The token is sent only
# in WebSocket handshake headers, never in URLs or action payloads.
- LANGBOT_BOX_CONTROL_TOKEN=${LANGBOT_BOX_CONTROL_TOKEN:-}
# Core process-wide blocking-work admission. These are native config
# overrides and are persisted with the effective data/config.yaml.
+1 -1
View File
@@ -1,7 +1,7 @@
# Kubernetes Deployment for LangBot
# This file provides Kubernetes deployment manifests for LangBot based on docker-compose.yaml
#
# Full deployment guide (zh/en/ja): https://langbot.app/docs -> Installation -> Kubernetes
# Full deployment guide (zh/en/ja): https://docs.langbot.app -> Installation -> Kubernetes
#
# Usage:
# kubectl -n langbot create secret generic langbot-plugin-runtime-control \
-17
View File
@@ -88,23 +88,6 @@ Each endpoint accepts **either**:
1. **User Token** (via `Authorization: Bearer <user_jwt_token>`) - for web UI and authenticated users
2. **API Key** (via `X-API-Key` or `Authorization: Bearer <api_key>`) - for external services
### Inspecting API Key Identity
`GET /api/v1/system/context` validates an API key (user JWT not accepted) and returns its bound identity without requiring resource permissions:
```json
{
"code": 0,
"msg": "ok",
"data": {
"instance_uuid": "...",
"workspace_uuid": "...",
"api_key_id": "...",
"permissions": ["..."]
}
}
```
## Example: Model Management
### List All LLM Models
-65
View File
@@ -1,65 +0,0 @@
# ChatGPT / Codex subscription
LangBot's **OpenAI Codex** model provider uses **Sign in with ChatGPT** and the account's Codex entitlement. It is separate from the existing OpenAI API-key provider: subscribing to ChatGPT does not supply an OpenAI Platform API key, and API-key billing is unchanged.
## Connect an account
1. Open **Models**, choose **Add Provider**, and select **OpenAI Codex**.
2. Enter a provider name and choose **Save and sign in**. This saves the provider before authorization, so an interrupted login can be retried from its settings.
3. Open the OpenAI authorization link and enter the one-time code displayed in LangBot. Sign in on OpenAI's site, not in LangBot.
4. If OpenAI asks you to enable device-code authorization, enable it in your ChatGPT account's security settings, or contact your workspace administrator.
5. Keep the LangBot dialog open until it confirms the connection, then finish the form.
6. Use the existing **Scan models** or **Add model** controls, test the model, and select it in a pipeline as usual. Only LLM models are supported by this provider.
The device-code flow also works when LangBot runs remotely or in Docker: the browser does not need to reach a localhost OAuth callback on the server. Serve the LangBot management panel over HTTPS when accessing it remotely.
The account's model catalog is authoritative. A model listed elsewhere or entered manually is not a guarantee that this account has access. Scan errors are reported rather than replaced with a fabricated available-model list.
## Reconnect and disconnect
Open the provider's existing settings to sign in again or disconnect. LangBot refreshes expiring access tokens automatically. A revoked or invalid refresh grant requires another sign-in; transient network failures are not proof that the grant was revoked.
**Disconnect** removes this provider's locally stored authorization. It does not log the account out of other applications or revoke the account globally. Canceling a pending sign-in is separate from disconnecting an existing account. Removing a provider also removes its authorization; the normal rule that models must be removed first still applies.
A saved provider can remain disconnected. Scanning or invoking it then returns a sign-in-required error; LangBot does not silently switch to paid API-key billing.
## Usage and deployment boundary
Calls consume the connected account's included Codex usage and remain subject to OpenAI's plan limits, model availability, workspace policies, and terms. Token counts recorded by LangBot are request usage, not a measurement of remaining subscription quota or an OpenAI invoice.
Use this integration for your own authorized account and trusted workflows. Third-party sign-in support is not permission to pool accounts, resell subscription quota, or redistribute one subscription as a shared API service. For a public or commercial multi-user service, use the appropriate OpenAI API or separately authorized enterprise arrangement. The provider remains a Workspace resource in LangBot: consider who can invoke its models before connecting a personal account.
## Credential handling and API surface
- OAuth credentials are stored server-side separately from provider API keys. Provider and model reads do not supply OAuth access, refresh, or ID tokens.
- Authorization uses a fixed OpenAI origin. The Codex provider does not accept a custom base URL or manually supplied API keys.
- Authentication controls require an authenticated LangBot browser user with `provider_secret.manage` in the selected Workspace. Pending attempts are scoped to the Workspace, provider, and initiating user.
- Browser storage must not contain OAuth tokens. Treat the server database and its backups as sensitive application data.
- MCP and LangBot API keys do not expose the browser-only OAuth controls. Agents may inspect configured providers and models with the existing tools, but a human connects the subscription in the management panel.
The provider-scoped authentication routes are under `/api/v1/provider/providers/{uuid}/codex`:
| Method | Suffix | Purpose |
| --- | --- | --- |
| GET | `/status` | Read local connection state without returning credentials |
| POST | `/device` | Start device authorization |
| POST | `/device/poll` | Poll the initiating user's authorization attempt |
| DELETE | `/device/{authorization_id}` | Cancel only that pending attempt |
| DELETE | `/auth` | Remove local authorization |
Use the returned polling interval and expiration time. An expired attempt must be restarted. These routes are not a general-purpose subscription-to-API gateway.
## References
- [OpenAI Codex authentication](https://developers.openai.com/codex/auth): ChatGPT versus API-key access and device-code login.
- [Hermes Agent providers](https://hermes-agent.nousresearch.com/docs/integrations/providers/): subscription device authentication and refresh recovery.
- [OpenClaw OpenAI provider](https://docs.openclaw.ai/providers/openai): subscription and API-key route distinctions.
- [New API](https://github.com/QuantumNous/new-api): reference for Codex protocol compatibility; its gateway/account-pooling product model is not adopted here.
## 中文快速说明
在「模型」中添加提供商,选择 **OpenAI Codex**,填写名称并点击「保存并登录」。打开 OpenAI 授权页面,输入 LangBot 显示的一次性验证码,完成授权后回到原对话框。随后照常扫描或添加模型、测试模型,并在流水线中选择它。
无需填写 API Key,也无需为远程服务器配置 localhost 回调。登录中断后可以从该提供商的设置中重试;断开连接只删除 LangBot 中保存的授权。调用消耗所登录账号的 Codex 额度,受账号实际权限和 OpenAI 限制约束,不会自动转用按量付费的 OpenAI API。
此功能用于自己的授权账号及可信工作流,不应将个人订阅作为面向多个用户转售或共享的 API 服务。提供商仍是 LangBot 工作空间内的资源,连接个人账号前请确认模型的使用范围。
+2 -2
View File
@@ -218,8 +218,8 @@ metadata:
spec:
categories: [popular, global]
help_links:
zh: https://langbot.app/docs/zh/platforms/http-bot
en: https://langbot.app/docs/en/platforms/http-bot
zh: https://docs.langbot.app/zh/platforms/http-bot
en: https://docs.langbot.app/en/platforms/http-bot
config:
- { name: inbound_secret, type: string, required: true, default: "" }
- { name: callback_url, type: string, required: false, default: "" }
+1 -18
View File
@@ -10,19 +10,6 @@ uvx langbot
This will automatically download and run the latest version of LangBot.
SeekDB support is optional and is not installed by the command above. If you
want to use the SeekDB vector database or the built-in SeekDB embedding model,
run LangBot with the `seekdb` extra:
```bash
uvx --from 'langbot[seekdb]@latest' langbot
```
The extra includes native dependencies whose supported operating systems may
be narrower than LangBot's. In particular, the current Apple Silicon wheels
require macOS 15 or later. The default Chroma backend does not have this
requirement.
## Install with pip/uv
You can also install LangBot as a regular Python package:
@@ -33,10 +20,6 @@ pip install langbot
# Using uv
uv pip install langbot
# Include optional SeekDB support
pip install 'langbot[seekdb]'
# or: uv pip install 'langbot[seekdb]'
```
Then run it:
@@ -118,7 +101,7 @@ uvx langbot
## System Requirements
- Python 3.11 or higher (lower than Python 4)
- Python 3.10.1 or higher
- Operating System: Linux, macOS, or Windows
## Differences from Source Installation
+44 -35
View File
@@ -16,20 +16,12 @@ This document describes how to use OceanBase SeekDB as the vector database backe
## Installation
SeekDB is an optional LangBot feature. A normal LangBot installation uses
Chroma by default and does not install `pyseekdb` or its native bindings.
SeekDB support is automatically included when you install LangBot. The required dependency `pyseekdb` is listed in `pyproject.toml`.
Choose the command that matches how you run LangBot:
If you need to install it manually:
```bash
# PyPI / uvx
uvx --from 'langbot[seekdb]@latest' langbot
# Installed package
pip install 'langbot[seekdb]'
# Source checkout
uv sync --extra seekdb
pip install pyseekdb
```
## ⚠️ Platform Compatibility
@@ -38,36 +30,31 @@ uv sync --extra seekdb
| Platform | Status | Notes |
|----------|--------|-------|
| Linux x86_64 / ARM64 | ✅ Supported | Full embedded mode support via `pylibseekdb` |
| macOS 15+ on Apple Silicon | ✅ Supported | Requires the macOS ARM64 `pylibseekdb` wheel |
| macOS 14 or earlier on Apple Silicon | ❌ Not currently supported | The published native wheel requires macOS 15+; follow [oceanbase/seekdb#1324](https://github.com/oceanbase/seekdb/issues/1324) |
| macOS on Intel | ❌ Not currently supported | No embedded binding is selected by `pyseekdb` |
| Windows | ❌ Not currently supported | No Windows `pylibseekdb` wheel is published |
| Linux | ✅ Supported | Full embedded mode support via `pylibseekdb` |
| macOS | ❌ Not Supported | `pylibseekdb` is Linux-only; use server mode instead |
| Windows | ❌ Not Supported | `pylibseekdb` is Linux-only; use server mode instead |
**Important**: Embedded mode requires a compatible `pylibseekdb` wheel. Do not
force-install or retag a wheel built for a newer macOS release: the bundled
binaries also declare macOS 15 as their minimum deployment target.
**Important**: Embedded mode requires the `pylibseekdb` library, which is only available on Linux. If you're on macOS or Windows, you must use server mode.
### Server Mode (Docker)
| Platform | Status | Notes |
|----------|--------|-------|
| Linux | ✅ Supported | Full Docker support |
| macOS | ✅ Supported by Docker Desktop | The previous slow-disk startup issue was fixed upstream in [oceanbase/seekdb#36](https://github.com/oceanbase/seekdb/issues/36) |
| Windows | ⚠️ Depends on the container runtime | Use a Linux container and follow the upstream image documentation |
| macOS | ⚠️ Known Issue | Docker container initialization failure - [See Issue #36](https://github.com/oceanbase/seekdb/issues/36) |
| Windows | ⚠️ Untested | Should work but not yet tested |
**macOS Users**: Currently, SeekDB Docker containers have an initialization issue on macOS ([oceanbase/seekdb#36](https://github.com/oceanbase/seekdb/issues/36)). Until this is resolved, we recommend:
- Using ChromaDB or Qdrant as alternatives
- Connecting to a remote SeekDB server on Linux if available
### Server Mode (Remote Connection)
| Platform | Status | Notes |
|----------|--------|-------|
| Linux | ✅ Supported | Install the `seekdb` extra and connect to the remote server |
| macOS 15+ on Apple Silicon | ✅ Supported | Install the `seekdb` extra and connect to the remote server |
| macOS 14 or earlier on Apple Silicon | ⚠️ Blocked by upstream packaging | `pyseekdb` currently requires the unavailable native wheel even for server-only use; follow [#1324](https://github.com/oceanbase/seekdb/issues/1324) |
| macOS on Intel / Windows | ✅ Server mode only | Embedded bindings are not available |
| All Platforms | ✅ Supported | Connect to SeekDB running on a remote Linux server |
Remote server mode does not use embedded storage at runtime. However, whether
the Python client can be installed still depends on `pyseekdb`'s package
metadata for the current platform.
**Recommendation for macOS/Windows users**: Deploy SeekDB on a Linux server and connect via server mode configuration.
## Configuration
@@ -183,23 +170,22 @@ Key methods:
### Import Error
If you see: `SeekDB support is not installed`
If you see: `ImportError: pyseekdb is not installed`
Solution:
```bash
uv sync --extra seekdb
# or: uvx --from 'langbot[seekdb]@latest' langbot
pip install pyseekdb
```
### Embedded Mode Is Unavailable on the Current Platform
### Embedded Mode Error on macOS/Windows
**Error**:
```
RuntimeError: Embedded Client is not available because pylibseekdb is not available.
Please install pylibseekdb (Linux only) or use RemoteServerClient (host/port) instead.
```
**Cause**: No compatible `pylibseekdb` wheel is installed for the current OS,
CPU architecture, Python version, and macOS deployment target.
**Cause**: `pylibseekdb` is only available on Linux platforms.
**Solution**: Use server mode instead:
1. Deploy SeekDB on a Linux server or VM
@@ -222,6 +208,29 @@ vdb:
use: chroma # or qdrant
```
### Docker Container Fails on macOS
**Symptoms**:
```bash
docker run -d -p 2881:2881 oceanbase/seekdb:latest
# Container exits immediately with code 30
```
**Error in logs**:
```
[ERROR] Code: Agent.SeekDB.Not.Exists
Message: initialize failed: init agent failed: SeekDB not exists in current directory.
```
**Cause**: This is a known issue with SeekDB Docker containers on macOS. See [oceanbase/seekdb#36](https://github.com/oceanbase/seekdb/issues/36).
**Status**: Under investigation by OceanBase team.
**Workaround Options**:
1. **Use alternatives**: ChromaDB or Qdrant work perfectly on macOS
2. **Remote server**: Deploy SeekDB on a Linux server and connect remotely
3. **Wait for fix**: Monitor the GitHub issue for updates
### Connection Error (Server Mode)
If SeekDB server is not reachable, check:
@@ -243,7 +252,7 @@ For large datasets:
- SeekDB GitHub: https://github.com/oceanbase/seekdb
- pyseekdb SDK: https://github.com/oceanbase/pyseekdb
- OceanBase Documentation: https://oceanbase.ai
- LangBot Documentation: https://langbot.app/docs
- LangBot Documentation: https://docs.langbot.app
## License
Binary file not shown.

Before

Width:  |  Height:  |  Size: 73 KiB

+1 -1
View File
@@ -6,7 +6,7 @@ Minimal, dependency-light clients for the LangBot **HTTP Bot** platform adapter.
They show the whole loop: signing a request, pushing a message, and receiving
multi-part replies on a callback endpoint.
Full guide: [docs.langbot.app — HTTP Bot](https://langbot.app/docs/en/usage/platforms/http-bot).
Full guide: [docs.langbot.app — HTTP Bot](https://docs.langbot.app/en/usage/platforms/http-bot).
Machine-readable contract: [`docs/http-bot-openapi.json`](../../docs/http-bot-openapi.json).
## Files
+1 -1
View File
@@ -6,7 +6,7 @@
它们完整展示了整条链路:对请求签名、推送一条消息、在回调端点接收
1→M 的多段回复。
完整指南:[docs.langbot.app —— HTTP Bot](https://langbot.app/docs/zh/usage/platforms/http-bot)。
完整指南:[docs.langbot.app —— HTTP Bot](https://docs.langbot.app/zh/usage/platforms/http-bot)。
机器可读的接口契约:[`docs/http-bot-openapi.json`](../../docs/http-bot-openapi.json)。
## 文件清单
+1 -1
View File
@@ -6,7 +6,7 @@ A single self-contained HTML page that demos the LangBot **Page Bot**
(`web_page_bot`) embeddable chat widget — the one you drop onto any website with
a single `<script>` tag.
Full guide: [docs.langbot.app — Page Bot](https://langbot.app/docs/en/usage/platforms/webpage).
Full guide: [docs.langbot.app — Page Bot](https://docs.langbot.app/en/usage/platforms/webpage).
## Files
+1 -1
View File
@@ -6,7 +6,7 @@
(`web_page_bot`) 的可嵌入聊天组件 —— 也就是你用一行 `<script>` 标签就能放到任意
网站上的那个组件。
完整指南:[docs.langbot.app —— 页面机器人](https://langbot.app/docs/zh/usage/platforms/webpage)。
完整指南:[docs.langbot.app —— 页面机器人](https://docs.langbot.app/zh/usage/platforms/webpage)。
## 文件清单
-145
View File
@@ -1,145 +0,0 @@
LangBot 用户许可协议
User License Agreement
飞牛 fnOS 平台发行版 | 最后更新:2026 年 8 月
感谢您使用 LangBot!本协议是您(用户)与 LangBot 开源项目(以下简称「LangBot」「我们」)之间,就您在飞牛 fnOS 平台(含飞牛 NAS 设备、飞牛 OS 及其应用中心,以下简称「平台」)上安装、运行、使用 LangBot 应用所订立的合法协议。
您一旦在飞牛应用中心勾选「我接受许可协议的条款」并继续安装、或以其他方式运行 LangBot,即表示您已阅读、理解并同意本协议的全部内容。如您不同意本协议,请不要安装或使用本软件。
---
一、软件性质
1.1 LangBot 是一款基于 LLM 的多平台智能对话机器人开源软件。主程序源代码按照 Apache License, Version 2.0 公开,您可在遵守开源协议的前提下自由使用、修改与再分发。
1.2 本发行版系 LangBot 社区为飞牛 fnOS 平台打包构建的自托管移植版本,与飞牛官方、飞牛硬件厂商不存在从属或关联关系。飞牛应用中心提供的分发渠道不构成对软件功能、可用性的任何担保。
---
二、使用授权
2.1 授予您一份有限的、非排他的、不可转让的个人使用许可:您可在一台或多台您合法拥有或管理的 fnOS 设备上安装、运行本软件,用于个人、家庭或组织内部合法用途。
2.2 您不得:
(a) 将本软件用于违反中国大陆地区法律法规或您所在司法管辖区法律的用途;
(b) 对机器人账号进行骚扰、诈骗、批量营销、发布违法违规内容等滥用行为;
(c) 逆向工程、反编译本软件所包含的第三方二进制(uv 等),但适用法律明确允许或对应开源许可证另作规定的除外;
(d) 试图干扰、过载或损害任何由本软件对外提供的服务或其基础设施。
---
三、服务可用性与「现状」提供
3.1 我们努力提供稳定可靠的软件,但**不保证运行无中断或无错误**。软件可能因维护、升级、不可抗力或其他原因暂时不可用。
3.2 本软件按「现状」「按可用」提供。我们不作任何明示或默示担保,包括但不限于对适销性、特定用途适用性与非侵权性的默示担保。我们不保证:
(a) 软件将满足您的具体需求;
(b) 运行不间断、及时、安全或无错误;
(c) 使用获得的结果准确或可靠;
(d) 任何错误都会被修复。
---
四、责任限制
在适用法律允许的最大范围内:
(a) 我们不对任何**间接、附带、特殊、惩罚性或后果性损失**承担责任,包括但不限于利润损失、数据丢失、商誉损失、业务中断或其他无形损失,无论是否已被告知该等损害发生的可能性;
(b) 无论基于合同、侵权、严格责任或其他任何理论,我们就本软件所引起的所有索赔,向您承担的**累计赔偿总额**不超过 15 美元(或等值当地货币)。
---
五、用户责任
5.1 您对通过本软件发送的所有内容和消息(包括您所配置并运行的机器人发出的消息)承担全部责任。
5.2 您必须遵守所有适用的法律法规,包括但不限于数据保护、隐私、消费者保护和反垃圾邮件法律。
5.3 您有责任保管好自己的账号凭据和各类 API Key,并**自行对重要数据进行备份**。我们对因服务中断、账号终止或系统故障等任何原因造成的数据丢失不承担责任。
5.4 您不得将本软件用于任何非法活动、发送垃圾信息、骚扰他人或侵害他人合法权益。
5.5 您使用机器人接入任何即时通讯平台(QQ、微信、飞书、钉钉、Telegram 等)前,须自行确认已获得该平台授权并遵守其开发者协议与社区规范;因违规接入导致的账号封禁、平台处罚,由您自行承担。
---
六、数据与隐私
6.1 本自托管版本默认情况下,所有配置、对话记录、知识库与插件数据均保存在您所安装的 fnOS 设备本地共享目录(langbot/data)中,不会被自动上传至除您显式配置的模型/服务提供商以外的任何第三方。您对自己的数据及备份负责。
6.2 LangBot 默认启用最少量的匿名遥测,用于帮助改进产品。详细政策见官方文档:
https://docs.langbot.app/zh/insight/data-collection-policy
6.3 启用遥测时仅可能发送:查询事件(适配器类型、运行器类型、模型名称、处理耗时、版本号、匿名工作区 UUID、插件/功能使用计数、不含用户内容的错误追踪)、每日一次的工作区心跳(部署概况、资源对象数量)、完全自愿的问卷回答。
6.4 我们绝对不收集:消息内容、用户名/手机号/平台账号 ID、API 密钥或凭据、IP 地址、文件或媒体内容。
6.5 关闭方式:进入 LangBot Web 管理界面 → 设置 → Space 遥测,关闭开关;或在配置文件 data/config.yaml 中设置 `space.disable_telemetry: true`。关闭后所有功能照常运行。
---
七、不可抗力
因不可抗力事件导致的履约失败或延迟,我们不承担责任。不可抗力包括但不限于:自然灾害(地震、洪水、飓风等)、战争、恐怖主义或内乱、政府行为或法规、网络攻击(DDoS、勒索软件等)、第三方服务故障(AI 模型提供商、即时通讯平台、飞牛平台运行时等)、电力故障或互联网连接中断、流行病或公共卫生紧急事件。
---
八、第三方服务
本软件可能依赖或集成第三方服务,包括但不限于:即时通讯平台(Telegram、Discord、微信、QQ、Slack 等)、AI 模型提供商(OpenAI、Anthropic、Google、深度求索等)、飞牛平台的应用中心运行时与依赖应用(如 Node.js)。我们不对第三方服务的可用性、准确性、可靠性或安全性负责;使用第三方服务须受其各自条款约束,第三方服务的变更可能不经通知即影响本软件功能。
---
九、软件修改与终止
9.1 我们保留随时修改、暂停或终止软件或其中任何部分的权利,无论是否事先通知。
9.2 我们保留随时修改本协议的权利。对重大变更我们会尽力在项目主页公告,继续使用软件即视为接受修订后的协议。
9.3 您可以在飞牛应用中心卸载本软件。卸载时向导会询问是否保留数据,您可自主选择。终止后您继续使用软件的权利立即终止,我们无义务保留您的数据。
---
十、赔偿
您同意赔偿、抗辩并使 LangBot 团队及其关联贡献者、管理人员、代理人、员工免受任何及所有因以下事项引起或与之相关的索赔、损失、损害、责任、成本和费用(包括合理的律师费):
(a) 您对软件的使用;
(b) 您违反本协议;
(c) 您违反任何适用的法律或法规;
(d) 您侵犯任何第三方权利;
(e) 您或您所配置的机器人通过本软件传输的内容。
---
十一、知识产权
11.1 本软件及其设计、代码、文档和品牌归 LangBot 项目所有并受知识产权法保护。
11.2 使用本软件并不授予您对软件的任何所有权。
11.3 您保留通过本软件创建和传输内容的所有权。
---
十二、争议解决
因本协议引起或与之相关的任何争议,应首先通过友好协商解决。协商在 30 日内未达成一致的,任何一方均可依适用法律规定向有管辖权的法院提起诉讼。
---
十三、可分割性
如本协议的任何条款被认定为不可执行或无效,该条款应在最小必要范围内予以限制或剔除,其余条款继续完全有效。
---
十四、完整协议
本协议连同我们的数据收集政策(https://docs.langbot.app/zh/insight/data-collection-policy)构成您与我们之间关于本软件的完整协议,并取代所有先前的协议与谅解。
---
附录:开源许可证声明
LangBot 主程序代码依照 Apache License 2.0 发布。详细条款见:
https://github.com/langbot-app/LangBot/blob/master/LICENSE
或 LangBot 源码包内的 LICENSE 文件。
数据收集政策:https://docs.langbot.app/zh/insight/data-collection-policy
-78
View File
@@ -1,78 +0,0 @@
# LangBot fnOS Packaging
This directory packages LangBot as a `.fpk` app for the fnOS App Store. It is a native deployment: no Docker involved — uv creates a Python virtual environment directly on the NAS, and Node.js v22 from the fnOS App Store provides the Box sandbox and npx MCP capabilities.
## Directory Structure
```
packaging/fnos/
├── manifest # App metadata (appname/version/port/dependency declarations)
├── build.sh # One-shot build script (shared by local and CI)
├── LICENSE
├── config/
│ ├── privilege # Privilege config (run-as: root)
│ └── resource # Persistent data share declaration (langbot/data)
├── cmd/ # Lifecycle scripts (fnOS invokes them with TRIM_* env vars)
│ ├── main # Service start/stop manager (start/stop/status, owns PID/log)
│ ├── install_init # Pre-install hook
│ ├── install_callback # Post-install hook: create venv, uv sync deps, seed config.yaml port
│ ├── upgrade_init # Pre-upgrade hook
│ ├── upgrade_callback # Post-upgrade hook
│ ├── uninstall_init # Pre-uninstall hook
│ ├── uninstall_callback # Post-uninstall hook (keeps data per wizard choice)
│ ├── config_init # Pre-config-change hook
│ └── config_callback # Post-config-change hook (apply new port etc.)
├── wizard/ # Install wizards (JSON forms; values passed as wizard_* env vars)
│ ├── install # On install: Node version, web port, deployment-time notice
│ ├── upgrade # On upgrade: Node version confirmation
│ └── uninstall # On uninstall: whether to keep data
├── app/
│ ├── ui/config # Desktop entry declaration (${wizard_port} placeholder, substituted by fnOS at install)
│ ├── desktop/langbot.main.url
│ ├── langbot/ # [generated] repo source synced via rsync (includes web/dist)
│ └── bin/ # [generated] offline uv binaries (x86_64/aarch64)
├── ICON.PNG / ICON_256.PNG # [generated] derived from res/logo-blue.png
└── langbot.fpk # [generated] final artifact
```
Paths marked `[generated]` are produced by `build.sh`, ignored via `.gitignore`; everything else is a git-tracked source file.
## Building
### Locally
```bash
bash packaging/fnos/build.sh
```
Dependencies: python3 + Pillow, node + npm (or pnpm), fnpack (official fnOS packaging CLI, download from https://developer.fnnas.com/docs/cli/fnpack).
The version comes from `version=` maintained in the manifest; it can also be injected: `FPK_VERSION=4.10.10-1 bash packaging/fnos/build.sh`.
### CI
[`.github/workflows/build-fnos-fpk.yaml`](../../.github/workflows/build-fnos-fpk.yaml) triggers automatically on Release publication and uploads `langbot-<version>-fnos.fpk` to the Release; it can also be triggered manually via workflow_dispatch.
Version sources (consistent with the other release workflows):
| Trigger | Version source |
|---|---|
| Release/tag auto build | Tag name (`v4.10.10-1` → in-package `4.10.10-1`; build.sh strips the `v` prefix on injection) |
| Manual workflow_dispatch / local build | Version maintained in manifest |
## Final Artifact
`langbot.fpk` (gzip + tar archive), containing:
- `manifest` — realigned and appended with a `checksum` field by fnpack
- `app.tgz` — app payload (source, web/dist, uv binaries, entry configs)
- `cmd/`, `config/`, `wizard/` — lifecycle scripts and wizards
- `ICON.PNG`, `ICON_256.PNG`, `LICENSE`
The first startup after installation takes about 5-10 minutes to finish dependency deployment (uv venv + sync); after that the web admin UI is reachable via the desktop icon or `http://<NAS-IP>:<port>` (default port 5300).
## References
- fnOS developer docs: https://developer.fnnas.com/
- fnOS app wizard: https://developer.fnnas.com/docs/core-concepts/wizard/
- fnpack CLI: https://developer.fnnas.com/docs/cli/fnpack
@@ -1,9 +0,0 @@
{
"title": "LangBot",
"icon": "images/icon-256.png",
"type": "url",
"protocol": "http",
"port": "${wizard_port}",
"url": "/",
"allUsers": true
}
-13
View File
@@ -1,13 +0,0 @@
{
".url": {
"langbot.main": {
"title": "LangBot",
"icon": "images/icon-{0}.png",
"type": "url",
"protocol": "http",
"port": "${wizard_port}",
"url": "/",
"allUsers": true
}
}
}
-180
View File
@@ -1,180 +0,0 @@
#!/usr/bin/env bash
# packaging/fnos/build.sh - Build LangBot fnOS FPK package (in-repo version)
# 在 LangBot 仓库内直接打包飞牛 fnOS 应用
# Usage:
# bash packaging/fnos/build.sh # 版本取 manifest 中 version=
# FPK_VERSION=4.10.9 bash packaging/fnos/build.sh # 注入版本(CI 用 release tag
# Dependencies: python3+Pillow, node+npm, fnpack
set -euo pipefail
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
SRC_ROOT="$(cd "${SCRIPT_DIR}/../.." && pwd)" # LangBot 仓库根
FPK_DIR="${SCRIPT_DIR}"
echo "==> LangBot fnOS FPK builder (in-repo)"
echo " Source root: ${SRC_ROOT}"
echo " FPK dir: ${FPK_DIR}"
# --- 0. Inject version (release tag) ---
if [ -n "${FPK_VERSION:-}" ]; then
sed -i "s/^version=.*/version=${FPK_VERSION#v}/" "${FPK_DIR}/manifest"
else
# 无注入版本时自动跟进仓库主版本(pyproject.toml
PY_VER=$(grep -m1 '^version = ' "${SRC_ROOT}/pyproject.toml" | cut -d'"' -f2)
if [ -n "${PY_VER}" ]; then
sed -i "s/^version=.*/version=${PY_VER}/" "${FPK_DIR}/manifest"
fi
fi
MANIFEST_VER=$(grep '^version=' "${FPK_DIR}/manifest" | cut -d= -f2)
echo " FPK version: ${MANIFEST_VER}"
# --- 1. Build frontend ---
echo "[1/5] Building frontend (web/dist)..."
cd "${SRC_ROOT}/web"
if command -v pnpm >/dev/null 2>&1; then
pnpm install --frozen-lockfile 2>/dev/null || pnpm install
pnpm build
else
npm install
npx vite build
fi
[ -d dist ] || { echo "ERROR: web/dist missing" >&2; exit 1; }
echo " Frontend built"
# --- 2. Sync source into packaging/fnos/app/langbot/ ---
echo "[2/5] Syncing source to app/langbot/..."
rm -rf "${FPK_DIR}/app/langbot"
mkdir -p "${FPK_DIR}/app/langbot"
cd "${SRC_ROOT}"
rsync -a \
--exclude='.git' \
--exclude='.venv' \
--exclude='__pycache__' \
--exclude='*.pyc' \
--exclude='web/node_modules' \
--exclude='web/.vite' \
--exclude='tests' \
--exclude='packaging' \
--exclude='.pytest_cache' \
--exclude='.mypy_cache' \
--exclude='.ruff_cache' \
--exclude='data' \
--exclude='*.log' \
--exclude='.dockerignore' \
--exclude='Dockerfile' \
--exclude='docker/' \
--exclude='kubernetes.yaml' \
--exclude='.github/' \
--exclude='docs/' \
--exclude='examples/' \
--exclude='res/' \
./ "${FPK_DIR}/app/langbot/"
[ -d "${FPK_DIR}/app/langbot/web/dist" ] || { echo "ERROR: web/dist missing after rsync!" >&2; exit 1; }
echo " Source synced ($(du -sh "${FPK_DIR}/app/langbot" | cut -f1))"
# --- 2.5 Download bundled uv binaries (offline install on NAS) ---
echo "[2.5/5] Downloading bundled uv binaries..."
UV_VERSION="0.12.9"
mkdir -p "${FPK_DIR}/app/bin"
for arch in x86_64 aarch64; do
out="${FPK_DIR}/app/bin/uv-${arch}"
if [ -x "${out}" ]; then
echo " uv-${arch} already present, skip"
continue
fi
tmp="$(mktemp -d)"
if curl -sSL -o "${tmp}/uv.tar.gz" \
"https://github.com/astral-sh/uv/releases/download/${UV_VERSION}/uv-${arch}-unknown-linux-gnu.tar.gz" \
&& tar xzf "${tmp}/uv.tar.gz" -C "${tmp}" \
&& cp "${tmp}/uv-${arch}-unknown-linux-gnu/uv" "${out}"; then
chmod +x "${out}"
echo " uv-${arch} downloaded (${UV_VERSION})"
else
echo " WARNING: failed to download uv for ${arch}, install will fall back to online install" >&2
fi
rm -rf "${tmp}"
done
# --- 3. Regenerate icons from res/logo-blue.png ---
echo "[3/5] Generating icons from res/logo-blue.png..."
export LOGO_SRC="${SRC_ROOT}/res/logo-blue.png"
export OUT_DIR="${FPK_DIR}"
python3 << 'PYEOF'
from PIL import Image
import os, sys
src = os.environ.get("LOGO_SRC")
out_dir = os.environ.get("OUT_DIR")
if not src or not out_dir:
print("ERROR: LOGO_SRC or OUT_DIR not set", file=sys.stderr)
sys.exit(1)
if not os.path.isfile(src):
print(f"ERROR: logo source not found: {src}", file=sys.stderr)
sys.exit(1)
img = Image.open(src).convert("RGBA")
for size, name in [(64, "ICON.PNG"), (256, "ICON_256.PNG")]:
img.resize((size, size), Image.LANCZOS).save(os.path.join(out_dir, name))
ui_dir = os.path.join(out_dir, "app/ui/images")
os.makedirs(ui_dir, exist_ok=True)
for size in [64, 256]:
img.resize((size, size), Image.LANCZOS).save(os.path.join(ui_dir, f"icon-{size}.png"))
desktop_dir = os.path.join(out_dir, "app/desktop/images")
os.makedirs(desktop_dir, exist_ok=True)
for size in [64, 256]:
img.resize((size, size), Image.LANCZOS).save(os.path.join(desktop_dir, f"icon-{size}.png"))
print(" Icons generated")
PYEOF
# --- 4. Validate structure ---
echo "[4/5] Validating FPK structure..."
ERRORS=0
[ -f "${FPK_DIR}/manifest" ] || { echo " MISSING: manifest"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/config/privilege" ] || { echo " MISSING: config/privilege"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/config/resource" ] || { echo " MISSING: config/resource"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/ICON.PNG" ] || { echo " MISSING: ICON.PNG"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/ICON_256.PNG" ] || { echo " MISSING: ICON_256.PNG"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/app/ui/config" ] || { echo " MISSING: app/ui/config"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/app/ui/images/icon-64.png" ] || { echo " MISSING: app/ui/images/icon-64.png"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/app/ui/images/icon-256.png" ] || { echo " MISSING: app/ui/images/icon-256.png"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/app/desktop/langbot.main.url" ] || { echo " MISSING: app/desktop/langbot.main.url"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/app/desktop/images/icon-64.png" ] || { echo " MISSING: app/desktop/images/icon-64.png"; ERRORS=$((ERRORS+1)); }
[ -f "${FPK_DIR}/app/desktop/images/icon-256.png" ] || { echo " MISSING: app/desktop/images/icon-256.png"; ERRORS=$((ERRORS+1)); }
[ -d "${FPK_DIR}/cmd" ] || { echo " MISSING: cmd/"; ERRORS=$((ERRORS+1)); }
[ -d "${FPK_DIR}/wizard" ] || { echo " MISSING: wizard/"; ERRORS=$((ERRORS+1)); }
for script in "${FPK_DIR}/cmd/"*; do
[ -x "${script}" ] || { echo " NOT EXECUTABLE: cmd/$(basename "$script")"; ERRORS=$((ERRORS+1)); }
done
if [ "${ERRORS}" -gt 0 ]; then
echo "FAILED: ${ERRORS} validation errors" >&2
exit 1
fi
echo " Structure OK"
# --- 5. Build FPK ---
echo "[5/5] Building .fpk..."
if ! command -v fnpack >/dev/null 2>&1; then
echo "ERROR: fnpack not found in PATH." >&2
echo " Download from https://developer.fnnas.com/docs/cli/fnpack" >&2
exit 1
fi
cd "${FPK_DIR}"
fnpack build
FPK_FILE=$(ls -t *.fpk 2>/dev/null | head -1)
if [ -n "${FPK_FILE}" ]; then
echo ""
echo "==> Done! FPK: ${FPK_DIR}/${FPK_FILE}"
echo " Size: $(du -sh "${FPK_FILE}" | cut -f1)"
else
echo "WARNING: fnpack finished but no .fpk found in ${FPK_DIR}" >&2
exit 1
fi
-6
View File
@@ -1,6 +0,0 @@
#!/bin/bash
# cmd/config_callback - post-config hook
# Applies wizard-provided settings (e.g. Node.js version) that cmd/main reads.
# Nothing to persist currently; cmd/main reads wizard_node_version at runtime.
exit 0
-4
View File
@@ -1,4 +0,0 @@
#!/bin/bash
# cmd/config_init - pre-config hook
exit 0
-149
View File
@@ -1,149 +0,0 @@
#!/bin/bash
# cmd/install_callback - post-install hook
# Prefers bundled uv binary, creates Python venv, syncs deps, verifies dist.
# Also validates the user-selected Node.js version is actually installed.
APP_DIR="${TRIM_APPDEST}/langbot"
# --- Resolve data directory ---
# LangBot loads config CWD-relative (data/config.yaml); cmd/main replaces
# APP_DIR/data with a symlink to this persistent dir on every start.
DATA_DIR="${TRIM_DATA_SHARE_PATHS%%:*}"
if [ -z "${DATA_DIR}" ]; then
DATA_DIR="${TRIM_PKGVAR}/data"
fi
cd "${APP_DIR}" || {
echo "App directory missing after install" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
# --- Ensure data directory exists ---
mkdir -p "${DATA_DIR}/plugins" "${DATA_DIR}/box" "${DATA_DIR}/logs" 2>/dev/null || true
# --- Pre-seed config.yaml with the user-selected web port ---
# Data root points at the persistent share (LANGBOT_DATA_ROOT is exported by
# cmd/main at start), so this config survives app upgrades.
# LangBot copies templates/config.yaml there on first boot only if missing,
# so we seed it ourselves with the chosen port.
PORT="${wizard_port:-5300}"
case "${PORT}" in
''|*[!0-9]*) PORT="5300" ;;
esac
TEMPLATE_FILE="${APP_DIR}/src/langbot/templates/config.yaml"
_patch_config() {
local cfg_dir="$1"
local cfg_file="${cfg_dir}/config.yaml"
mkdir -p "${cfg_dir}" 2>/dev/null || true
if [ ! -f "${cfg_file}" ] && [ -f "${TEMPLATE_FILE}" ]; then
cp "${TEMPLATE_FILE}" "${cfg_file}"
fi
if [ -f "${cfg_file}" ]; then
sed -i -E "/^api:/,/^[a-z_]+:/ s/^([[:space:]]*port:).*/\1 ${PORT}/" "${cfg_file}"
sed -i "s#webhook_prefix: 'http://127\.0\.0\.1:[0-9]*'#webhook_prefix: 'http://127.0.0.1:${PORT}'#" "${cfg_file}"
fi
}
# Seed the persistent dir AND the CWD-relative APP_DIR/data (cmd/main merges
# the latter into the persistent dir via symlink on first start, so the port
# survives regardless of which path ends up being read).
_patch_config "${DATA_DIR}"
if [ "${APP_DIR}/data" != "${DATA_DIR}" ]; then
_patch_config "${APP_DIR}/data"
fi
# --- Fallback patch for desktop entry port ---
# fnOS natively substitutes ${wizard_port} in ui/config at install time.
# This only kicks in if the placeholder somehow survived (e.g. CLI install).
UI_CONFIG="${TRIM_APPDEST}/ui/config"
if [ -f "${UI_CONFIG}" ] && grep -q 'wizard_port\|{port}' "${UI_CONFIG}"; then
sed -i "s/\${wizard_port}/${PORT}/g; s/{port}/${PORT}/g" "${UI_CONFIG}"
fi
# --- Validate user-selected Node.js version is installed ---
NODE_VERSION="${wizard_node_version:-22}"
if [ ! -d "/var/apps/nodejs_v${NODE_VERSION}" ]; then
echo "Node.js v${NODE_VERSION} 未安装:请先在应用中心安装 nodejs_v${NODE_VERSION},再重新安装本应用。" > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
# --- Python check ---
PYTHON_BIN="python3"
if ! command -v "${PYTHON_BIN}" >/dev/null 2>&1; then
PYTHON_BIN="python"
fi
if ! command -v "${PYTHON_BIN}" >/dev/null 2>&1; then
echo "Python not found on this system" > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
PY_VER=$("${PYTHON_BIN}" -c 'import sys; print(f"{sys.version_info.major}.{sys.version_info.minor}")' 2>/dev/null)
if [ -z "${PY_VER}" ]; then
echo "Python 3.11+ required but not found on this system" > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
PY_MAJOR=$(echo "${PY_VER}" | cut -d. -f1)
PY_MINOR=$(echo "${PY_VER}" | cut -d. -f2)
if [ "${PY_MAJOR}" -lt 3 ] || { [ "${PY_MAJOR}" -eq 3 ] && [ "${PY_MINOR}" -lt 11 ]; }; then
echo "Python 3.11+ required, found ${PY_VER}" > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
# --- Resolve uv: bundled binary first, then online fallbacks ---
UV_BIN=""
ARCH=$(uname -m)
case "${ARCH}" in
x86_64) BUNDLED_UV="${TRIM_APPDEST}/bin/uv-x86_64" ;;
aarch64) BUNDLED_UV="${TRIM_APPDEST}/bin/uv-aarch64" ;;
*) BUNDLED_UV="" ;;
esac
if [ -n "${BUNDLED_UV}" ] && [ -x "${BUNDLED_UV}" ]; then
mkdir -p "${TRIM_PKGVAR}/bin"
cp "${BUNDLED_UV}" "${TRIM_PKGVAR}/bin/uv" && chmod +x "${TRIM_PKGVAR}/bin/uv"
UV_BIN="${TRIM_PKGVAR}/bin/uv"
fi
if [ -z "${UV_BIN}" ] && command -v uv >/dev/null 2>&1; then
UV_BIN="uv"
fi
if [ -z "${UV_BIN}" ]; then
"${PYTHON_BIN}" -m pip install --user --no-cache-dir uv 2>/dev/null || \
"${PYTHON_BIN}" -m pip install --no-cache-dir uv 2>/dev/null || \
curl -LsSf https://astral.sh/uv/install.sh | sh 2>/dev/null || true
export PATH="${HOME}/.local/bin:${PATH}"
if command -v uv >/dev/null 2>&1; then
UV_BIN="uv"
elif [ -x "${HOME}/.local/bin/uv" ]; then
UV_BIN="${HOME}/.local/bin/uv"
fi
fi
if [ -z "${UV_BIN}" ]; then
echo "无法获取 uv:内置二进制缺失且在线安装失败。请检查网络后重新安装。" > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
# --- Create venv via uv ---
if [ ! -d ".venv" ]; then
"${UV_BIN}" venv .venv --python "${PYTHON_BIN}" || {
echo "Failed to create Python virtual environment via uv" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
fi
# --- Sync dependencies ---
"${UV_BIN}" sync --extra seekdb || {
echo "Dependency sync failed. Check network connectivity." > "${TRIM_TEMP_LOGFILE}"
exit 1
}
# --- Verify frontend dist ---
if [ ! -d "web/dist" ] || [ -z "$(ls -A web/dist 2>/dev/null)" ]; then
echo "Frontend dist missing! Web UI will not be available." > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
exit 0
-5
View File
@@ -1,5 +0,0 @@
#!/bin/bash
# cmd/install_init - pre-install hook
# Nothing special to do before extraction.
exit 0
-194
View File
@@ -1,194 +0,0 @@
#!/bin/bash
# cmd/main - LangBot lifecycle manager for fnOS
# Handles start / stop / status via standalone-runtime Python process.
# Node.js path is injected into PATH so Box sandbox npx MCP servers can run.
PID_FILE="${TRIM_PKGVAR}/langbot.pid"
APP_DIR="${TRIM_APPDEST}/langbot"
LOG_FILE="${TRIM_PKGVAR}/langbot.log"
# --- Locate fnOS Node.js bin path ---
# fnOS appname-based path: /var/apps/nodejs_vXX/target/bin
# This is a stable symlink regardless of which volume the app is on.
NODE_VERSION="${wizard_node_version:-22}"
NODE_BIN_DIR="/var/apps/nodejs_v${NODE_VERSION}/target/bin"
if [ -d "${NODE_BIN_DIR}" ]; then
export PATH="${NODE_BIN_DIR}:${PATH}"
fi
# --- Persistent data root ---
# LangBot loads data/config.yaml CWD-RELATIVE (see core/stages/load_config.py:
# load_yaml_config('data/config.yaml', ...)) and resolves its data root to
# <CWD>/data in source-install mode — it does NOT honour LANGBOT_DATA_ROOT for
# config.yaml. So the real fix is the symlink below: APP_DIR/data -> DATA_DIR.
DATA_DIR="${TRIM_DATA_SHARE_PATHS%%:*}"
if [ -z "${DATA_DIR}" ]; then
DATA_DIR="${TRIM_PKGVAR}/data"
fi
export LANGBOT_DATA_ROOT="${DATA_DIR}"
mkdir -p "${DATA_DIR}" 2>/dev/null || true
# --- Locate Python ---
PYTHON_BIN="python3"
! command -v "${PYTHON_BIN}" >/dev/null 2>&1 && PYTHON_BIN="python"
# --- Locate uv ---
# install_callback puts the bundled uv binary at ${TRIM_PKGVAR}/bin/uv
UV_BIN="${TRIM_PKGVAR}/bin/uv"
if [ ! -x "${UV_BIN}" ]; then
UV_BIN="uv"
fi
if ! command -v "${UV_BIN}" >/dev/null 2>&1; then
UV_BIN="${HOME}/.local/bin/uv"
fi
if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then
UV_BIN="${HOME}/.cargo/bin/uv"
fi
if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then
UV_BIN="${APP_DIR}/.venv/bin/uv"
fi
case $1 in
start)
if [ -f "${PID_FILE}" ]; then
PID=$(cat "${PID_FILE}" | tr -d '[:space:]')
if [ -n "${PID}" ] && kill -0 "${PID}" 2>/dev/null; then
exit 0
fi
rm -f "${PID_FILE}"
fi
if [ ! -d "${APP_DIR}" ]; then
echo "LangBot app directory missing: ${APP_DIR}" > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
cd "${APP_DIR}" || {
echo "Cannot enter app directory: ${APP_DIR}" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
if [ ! -d ".venv" ]; then
echo "Python virtual environment not found. Please reinstall LangBot." > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
# --- Unify data location: APP_DIR/data -> symlink to persistent DATA_DIR ---
# LangBot reads config CWD-relative (data/config.yaml), so without this a
# fresh data/ with a default 5300 config gets recreated inside target/ on
# every install/upgrade. The symlink keeps everything on the persistent
# share; it is recreated here on each start (upgrades wipe target/).
APP_DATA="${APP_DIR}/data"
if [ -L "${APP_DATA}" ]; then
# already a symlink; re-point if the persistent dir changed
[ "$(readlink "${APP_DATA}")" != "${DATA_DIR}" ] && ln -sfn "${DATA_DIR}" "${APP_DATA}"
elif [ -d "${APP_DATA}" ]; then
# legacy real dir (created by LangBot before this fix): merge into the
# persistent dir without overwriting newer files already there
mkdir -p "${DATA_DIR}"
cp -an "${APP_DATA}/." "${DATA_DIR}/" 2>/dev/null || cp -a "${APP_DATA}/." "${DATA_DIR}/"
rm -rf "${APP_DATA}"
ln -s "${DATA_DIR}" "${APP_DATA}"
else
ln -s "${DATA_DIR}" "${APP_DATA}"
fi
# Ensure LangBot's actual listen port always matches what fnOS shows in
# "应用设置 → 访问端口" (which is the single source of truth from the user's
# POV). Two sources, checked in priority order:
# 1. ${wizard_port} — only set during install/upgrade callbacks (not on
# normal `start`; kept for completeness).
# 2. target/ui/config — read the "port" field written by fnOS after
# ${wizard_port} substitution AND any later edit the user made via
# "应用设置 → 自定义 URL" pencil button.
# Without this: user picks 5303 in wizard, LangBot still listens on its
# default 5300, desktop shortcut hits 5303 → connection refused.
CONFIG_FILE="${DATA_DIR}/config.yaml"
_patch_port() {
local _port="$1"
case "${_port}" in
''|*[!0-9]*) return ;;
esac
if [ ! -f "${CONFIG_FILE}" ]; then
local _tmpl="${APP_DIR}/src/langbot/templates/config.yaml"
[ -f "${_tmpl}" ] && cp "${_tmpl}" "${CONFIG_FILE}"
fi
if [ -f "${CONFIG_FILE}" ]; then
sed -i -E "/^api:/,/^[a-z_]+:/ s/^([[:space:]]*port:).*/\1 ${_port}/" "${CONFIG_FILE}"
sed -i "s#webhook_prefix: 'http://127\.0\.0\.1:[0-9]*'#webhook_prefix: 'http://127.0.0.1:${_port}'#" "${CONFIG_FILE}"
fi
}
if [ -n "${wizard_port:-}" ]; then
_patch_port "${wizard_port}"
fi
if [ -f "${TRIM_APPDEST}/ui/config" ]; then
_port_from_ui=$(python3 -c 'import json,sys
try:
d = json.load(open(sys.argv[1]))
for _name, _entry in (d.get(".url") or {}).items():
p = _entry.get("port")
if isinstance(p, (int, float)):
print(int(p))
elif isinstance(p, str) and p.isdigit():
print(int(p))
break
except Exception:
pass
' "${TRIM_APPDEST}/ui/config" 2>/dev/null)
if [ -n "${_port_from_ui}" ]; then
_patch_port "${_port_from_ui}"
fi
fi
# Native deployment: no --standalone-runtime flag, LangBot spawns the
# plugin runtime as a stdio subprocess (same as official `uv run main.py`).
# (--standalone-runtime would require an external runtime at
# ws://langbot_plugin_runtime:5400, which only exists in Docker Compose.)
# --standalone-box omitted: Box sandbox defaults off, users enable via Web UI
nohup "${UV_BIN}" run --no-sync main.py \
> "${LOG_FILE}" 2>&1 &
echo $! > "${PID_FILE}"
sleep 3
if [ -f "${PID_FILE}" ]; then
PID=$(cat "${PID_FILE}" | tr -d '[:space:]')
if [ -n "${PID}" ] && kill -0 "${PID}" 2>/dev/null; then
exit 0
fi
fi
echo "LangBot failed to start. Check ${LOG_FILE}" > "${TRIM_TEMP_LOGFILE}"
exit 1
;;
stop)
if [ -f "${PID_FILE}" ]; then
PID=$(cat "${PID_FILE}" | tr -d '[:space:]')
if [ -n "${PID}" ]; then
kill "${PID}" 2>/dev/null
for _ in 1 2 3 4 5 6 7 8 9 10; do
kill -0 "${PID}" 2>/dev/null || break
sleep 1
done
kill -9 "${PID}" 2>/dev/null
fi
rm -f "${PID_FILE}"
fi
exit 0
;;
status)
if [ -f "${PID_FILE}" ]; then
PID=$(cat "${PID_FILE}" | tr -d '[:space:]')
if [ -n "${PID}" ] && kill -0 "${PID}" 2>/dev/null; then
exit 0
fi
rm -f "${PID_FILE}"
fi
exit 3
;;
*)
exit 1
;;
esac
-21
View File
@@ -1,21 +0,0 @@
#!/bin/bash
# cmd/uninstall_callback - post-uninstall hook
# The system preserves var/ and shares/ by default. Honor the user's
# wizard_keep_data choice: delete data only when explicitly requested.
if [ "${wizard_keep_data:-yes}" = "no" ]; then
# 应用运行数据(pid、日志等)
if [ -n "${TRIM_PKGVAR}" ]; then
rm -rf "${TRIM_PKGVAR:?}"/langbot.pid \
"${TRIM_PKGVAR:?}"/langbot.log \
"${TRIM_PKGVAR:?}"/bin 2>/dev/null || true
fi
# 共享数据目录(langbot/data
DATA_DIR="${TRIM_DATA_SHARE_PATHS%%:*}"
if [ -n "${DATA_DIR}" ]; then
rm -rf "${DATA_DIR:?}" 2>/dev/null || true
fi
fi
exit 0
-20
View File
@@ -1,20 +0,0 @@
#!/bin/bash
# cmd/uninstall_init - pre-uninstall hook
# Stop LangBot before files are removed.
PID_FILE="${TRIM_PKGVAR}/langbot.pid"
if [ -f "${PID_FILE}" ]; then
PID=$(cat "${PID_FILE}" | tr -d '[:space:]')
if [ -n "${PID}" ] && kill -0 "${PID}" 2>/dev/null; then
kill "${PID}" 2>/dev/null
for _ in 1 2 3 4 5 6 7 8 9 10; do
kill -0 "${PID}" 2>/dev/null || break
sleep 1
done
kill -9 "${PID}" 2>/dev/null
fi
rm -f "${PID_FILE}"
fi
exit 0
-77
View File
@@ -1,77 +0,0 @@
#!/bin/bash
# cmd/upgrade_callback - post-upgrade hook
# Re-sync dependencies after code replacement using uv.
APP_DIR="${TRIM_APPDEST}/langbot"
# Persistent data root (must match cmd/main)
DATA_DIR="${TRIM_DATA_SHARE_PATHS%%:*}"
[ -z "${DATA_DIR}" ] && DATA_DIR="${TRIM_PKGVAR}/data"
# Apply port from upgrade wizard (config persists across upgrades; this
# only rewrites it when the user changed the value in the upgrade wizard)
CONFIG_FILE="${DATA_DIR}/config.yaml"
if [ -n "${wizard_port:-}" ] && [ -f "${CONFIG_FILE}" ]; then
case "${wizard_port}" in
''|*[!0-9]*) ;;
*)
sed -i -E "/^api:/,/^[a-z_]+:/ s/^([[:space:]]*port:).*/\1 ${wizard_port}/" "${CONFIG_FILE}"
sed -i "s#webhook_prefix: 'http://127\.0\.0\.1:[0-9]*'#webhook_prefix: 'http://127.0.0.1:${wizard_port}'#" "${CONFIG_FILE}"
;;
esac
fi
cd "${APP_DIR}" || {
echo "App directory missing after upgrade" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
# Find uv (bundled first, then PATH / ~/.local/bin / ~/.cargo/bin)
UV_BIN="${TRIM_PKGVAR}/bin/uv"
if [ ! -x "${UV_BIN}" ]; then
UV_BIN="uv"
fi
if ! command -v "${UV_BIN}" >/dev/null 2>&1; then
UV_BIN="${HOME}/.local/bin/uv"
fi
if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then
UV_BIN="${HOME}/.cargo/bin/uv"
fi
PYTHON_BIN="python3"
! command -v "${PYTHON_BIN}" >/dev/null 2>&1 && PYTHON_BIN="python"
# Re-sync deps
if [ -d ".venv" ]; then
"${UV_BIN}" sync --extra seekdb 2>/dev/null || {
echo "Dependency sync failed after upgrade" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
else
# Venv was lost, recreate via uv
if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then
"${PYTHON_BIN}" -m pip install --user --no-cache-dir uv 2>/dev/null || \
"${PYTHON_BIN}" -m pip install --no-cache-dir uv 2>/dev/null || {
echo "Failed to install uv" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
export PATH="${HOME}/.local/bin:${PATH}"
UV_BIN="uv"
fi
"${UV_BIN}" venv .venv --python "${PYTHON_BIN}" || {
echo "Failed to recreate virtual environment" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
"${UV_BIN}" sync --extra seekdb || {
echo "Dependency sync failed" > "${TRIM_TEMP_LOGFILE}"
exit 1
}
fi
# Verify frontend dist still present
if [ ! -d "web/dist" ] || [ -z "$(ls -A web/dist 2>/dev/null)" ]; then
echo "Frontend dist missing after upgrade! Web UI will not be available." > "${TRIM_TEMP_LOGFILE}"
exit 1
fi
exit 0
-20
View File
@@ -1,20 +0,0 @@
#!/bin/bash
# cmd/upgrade_init - pre-upgrade hook
# Stop the running LangBot process before files are replaced.
PID_FILE="${TRIM_PKGVAR}/langbot.pid"
if [ -f "${PID_FILE}" ]; then
PID=$(cat "${PID_FILE}" | tr -d '[:space:]')
if [ -n "${PID}" ] && kill -0 "${PID}" 2>/dev/null; then
kill "${PID}" 2>/dev/null
for _ in 1 2 3 4 5 6 7 8 9 10; do
kill -0 "${PID}" 2>/dev/null || break
sleep 1
done
kill -9 "${PID}" 2>/dev/null
fi
rm -f "${PID_FILE}"
fi
exit 0
-5
View File
@@ -1,5 +0,0 @@
{
"defaults": {
"run-as": "root"
}
}
-9
View File
@@ -1,9 +0,0 @@
{
"data-share": {
"shares": [
{
"name": "langbot/data"
}
]
}
}
-16
View File
@@ -1,16 +0,0 @@
appname=langbot
version=4.10.10
display_name=LangBot
desc=基于 LLM 的多平台智能对话机器人,支持 QQ、微信、飞书、钉钉、Telegram 等十余种即时通讯平台,内置 Web 管理界面和 AI Agent 能力。
platform=all
source=thirdparty
maintainer=LangBot
maintainer_url=https://langbot.app
service_port=5300
checkport=true
os_min_version=0.9.0
desktop_uidir=ui
desktop_applaunchname=langbot.main
ctl_stop=true
install_dep_apps=nodejs_v22
changelog=飞牛 fnOS 增强版首版:原生 Python 部署,依赖 Node.js v22 以启用 Box 沙箱 + npx MCP 能力
-47
View File
@@ -1,47 +0,0 @@
[
{
"stepTitle": "运行环境",
"items": [
{
"type": "tips",
"helpText": "LangBot 依赖 Node.js v22 运行 Box 沙箱和 npx MCP。请先在应用中心安装 Node.js v22。"
},
{
"type": "select",
"field": "wizard_node_version",
"label": "Node.js 版本",
"initValue": "22",
"options": [
{ "label": "Node.js v22 (推荐, LTS)", "value": "22" },
{ "label": "Node.js v24", "value": "24" },
{ "label": "Node.js v20", "value": "20" }
]
}
]
},
{
"stepTitle": "访问配置",
"items": [
{
"type": "text",
"field": "wizard_port",
"label": "Web 访问端口",
"initValue": "5300",
"rules": [
{ "required": true, "message": "请输入访问端口" },
{ "pattern": "^[0-9]+$", "message": "端口只能是数字" },
{ "min": 1, "max": 5, "message": "端口号长度不正确" }
]
}
]
},
{
"stepTitle": "安装说明",
"items": [
{
"type": "tips",
"helpText": "安装完成后,LangBot 首次启动约需 5-10 分钟完成依赖部署与初始化,部署完成后即可打开网页端使用。"
}
]
}
]
-17
View File
@@ -1,17 +0,0 @@
[
{
"stepTitle": "数据保留",
"items": [
{
"type": "radio",
"field": "wizard_keep_data",
"label": "是否保留 LangBot 数据(插件、配置、日志)",
"initValue": "yes",
"options": [
{ "label": "保留数据(重新安装后可继续使用)", "value": "yes" },
{ "label": "彻底删除全部数据", "value": "no" }
]
}
]
}
]
-38
View File
@@ -1,38 +0,0 @@
[
{
"stepTitle": "运行环境",
"items": [
{
"type": "tips",
"helpText": "LangBot 依赖 Node.js v22 运行 Box 沙箱和 npx MCP。如需更换版本,请先在应用中心安装对应版本。"
},
{
"type": "select",
"field": "wizard_node_version",
"label": "Node.js 版本",
"initValue": "22",
"options": [
{ "label": "Node.js v22 (推荐, LTS)", "value": "22" },
{ "label": "Node.js v24", "value": "24" },
{ "label": "Node.js v20", "value": "20" }
]
}
]
},
{
"stepTitle": "访问配置",
"items": [
{
"type": "text",
"field": "wizard_port",
"label": "Web 访问端口",
"initValue": "5300",
"rules": [
{ "required": true, "message": "请输入访问端口" },
{ "pattern": "^[0-9]+$", "message": "端口只能是数字" },
{ "min": 1, "max": 5, "message": "端口号长度不正确" }
]
}
]
}
]
+5 -10
View File
@@ -1,6 +1,6 @@
[project]
name = "langbot"
version = "4.10.11"
version = "4.10.7"
description = "Production-grade platform for building agentic IM bots"
readme = "README.md"
license-files = ["LICENSE"]
@@ -22,6 +22,7 @@ dependencies = [
"discord-py>=2.5.2",
"pynacl>=1.5.0", # Required for Discord voice support
"gewechat-client>=0.1.5",
"itchat-uos>=1.5.0.dev",
"lark-oapi>=1.5.5",
"mcp>=1.25.0,<2.0.0",
"nakuru-project-idk>=0.0.2.1",
@@ -70,7 +71,8 @@ dependencies = [
"langchain-text-splitters>=1.1.2",
"chromadb>=1.0.0,<2.0.0",
"qdrant-client (>=1.15.1,<2.0.0)",
"langbot-plugin==0.5.8",
"pyseekdb==1.1.0.post3",
"langbot-plugin @ git+https://github.com/langbot-app/langbot-plugin-sdk.git@7b559da430a50f80a7d30c9d3d66f088503ddbb3",
"asyncpg>=0.30.0",
"line-bot-sdk>=3.19.0",
"matrix-nio>=0.25.2",
@@ -81,7 +83,6 @@ dependencies = [
"botocore>=1.42.39",
"litellm>=1.0.0",
"valkey-glide>=2.4.1,<3.0.0; sys_platform != 'win32'", # No Windows wheels are published
"webauthn>=3.0.0",
]
keywords = [
"bot",
@@ -108,15 +109,9 @@ classifiers = [
"Topic :: Communications :: Chat",
]
[project.optional-dependencies]
seekdb = [
"pyseekdb==1.4.0.post1",
"pylibseekdb==1.4.0; sys_platform == 'linux' or (sys_platform == 'darwin' and platform_machine == 'arm64')",
]
[project.urls]
Homepage = "https://langbot.app"
Documentation = "https://langbot.app/docs"
Documentation = "https://docs.langbot.app"
Repository = "https://github.com/langbot-app/LangBot"
[project.scripts]
+1 -2
View File
@@ -1349,8 +1349,7 @@
"local-agent",
"tools",
"e2b",
"nsjail",
"host"
"nsjail"
],
"automation": "",
"setup_automation": [],
+2 -8
View File
@@ -48,7 +48,7 @@ tools, skill add/edit, and stdio MCP are disabled. Set `box.enabled: false`
## Kubernetes
See `docker/kubernetes.yaml` and the deployment guide at
https://langbot.app/docs. `docker/deploy-k8s-test.sh` is a test helper.
https://docs.langbot.app. `docker/deploy-k8s-test.sh` is a test helper.
## config.yaml (generated at `data/config.yaml` on first run)
@@ -63,7 +63,7 @@ Key settings:
| `api.global_api_key` | **Global API key** for the HTTP API + MCP server. Non-empty = accepted with no login/DB record; no `lbk_` prefix required. Empty = disabled. Plaintext — trusted/internal only, serve over HTTPS. |
| `plugin.runtime_ws_url` | Standalone plugin runtime WS URL (e.g. `ws://langbot_plugin_runtime:5400/control/ws`) |
| `box.enabled` | Master switch for the Box sandbox runtime |
| `box.backend` | `local` (Docker/nsjail autopick) / `docker` / `nsjail` / `e2b` / explicit unsafe `host`; env override `BOX__BACKEND` |
| `box.backend` | `local` (Docker/nsjail autopick) / `docker` / `nsjail` / `e2b`; env override `BOX__BACKEND` |
| `box.runtime.endpoint` | External Box runtime URL (e.g. `ws://127.0.0.1:5410`); empty = local auto-managed |
Many keys have `ENV__SUBKEY` overrides (e.g. `BOX__BACKEND`, `BOX__ENABLED`).
@@ -75,10 +75,6 @@ Many keys have `ENV__SUBKEY` overrides (e.g. `BOX__BACKEND`, `BOX__ENABLED`).
with `--standalone-runtime`.
- Box has a parallel `--standalone-box` flag; the Docker box host is
`langbot_box:5410`.
- `box.backend: host` runs commands directly as the Box Runtime system user.
It is never auto-selected, provides no sandbox isolation, and is only for
trusted local development. A WebSocket-controlled host backend requires
`LANGBOT_BOX_CONTROL_TOKEN`; local stdio control is allowed.
## Global API key — enabling for agents/automation
@@ -97,7 +93,5 @@ login session. See `langbot-mcp-ops` for using it, and `docs/API_KEY_AUTH.md`.
- "No supported sandbox backend (Docker / nsjail / E2B)" with Docker running
usually means the user isn't in the `docker` group →
`sudo usermod -aG docker <user>` and restart in a new shell.
- Do not use `box.backend: host` as a production fallback. It cannot enforce
image, filesystem, network, PID, CPU, memory, or storage isolation.
- Box root host/container path mismatch breaks sandbox container creation.
- Don't commit a non-empty `api.global_api_key` to version control.
-36
View File
@@ -43,8 +43,6 @@ Two kinds of key are accepted:
Invalid, revoked, or expired keys get `401 Unauthorized`. A valid key whose
scopes do not authorize a tool gets `403 Forbidden`.
To inspect key identity and permissions, call `GET /api/v1/system/context` with the API key.
## Client configuration
```json
@@ -77,8 +75,6 @@ shape as the corresponding HTTP API request body. Discover resources with the
`list_*` / `get_*` tools before mutating; identifiers are UUIDs. Reads require
`resource.view`; mutations require `resource.manage`. All service calls inherit
the immutable Workspace context authenticated at the MCP transport boundary.
Pass `is_default: true` to `create_pipeline` only when the Workspace does not
already have a default pipeline.
## How to use
@@ -88,38 +84,6 @@ already have a default pipeline.
4. Use `list_*` tools to discover, then `get_*` / `create_*` / `update_*` /
`delete_*` as needed.
## ChatGPT / Codex subscription providers
`list_model_providers` can return the `openai-codex` requester. Its OAuth
credentials are server-only and are not provider API keys. Never ask a user
to paste ChatGPT access tokens, refresh tokens, or a Codex auth cache into an
MCP tool or model configuration.
A human connects or disconnects the subscription through **Models → provider
settings** in the LangBot web UI. The provider-scoped `/codex/*` authentication
routes deliberately require a browser-user session and are not exposed as MCP
tools or authorized by a LangBot API key. Once connected, models are managed
and selected through the normal provider/model workflow. A disconnected
provider must be reauthorized; do not silently replace it with API-key billing.
See [ChatGPT / Codex subscription](../../../docs/CODEX_SUBSCRIPTION.md) for setup,
usage limits, and the personal-account versus shared-service boundary.
## Provider deletion
The curated MCP surface currently lists providers but has no provider-deletion
tool. In the web UI, **Edit Provider → Delete** asks for confirmation before
removing that provider and all its LLM, embedding, and rerank models. This is
irreversible; never interpret a request to edit a provider as authorization to
delete it.
The equivalent HTTP operation is
`DELETE /api/v1/provider/providers/{uuid}?cascade=true`, requiring
`resource.manage` in the authenticated Workspace. Omitting `cascade` preserves
the existing refusal to delete providers that still have models. Cloud-managed
providers remain protected. Cascade deletion removes stored Codex authorization
state as well; it is not the same operation as disconnecting an account.
## Implementation & maintenance (for LangBot developers)
- Server: `src/langbot/pkg/api/mcp/server.py` (FastMCP). Tools call the service
+1 -34
View File
@@ -25,10 +25,7 @@ CLI uses. Create one in your Space account (Profile → Personal Access Tokens),
then send it as a Bearer token:
```
Authorization: Bearer <your-pat>
```
Requests without a valid PAT get `401 Unauthorized`.
Authorization: Bearer lbpat_...uests without a valid PAT get `401 Unauthorized`.
## Client configuration
@@ -69,36 +66,6 @@ All tools are read-only.
state (available, unprobed, unavailable), then Space recommendation. Each
item includes `availability.up`, `last_probed_at`, latency, and HTTP status.
## Runner usage recommendations
Use `search_plugins` with `runner_usage: "agent"` for Agent, pipeline, and
setup-wizard recommendations, or `runner_usage: "event"` for event processors.
The component kind remains `Runner`. Only these two exact values are accepted;
omit the optional field to preserve unfiltered browsing.
```json
{"query":"", "runner_usage":"agent", "page":1, "page_size":100}
```
Plugin results include `latest_version` and `runner_usages: string[]`, the
explicit union of usages in that latest installable version. Only recommend a
plugin when this array explicitly contains the target usage. Missing, empty,
malformed, or unknown usages must never mean agent-compatible. Event-only
plugins must never enter Agent recommendations. Empty filtered results are
valid while legacy packages await corrected releases; never remove the filter
to fill a recommendation list.
REST callers use `runner_usage` on both
`POST /api/v1/marketplace/extensions/search` and the compatibility
`POST /api/v1/marketplace/plugins/search`; preserve it during fallback. Add
`"type_filter":"plugin", "component_filter":"Runner"` on the unified endpoint.
Usage is ANDed with other filters before pagination and `total`; MCP/Skill items
do not match. Invalid REST values return HTTP 400.
Open the same filter in the webpage:
`https://space.langbot.app/market?type=plugin&component=Runner&runner_usage=agent`
(or `runner_usage=event`). Switch All / Agent / Event in the Runner usage row.
## Implementation & maintenance (for Space developers)
- Server: `internal/controller/mcp/server.go` (official Go MCP SDK
@@ -13,7 +13,6 @@ tags:
- tools
- e2b
- nsjail
- host
skills:
- langbot-env-setup
- langbot-testing
@@ -24,7 +23,7 @@ env:
- LANGBOT_LOCAL_AGENT_PIPELINE_NAME
preconditions:
- "LANGBOT_LOCAL_AGENT_PIPELINE_URL or LANGBOT_LOCAL_AGENT_PIPELINE_NAME points to the local-agent pipeline under test."
- "LangBot is started with the Box backend intended for this run, such as e2b, nsjail, or explicit host development mode."
- "LangBot is started with the sandbox backend intended for this run, such as e2b or nsjail."
- "The selected model route supports tool/function calling strongly enough to invoke sandbox tools."
steps:
- "Start LangBot with the target sandbox backend and confirm the Box status UI or LANGBOT_BACKEND_URL /api/v1/box/status reports the expected backend."
@@ -34,7 +33,7 @@ steps:
checks:
- "UI: Debug Chat final assistant response contains E2E_OK:<skill-name>."
- "Logs: The model called exec, register_skill, activate, then exec again from the activated skill path."
- "Logs: The selected backend name is the expected one, such as e2b, nsjail, or host."
- "Logs: The selected backend name is the expected one, such as e2b or nsjail."
- "Skill store: The registered package and activated writeback match references/sandbox-skill-authoring.md."
- "Box status: recent_error_count is 0 after the run."
evidence_required:
@@ -4,7 +4,7 @@
Verify that Local Agent can use sandbox tools to create, register, activate, and use a LangBot skill package through the same path a user would exercise in Debug Chat.
This flow applies to Docker, nsjail, E2B, and the explicit host development backend. Host runs commands directly as the Box Runtime user and must never be treated as sandbox-isolation coverage. API calls are useful diagnostics, but the primary pass/fail signal is the model-driven Debug Chat tool sequence.
This flow applies to Docker, nsjail, and E2B backends. API calls are useful diagnostics, but the primary pass/fail signal is the model-driven Debug Chat tool sequence.
## Preconditions
@@ -13,7 +13,6 @@ This flow applies to Docker, nsjail, E2B, and the explicit host development back
- `BOX_BACKEND=e2b` when validating E2B.
- `BOX_BACKEND=nsjail` when validating nsjail.
- `BOX_BACKEND=local` or `docker` when validating local container fallback.
- `BOX_BACKEND=host` only when validating explicit, trusted local direct execution.
3. Confirm `/api/v1/box/status` reports `available: true` and the expected backend name.
4. Confirm Debug Chat uses a model with function-calling ability.
5. Confirm backend logs say native sandbox tools are available.
@@ -72,7 +71,7 @@ Backend logs should show:
- `register_skill`
- `activate`
- a second `exec` whose workdir is `/workspace/.skills/<skill-name>`
- `backend=e2b`, `backend=nsjail`, `backend=host`, or the expected local backend
- `backend=e2b`, `backend=nsjail`, or the expected local backend
After the run, verify the skill store through the UI or API:
@@ -126,8 +125,6 @@ For E2B raw HTTP diagnostics, include a valid template id such as `base`; a miss
- Session metadata should keep LangBot logical paths such as `/workspace`; storing provider-internal paths can make later requests look incompatible.
- nsjail versions differ. Some expose only `--disable_clone_new*` flags and use `--bindmount` instead of `--rw_bind`.
- On WSL, cgroup v2 may exist but not be writable. The backend should warn and fall back to rlimits rather than fail the sandbox.
- The host backend does not honor sandbox image, network, rootfs, process, or
resource isolation. Use a disposable workspace and low-privilege account.
- If `ALL_PROXY` uses a SOCKS URL and `socksio` is not installed, some Python HTTP clients can fail during startup. Prefer consistent HTTP proxy variables unless SOCKS support is installed.
## Related Troubleshooting
@@ -3,7 +3,7 @@ title: "Native sandbox tools are unavailable even though a backend is configured
date: 2026-05-18
symptoms:
- "Backend logs show Native sandbox tools (exec/read/write/edit/glob/grep) are NOT available."
- "The Box runtime later reports that E2B, nsjail, Docker, or explicit host mode is configured."
- "The Box runtime later reports that E2B, nsjail, or Docker is configured."
- "Debug Chat does not expose exec, register_skill, or activate as usable tools."
patterns:
- "Native sandbox tools ... are NOT available"
@@ -19,7 +19,6 @@ fix_steps:
- "Ensure the Box runtime reselects a backend when get_backend_info is called and the cached backend is empty."
- "For E2B, verify the key without printing it and confirm any required template setting."
- "For nsjail, run nsjail --help and confirm the binary is on PATH for the LangBot process."
- "For trusted local development only, explicitly set box.backend=host; never use host as a production sandbox fallback."
verification: "Run sandbox-skill-authoring-e2e. Logs should show Native sandbox tools are available and /api/v1/box/status should report available=true with the expected backend."
related_cases:
- sandbox-skill-authoring-e2e
+1 -1
View File
@@ -16,7 +16,7 @@ asciiart = r"""
|___/
⭐️ Open Source 开源地址: https://github.com/langbot-app/LangBot
📖 Documentation 文档地址: https://langbot.app/docs
📖 Documentation 文档地址: https://docs.langbot.app
"""
+11 -38
View File
@@ -1,14 +1,13 @@
from __future__ import annotations
import asyncio
import json
import os
import typing
from pathlib import Path
import httpx
import typing
import json
from .errors import DifyAPIError
from pathlib import Path
import os
_MAX_DIFY_RESPONSE_BYTES = 1024 * 1024
_MAX_DIFY_SSE_LINE_BYTES = 1024 * 1024
@@ -16,32 +15,6 @@ _MAX_DIFY_STREAM_BYTES = 16 * 1024 * 1024
_MAX_DIFY_UPLOAD_BYTES = 10 * 1024 * 1024
def _decode_sse_data(line: bytes) -> dict[str, typing.Any] | None:
data = line[5:].strip()
if not data or data == b'[DONE]':
return None
try:
payload = json.loads(data.decode('utf-8'))
except (json.JSONDecodeError, UnicodeDecodeError) as exc:
raise DifyAPIError('Dify SSE data line is not valid JSON') from exc
if not isinstance(payload, dict):
raise DifyAPIError('Dify SSE event is not a JSON object')
return payload
def _decode_upload_response(body: bytes) -> dict[str, typing.Any]:
try:
response = json.loads(body)
except (json.JSONDecodeError, UnicodeDecodeError) as exc:
raise DifyAPIError('Dify upload response is not valid JSON') from exc
if not isinstance(response, dict):
raise DifyAPIError('Dify upload response is not a JSON object')
payload = response.get('data', response)
if not isinstance(payload, dict) or not isinstance(payload.get('id'), str) or not payload['id']:
raise DifyAPIError('Dify upload response does not contain a valid file id')
return payload
async def _read_limited_response(
response: httpx.Response,
*,
@@ -83,16 +56,16 @@ async def _iter_sse_json(
line = raw_line.rstrip(b'\r').strip()
if not line or not line.startswith(b'data:'):
continue
payload = _decode_sse_data(line)
if payload is not None:
payload = json.loads(line[5:].decode('utf-8', errors='replace'))
if isinstance(payload, dict):
yield payload
if len(buffer) > _MAX_DIFY_SSE_LINE_BYTES:
raise DifyAPIError('Dify SSE event exceeds the runtime limit')
line = bytes(buffer).rstrip(b'\r').strip()
if line.startswith(b'data:'):
payload = _decode_sse_data(line)
if payload is not None:
payload = json.loads(line[5:].decode('utf-8', errors='replace'))
if isinstance(payload, dict):
yield payload
@@ -269,7 +242,7 @@ class AsyncDifyServiceClient:
file: httpx._types.FileTypes,
user: str,
timeout: float = 30.0,
) -> dict[str, typing.Any]:
) -> str:
# 处理 Path 对象
if isinstance(file, Path):
if not file.exists():
@@ -298,6 +271,6 @@ class AsyncDifyServiceClient:
timeout=timeout,
) as response:
body = await _read_limited_response(response)
if response.status_code not in (200, 201):
if response.status_code != 201:
raise DifyAPIError(f'{response.status_code} {body.decode(errors="replace")}')
return _decode_upload_response(body)
return json.loads(body)
+2 -3
View File
@@ -697,10 +697,9 @@ class DingTalkClient:
if not await self.check_access_token():
await self.get_access_token()
template_params = dict(card_param_map or {})
cardData: dict = {'cardParamMap': _stringify_card_param_map(card_param_map)}
if card_data_config is not None:
template_params['config'] = card_data_config
cardData: dict = {'cardParamMap': _stringify_card_param_map(template_params)}
cardData['config'] = json.dumps(card_data_config)
body: dict = {
'cardTemplateId': card_template_id,
-63
View File
@@ -422,69 +422,6 @@ class QQOfficialClient:
await self.logger.error(f'Failed to send private message: {response_data}')
raise ValueError(response)
async def _send_markdown_msg(
self,
target_type: str,
target_id: str,
content: str,
msg_id: Optional[str] = None,
event_id: Optional[str] = None,
msg_seq: int = 1,
) -> None:
"""Send a Markdown message to a C2C user or QQ group."""
if not await self.check_access_token():
await self.get_access_token()
if target_type == 'c2c':
url = f'{self.base_url}/v2/users/{target_id}/messages'
elif target_type == 'group':
url = f'{self.base_url}/v2/groups/{target_id}/messages'
else:
raise ValueError(f'Unsupported Markdown target type: {target_type}')
data: dict[str, Any] = {
'msg_type': 2,
'markdown': {'content': content},
'msg_seq': msg_seq,
}
if msg_id:
data['msg_id'] = msg_id
if event_id:
data['event_id'] = event_id
async with self._http_client_context() as client:
headers = {
'Authorization': f'QQBot {self.access_token}',
'Content-Type': 'application/json',
}
response = await client.post(url, headers=headers, json=data)
if response.status_code != 200:
response_data = await httpclient.parse_json_response(response)
await self.logger.error(f'Failed to send Markdown message: {response_data}')
raise ValueError(response)
async def send_private_markdown_msg(
self,
user_openid: str,
content: str,
msg_id: Optional[str] = None,
event_id: Optional[str] = None,
msg_seq: int = 1,
) -> None:
"""Send a Markdown C2C message."""
await self._send_markdown_msg('c2c', user_openid, content, msg_id, event_id, msg_seq)
async def send_group_markdown_msg(
self,
group_openid: str,
content: str,
msg_id: Optional[str] = None,
event_id: Optional[str] = None,
msg_seq: int = 1,
) -> None:
"""Send a Markdown QQ group message."""
await self._send_markdown_msg('group', group_openid, content, msg_id, event_id, msg_seq)
async def send_group_text_msg(
self,
group_openid: str,
@@ -46,14 +46,6 @@ CMD_RESPOND_MSG = 'aibot_respond_msg'
CMD_RESPOND_WELCOME = 'aibot_respond_welcome_msg'
CMD_RESPOND_UPDATE = 'aibot_respond_update_msg'
CMD_SEND_MSG = 'aibot_send_msg'
# Media upload protocol (3 steps: init -> chunk * N -> finish). The
# command names below match the WeCom AI Bot long-connection protocol.
CMD_UPLOAD_INIT = 'aibot_upload_media_init'
CMD_UPLOAD_CHUNK = 'aibot_upload_media_chunk'
CMD_UPLOAD_FINISH = 'aibot_upload_media_finish'
# Default upload chunk size: 512 KB before base64 encoding.
_UPLOAD_CHUNK_SIZE = 512 * 1024
_DEDUP_CACHE_MAX = 4096
_STREAM_CACHE_MAX = 1024
@@ -503,145 +495,6 @@ class WecomBotWsClient:
body['chatid'] = chat_id
return await self._send_reply(req_id, body, cmd=CMD_SEND_MSG)
# ------------------------------------------------------------------
# Media upload (image / voice / file)
# ------------------------------------------------------------------
async def upload_media(
self,
data: bytes,
filename: str = 'attachment',
media_type: str = 'file',
) -> Optional[dict]:
"""Upload *data* to the WeCom AI Bot CDN and return the parsed ACK.
Implements the three-step protocol documented for the WeCom
AI Bot:
1. ``aibot_upload_media_init`` — declare media type, file name,
size, MD5 and chunk count; receive ``upload_id``.
2. ``aibot_upload_media_chunk`` — send each chunk (base64-encoded
bytes) until done; receive per-chunk ACK.
3. ``aibot_upload_media_finish`` — finalize the upload; receive
``media_id``.
Returns a dict with the final ``media_id`` (and the raw
``finish`` ACK) on success, or ``None`` on any failure. The
caller is expected to ignore the result and continue
gracefully — the framework will keep working without media
delivery.
"""
import base64 as _b64
import hashlib as _hl
if not data:
return None
file_size = len(data)
file_md5 = _hl.md5(data).hexdigest()
total_chunks = (file_size + _UPLOAD_CHUNK_SIZE - 1) // _UPLOAD_CHUNK_SIZE
if total_chunks == 0:
total_chunks = 1
# Step 1: init.
init_req_id = _generate_req_id(CMD_UPLOAD_INIT)
init_body = {
'type': media_type,
'filename': filename,
'total_size': file_size,
'total_chunks': total_chunks,
'md5': file_md5,
}
init_ack = await self._send_reply(
init_req_id,
init_body,
cmd=CMD_UPLOAD_INIT,
)
if not init_ack or init_ack.get('errcode', 0) != 0:
await self.logger.warning(f'upload_media init failed: ack={init_ack!r}')
return None
upload_id = (
init_ack.get('upload_id')
or init_ack.get('body', {}).get('upload_id')
or init_ack.get('data', {}).get('upload_id')
)
if not upload_id:
await self.logger.warning(f'upload_media init returned no upload_id: ack={init_ack!r}')
return None
# Step 2: chunks.
for index in range(total_chunks):
start = index * _UPLOAD_CHUNK_SIZE
end = min(start + _UPLOAD_CHUNK_SIZE, file_size)
chunk_bytes = data[start:end]
chunk_req_id = _generate_req_id(CMD_UPLOAD_CHUNK)
chunk_body = {
'upload_id': upload_id,
'chunk_index': index,
'base64_data': _b64.b64encode(chunk_bytes).decode('ascii'),
}
chunk_ack = await self._send_reply(
chunk_req_id,
chunk_body,
cmd=CMD_UPLOAD_CHUNK,
)
if not chunk_ack or chunk_ack.get('errcode', 0) != 0:
await self.logger.warning(f'upload_media chunk {index} failed: ack={chunk_ack!r}')
return None
# Step 3: finish.
finish_req_id = _generate_req_id(CMD_UPLOAD_FINISH)
finish_body = {'upload_id': upload_id}
finish_ack = await self._send_reply(
finish_req_id,
finish_body,
cmd=CMD_UPLOAD_FINISH,
)
if not finish_ack or finish_ack.get('errcode', 0) != 0:
await self.logger.warning(f'upload_media finish failed: ack={finish_ack!r}')
return None
media_id = (
finish_ack.get('media_id')
or finish_ack.get('body', {}).get('media_id')
or finish_ack.get('data', {}).get('media_id')
)
if not media_id:
await self.logger.warning(f'upload_media finish returned no media_id: ack={finish_ack!r}')
return None
return {'media_id': media_id, 'ack': finish_ack}
async def _reply_media(
self,
req_id: str,
media_id: str,
kind: str,
) -> Optional[dict]:
"""Send a media reply (image / voice / file) referencing *media_id*.
``kind`` is one of ``'image'``, ``'voice'``, ``'file'``. Uses
the standard ``aibot_respond_msg`` command with a per-kind
body key (matches the convention documented for the WeCom
AI Bot SDK).
"""
if kind not in {'image', 'voice', 'file'}:
await self.logger.warning(f'_reply_media called with unknown kind={kind!r}')
return None
body = {
'msgtype': kind,
kind: {'media_id': media_id},
}
return await self._send_reply(req_id, body, cmd=CMD_RESPOND_MSG)
async def reply_image(self, req_id: str, media_id: str) -> Optional[dict]:
return await self._reply_media(req_id, media_id, 'image')
async def reply_file(self, req_id: str, media_id: str) -> Optional[dict]:
return await self._reply_media(req_id, media_id, 'file')
async def reply_voice(self, req_id: str, media_id: str) -> Optional[dict]:
return await self._reply_media(req_id, media_id, 'voice')
async def push_stream_chunk(self, msg_id: str, content: str, is_final: bool = False) -> bool:
"""Push a streaming chunk for a given message ID.
@@ -936,13 +789,6 @@ class WecomBotWsClient:
'chat_type': message_data.get('type', 'single'),
}
self._prune_stream_state()
# Send an initial empty stream frame so the WeCom client
# shows its built-in loading spinner while the pipeline
# processes the message (e.g. RAG retrieval).
try:
await self.reply_stream(req_id, stream_id, '', finish=False)
except Exception:
await self.logger.warning(f'Failed to send initial stream frame: {traceback.format_exc()}')
message_data['stream_id'] = stream_id
message_data['req_id'] = req_id
@@ -295,34 +295,6 @@ class WecomCSClient:
raise Exception('Failed to send message')
return data
@_bounded_token_retry
async def send_image_msg(self, open_kfid: str, external_userid: str, msgid: str, media_id: str):
if not await self.check_access_token():
self.access_token = await self.get_access_token(self.secret)
url = f'{self.base_url}/kf/send_msg?access_token={self.access_token}'
payload = {
'touser': external_userid,
'open_kfid': open_kfid,
'msgid': msgid,
'msgtype': 'image',
'image': {
'media_id': media_id,
},
}
async with self._http_client_context() as client:
response = await client.post(url, json=payload)
data = await httpclient.parse_json_response(response)
if data['errcode'] == 40014 or data['errcode'] == 42001:
self.access_token = await self.get_access_token(self.secret)
return await self.send_image_msg(open_kfid, external_userid, msgid, media_id)
if data['errcode'] != 0:
await self.logger.error(f'发送图片失败:{data}')
raise Exception('Failed to send image message')
return data
async def handle_callback_request(self):
"""处理回调请求(独立端口模式,使用全局 request)。"""
return await self._handle_callback_internal(request)
@@ -15,7 +15,6 @@ from ....workspace.collaboration import MembershipPermissionError, WorkspaceColl
from ....workspace.errors import WorkspaceNotFoundError
from ....cloud.entitlements import EntitlementUnavailableError
from ....core.errors import TaskCapacityError
from ....provider.modelmgr.codex_errors import CodexProviderError
from ..authz import (
AuthenticationDeniedError,
AuthorizationError,
@@ -248,8 +247,6 @@ class RouterGroup(abc.ABC):
return await f(*args, **kwargs)
except Exception as e: # 自动 500
if isinstance(e, CodexProviderError):
return self.http_status(e.status_code, e.error_code, str(e))
if isinstance(e, AuthorizationError):
return self.http_status(e.status_code, e.error_code, str(e))
if isinstance(e, WorkspaceNotFoundError):
@@ -5,7 +5,6 @@ import quart
from ...authz import Permission
from ...context import RequestContext
from ...service.monitoring_traffic import get_traffic_series
from .. import group
@@ -219,7 +218,6 @@ class MonitoringRouterGroup(group.RouterGroup):
pipeline_ids = quart.request.args.getlist('pipelineId')
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
user_query = quart.request.args.get('userQuery')
is_active_str = quart.request.args.get('isActive')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
@@ -239,7 +237,6 @@ class MonitoringRouterGroup(group.RouterGroup):
pipeline_ids=pipeline_ids if pipeline_ids else None,
start_time=start_time,
end_time=end_time,
user_query=user_query,
is_active=is_active,
limit=limit,
offset=offset,
@@ -378,14 +375,6 @@ class MonitoringRouterGroup(group.RouterGroup):
return self.success(
data={
'traffic': await get_traffic_series(
self.ap,
request_context,
bot_ids=bot_ids or None,
pipeline_ids=pipeline_ids or None,
start_time=start_time,
end_time=end_time,
),
'overview': overview,
'messages': messages,
'llmCalls': llm_calls,
@@ -407,15 +396,7 @@ class MonitoringRouterGroup(group.RouterGroup):
@self.route('/sessions/<session_id>/analysis', methods=['GET'], permission=Permission.RESOURCE_VIEW)
async def get_session_analysis(session_id: str, request_context: RequestContext) -> str:
"""Get detailed analysis for a specific session"""
start_time = parse_iso_datetime(quart.request.args.get('startTime'))
end_time = parse_iso_datetime(quart.request.args.get('endTime'))
analysis = await self.ap.monitoring_service.get_session_analysis(
request_context,
session_id,
start_time=start_time,
end_time=end_time,
bot_id=quart.request.args.get('botId'),
)
analysis = await self.ap.monitoring_service.get_session_analysis(request_context, session_id)
# Always return success with the analysis data
# The frontend will handle the 'found: false' case
@@ -39,13 +39,7 @@ class PipelinesRouterGroup(group.RouterGroup):
permission=Permission.RESOURCE_MANAGE,
)
async def _(request_context: RequestContext) -> str:
pipeline_data = await quart.request.json
create_as_default = pipeline_data.get('is_default') is True
pipeline_uuid = await self.ap.pipeline_service.create_pipeline(
request_context,
pipeline_data,
default=create_as_default,
)
pipeline_uuid = await self.ap.pipeline_service.create_pipeline(request_context, await quart.request.json)
return self.success(data={'uuid': pipeline_uuid})
@self.route(
@@ -1,6 +1,7 @@
import asyncio
import dataclasses
import mimetypes
import os
import quart
@@ -1133,3 +1134,224 @@ class AdaptersRouterGroup(group.RouterGroup):
if session and session.get('task') and not session['task'].done():
session['task'].cancel()
return self.success(data={})
# -----------------------------------------------------------------------
# Itchat WeChat QR Code Login
# -----------------------------------------------------------------------
_itchat_login_sessions: dict = {}
_ITCHAT_SESSION_TTL = 600 # 10 minutes (allows multiple QR regenerations)
def _cleanup_expired_itchat_sessions():
import time
now = time.time()
expired = [
sid for sid, s in _itchat_login_sessions.items() if now - s.get('created_at', 0) > _ITCHAT_SESSION_TTL
]
for sid in expired:
session = _itchat_login_sessions.pop(sid, None)
if session:
core = session.get('core')
if core:
try:
core.alive = False
core.isLogging = False
except Exception:
pass
@self.route('/itchat/login', methods=['POST'])
async def _() -> str:
"""Start itchat WeChat QR code login. Returns session_id + QR code data URL."""
import uuid
import time
import base64
import threading
_cleanup_expired_itchat_sessions()
session_id = str(uuid.uuid4())
loop = asyncio.get_running_loop()
status_dir = os.path.join('data', 'itchat')
os.makedirs(status_dir, exist_ok=True)
qr_path = os.path.join(status_dir, f'{session_id}-QR.png')
session = {
'status': 'pending',
'qr_data_url': None,
'expire_at': None,
'nickname': None,
'error': None,
'created_at': time.time(),
'thread': None,
'logged_in': threading.Event(),
'core': None,
}
_itchat_login_sessions[session_id] = session
def _run_itchat_login():
try:
from itchat.core import Core
from itchat.content import TEXT as _TEXT
from langbot.pkg.platform.sources.itchat import ItchatAdapter
for f in (qr_path,):
try:
os.remove(f)
except OSError:
pass
_core = Core()
session['core'] = _core
def on_login():
try:
_core.get_friends(update=True)
user_info = _core.loginInfo.get('User', {})
nick = ItchatAdapter._get_obj_value(user_info, 'NickName', 'unknown')
wxid = ItchatAdapter._get_obj_value(user_info, 'UserName')
except Exception:
nick = 'unknown'
wxid = ''
session['nickname'] = nick
session['wxid'] = wxid
print(f'[itchat-login] Login success: {nick}', flush=True)
# Dump login status so the adapter can hot-reload it
try:
if not wxid:
raise ValueError('Unable to detect WeChat wxid after login')
account_status_path = ItchatAdapter.login_status_path_for_account(wxid)
_core.dump_login_status(account_status_path)
session['login_status_path'] = account_status_path
session['status'] = 'success'
print(f'[itchat-login] Session saved to {account_status_path}', flush=True)
except Exception as e:
session['status'] = 'error'
session['error'] = str(e)
print(f'[itchat-login] Failed to save session: {e}', flush=True)
finally:
session['logged_in'].set()
# Stop the message loop - we only needed the session for QR login
_core.alive = False
def on_qr(**kwargs):
qr_bytes = kwargs.get('qrcode', b'')
status = kwargs.get('status', '')
print(f'[itchat-login] QR callback: status={status}, bytes={len(qr_bytes)}', flush=True)
if status == '200':
return
# Only update QR image on new QR generation (status='0')
# or when status changes to '408' (timeout, QR may refresh)
if qr_bytes and status == '0':
b64 = base64.b64encode(qr_bytes).decode('utf-8')
def _update():
session['qr_data_url'] = f'data:image/png;base64,{b64}'
session['expire_at'] = time.time() + 120
session['status'] = 'waiting'
loop.call_soon_threadsafe(_update)
# Register a dummy text handler
@_core.msg_register([_TEXT])
def _dummy(msg):
pass
print('[itchat-login] Step 3: Calling auto_login...', flush=True)
_core.auto_login(
hotReload=False,
loginCallback=on_login,
qrCallback=on_qr,
)
print('[itchat-login] Step 4: auto_login returned, starting run...', flush=True)
_core.run(blockThread=True)
print('[itchat-login] Step 5: run() returned', flush=True)
except SystemExit as e:
print(f'[itchat-login] SystemExit: {e}', flush=True)
session['status'] = 'error'
session['error'] = f'itchat exited: {e}'
session['logged_in'].set()
except Exception as e:
import traceback
print(f'[itchat-login] Exception: {traceback.format_exc()}', flush=True)
session['status'] = 'error'
session['error'] = str(e)
session['logged_in'].set()
t = threading.Thread(target=_run_itchat_login, daemon=True)
t.start()
session['thread'] = t
# Wait for QR code to be ready (max 15 seconds)
for _ in range(30):
if session['qr_data_url'] or session['error'] or session['status'] == 'success':
break
await asyncio.sleep(0.5)
if session['error']:
return self.http_status(502, -1, session['error'])
if session['status'] == 'success':
return self.success(
data={
'session_id': session_id,
'status': 'success',
'nickname': session['nickname'],
'wxid': session.get('wxid', ''),
}
)
if not session['qr_data_url']:
session['status'] = 'error'
session['error'] = 'Timeout waiting for QR code'
return self.http_status(504, -1, 'Timeout waiting for QR code')
return self.success(
data={
'session_id': session_id,
'qr_data_url': session['qr_data_url'],
'expire_at': session['expire_at'],
}
)
@self.route('/itchat/login/status/<session_id>', methods=['GET'])
async def _(session_id: str) -> str:
"""Poll itchat login status."""
session = _itchat_login_sessions.get(session_id)
if not session:
return self.http_status(404, -1, 'Session not found')
data = {
'status': session['status'],
'qr_data_url': session['qr_data_url'],
'expire_at': session['expire_at'],
}
if session['status'] == 'success':
data['nickname'] = session.get('nickname', '')
data['wxid'] = session.get('wxid', '')
_itchat_login_sessions.pop(session_id, None)
elif session['status'] == 'error':
data['error'] = session['error']
_itchat_login_sessions.pop(session_id, None)
return self.success(data=data)
@self.route('/itchat/login/<session_id>', methods=['DELETE'])
async def _(session_id: str) -> str:
"""Cancel and clean up an itchat login session."""
session = _itchat_login_sessions.pop(session_id, None)
if session:
core = session.get('core')
if core:
try:
core.alive = False
core.isLogging = False
except Exception:
pass
thread = session.get('thread')
if thread and thread.is_alive():
# Thread is daemon, will die with the process
pass
return self.success(data={})
@@ -113,24 +113,6 @@ class BotsRouterGroup(group.RouterGroup):
)
return self.success(data={'sent': True})
@self.route(
'/<bot_uuid>/test-inbound',
methods=['POST'],
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.RESOURCE_MANAGE,
)
async def _(bot_uuid: str, request_context: RequestContext) -> str:
json_data = await quart.request.get_json(silent=True) or {}
try:
result = await self.ap.bot_service.send_http_bot_test_message(
request_context,
bot_uuid,
str(json_data.get('message') or ''),
)
except ValueError as exc:
return self.http_status(400, -1, str(exc))
return self.success(data=result)
@self.route(
'/<bot_uuid>/admins',
methods=['GET'],
@@ -3,7 +3,6 @@ import quart
from ....authz import Permission, has_permission
from ....context import RequestContext
from ... import group
from .query import resolve_include_secret
@group.group_class('models/llm', '/api/v1/provider/models/llm')
@@ -17,12 +16,7 @@ class LLMModelsRouterGroup(group.RouterGroup):
)
async def _(request_context: RequestContext) -> str:
provider_uuid = quart.request.args.get('provider_uuid')
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
include_secret = has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE)
if provider_uuid:
models = await self.ap.llm_model_service.get_llm_models_by_provider(
request_context,
@@ -59,16 +53,10 @@ class LLMModelsRouterGroup(group.RouterGroup):
permission=Permission.RESOURCE_VIEW,
)
async def _(model_uuid: str, request_context: RequestContext) -> str:
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
model = await self.ap.llm_model_service.get_llm_model(
request_context,
model_uuid,
include_secret=include_secret,
include_secret=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if model is None:
return self.http_status(404, -1, 'model not found')
@@ -123,12 +111,7 @@ class EmbeddingModelsRouterGroup(group.RouterGroup):
)
async def _(request_context: RequestContext) -> str:
provider_uuid = quart.request.args.get('provider_uuid')
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
include_secret = has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE)
if provider_uuid:
models = await self.ap.embedding_models_service.get_embedding_models_by_provider(
request_context,
@@ -165,16 +148,10 @@ class EmbeddingModelsRouterGroup(group.RouterGroup):
permission=Permission.RESOURCE_VIEW,
)
async def _(model_uuid: str, request_context: RequestContext) -> str:
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
model = await self.ap.embedding_models_service.get_embedding_model(
request_context,
model_uuid,
include_secret=include_secret,
include_secret=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if model is None:
return self.http_status(404, -1, 'model not found')
@@ -231,12 +208,7 @@ class RerankModelsRouterGroup(group.RouterGroup):
)
async def _(request_context: RequestContext) -> str:
provider_uuid = quart.request.args.get('provider_uuid')
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
include_secret = has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE)
if provider_uuid:
models = await self.ap.rerank_models_service.get_rerank_models_by_provider(
request_context,
@@ -273,16 +245,10 @@ class RerankModelsRouterGroup(group.RouterGroup):
permission=Permission.RESOURCE_VIEW,
)
async def _(model_uuid: str, request_context: RequestContext) -> str:
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
model = await self.ap.rerank_models_service.get_rerank_model(
request_context,
model_uuid,
include_secret=include_secret,
include_secret=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if model is None:
return self.http_status(404, -1, 'model not found')
@@ -3,86 +3,11 @@ import quart
from ....authz import Permission, has_permission
from ....context import RequestContext
from ... import group
from .query import resolve_include_secret
@group.group_class('models/providers', '/api/v1/provider/providers')
class ModelProvidersRouterGroup(group.RouterGroup):
async def initialize(self) -> None:
# Subscription authorization is an interactive, browser-user-only surface.
@self.route(
'/<provider_uuid>/codex/status',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.PROVIDER_SECRET_MANAGE,
)
async def codex_status(provider_uuid: str, request_context: RequestContext):
try:
return self.success(
data=await self.ap.provider_service.codex_auth.status(request_context, provider_uuid)
)
except ValueError as exc:
return self.http_status(400, -1, str(exc))
@self.route(
'/<provider_uuid>/codex/device',
methods=['POST'],
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.PROVIDER_SECRET_MANAGE,
)
async def codex_device(provider_uuid: str, request_context: RequestContext):
try:
return self.success(
data=await self.ap.provider_service.codex_auth.start(request_context, provider_uuid)
)
except ValueError as exc:
return self.http_status(400, -1, str(exc))
@self.route(
'/<provider_uuid>/codex/device/poll',
methods=['POST'],
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.PROVIDER_SECRET_MANAGE,
)
async def codex_poll(provider_uuid: str, request_context: RequestContext):
body = await quart.request.get_json()
if not isinstance(body, dict):
return self.http_status(400, -1, 'JSON object required')
try:
return self.success(
data=await self.ap.provider_service.codex_auth.poll(
request_context, provider_uuid, body.get('authorization_id')
)
)
except ValueError as exc:
return self.http_status(400, -1, str(exc))
@self.route(
'/<provider_uuid>/codex/auth',
methods=['DELETE'],
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.PROVIDER_SECRET_MANAGE,
)
async def codex_disconnect(provider_uuid: str, request_context: RequestContext):
try:
await self.ap.provider_service.codex_auth.disconnect(request_context, provider_uuid)
return self.success()
except ValueError as exc:
return self.http_status(400, -1, str(exc))
@self.route(
'/<provider_uuid>/codex/device/<authorization_id>',
methods=['DELETE'],
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.PROVIDER_SECRET_MANAGE,
)
async def codex_cancel(provider_uuid: str, authorization_id: str, request_context: RequestContext):
try:
await self.ap.provider_service.codex_auth.cancel(request_context, provider_uuid, authorization_id)
return self.success()
except ValueError as exc:
return self.http_status(400, -1, str(exc))
@self.route(
'',
methods=['GET'],
@@ -90,15 +15,9 @@ class ModelProvidersRouterGroup(group.RouterGroup):
permission=Permission.RESOURCE_VIEW,
)
async def _(request_context: RequestContext) -> str:
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
providers = await self.ap.provider_service.get_providers(
request_context,
include_secret=include_secret,
include_secret=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
for provider in providers:
counts = await self.ap.provider_service.get_provider_model_counts(request_context, provider['uuid'])
@@ -128,16 +47,10 @@ class ModelProvidersRouterGroup(group.RouterGroup):
permission=Permission.RESOURCE_VIEW,
)
async def _(provider_uuid: str, request_context: RequestContext) -> str:
include_secret, error = resolve_include_secret(
quart.request.args.get('include_secret'),
permitted=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if error:
return self.http_status(400, -1, error)
provider = await self.ap.provider_service.get_provider(
request_context,
provider_uuid,
include_secret=include_secret,
include_secret=has_permission(request_context, Permission.PROVIDER_SECRET_MANAGE),
)
if provider is None:
return self.http_status(404, -1, 'provider not found')
@@ -169,15 +82,7 @@ class ModelProvidersRouterGroup(group.RouterGroup):
)
async def _(provider_uuid: str, request_context: RequestContext) -> str:
try:
cascade_values = quart.request.args.getlist('cascade')
if cascade_values:
if len(cascade_values) != 1 or cascade_values[0] not in ('true', 'false'):
return self.http_status(400, -1, 'cascade must be a single true or false value')
await self.ap.provider_service.delete_provider(
request_context, provider_uuid, cascade=cascade_values[0] == 'true'
)
else:
await self.ap.provider_service.delete_provider(request_context, provider_uuid)
await self.ap.provider_service.delete_provider(request_context, provider_uuid)
return self.success()
except ValueError as e:
return self.http_status(400, -1, str(e))
@@ -1,15 +0,0 @@
from __future__ import annotations
def resolve_include_secret(raw_value: str | None, *, permitted: bool) -> tuple[bool, str | None]:
"""Resolve the optional secret projection query parameter."""
if raw_value is None:
return permitted, None
value = raw_value.strip().lower()
if value == 'false':
return False, None
if value == 'true':
return permitted, None
return False, 'include_secret must be either true or false'
@@ -7,116 +7,14 @@ from .. import group
from .....utils import constants
from .....entity.persistence.metadata import WorkspaceMetadata
from ...authz import Permission
from ...context import PrincipalType, RequestContext
from ...context import RequestContext
from .....provider.tools.loaders.mcp_policy import stdio_mcp_enabled
from .....workspace.invitation_delivery import InvitationDeliveryService
SYSTEM_CAPABILITY_OPERATIONS = (
'bot.list',
'bot.get',
'bot.create',
'bot.update',
'bot.delete',
'pipeline.list',
'pipeline.get',
'pipeline.create',
'pipeline.update',
'pipeline.delete',
'pipeline.copy',
'task.list',
'task.get',
'knowledge_base.list',
'knowledge_base.get',
'knowledge_base.create',
'knowledge_base.update',
'knowledge_base.delete',
'knowledge_base.file.list',
'knowledge_base.file.store',
'knowledge_base.file.delete',
'knowledge_base.retrieve',
'file.document.upload',
'plugin.install.github',
'plugin.install.marketplace',
'plugin.install.local',
'plugin.upgrade',
'plugin.get',
'plugin.list',
'plugin.config.get',
'plugin.config.update',
'plugin.logs',
'plugin.delete',
'provider.list',
'provider.get',
'provider.create',
'provider.update',
'provider.delete',
'provider.scan_models',
'model.llm.list',
'model.llm.get',
'model.llm.create',
'model.llm.update',
'model.llm.delete',
'model.llm.test',
'model.embedding.list',
'model.embedding.get',
'model.embedding.create',
'model.embedding.update',
'model.embedding.delete',
'model.embedding.test',
'model.rerank.list',
'model.rerank.get',
'model.rerank.create',
'model.rerank.update',
'model.rerank.delete',
'model.rerank.test',
'skill.list',
'skill.get',
'skill.create',
'skill.update',
'skill.delete',
'skill.files.list',
'skill.files.read',
'skill.files.write',
'skill.preview',
'skill.install.github',
'skill.install.upload',
'mcp_server.list',
'mcp_server.get',
'mcp_server.create',
'mcp_server.update',
'mcp_server.delete',
'mcp_server.resources',
'mcp_server.resource_templates',
'mcp_server.resource_read',
'mcp_server.logs',
'mcp_server.test',
)
@group.group_class('system', '/api/v1/system')
class SystemRouterGroup(group.RouterGroup):
async def initialize(self) -> None:
@self.route('/context', methods=['GET'], auth_type=group.AuthType.API_KEY)
async def _(request_context: RequestContext) -> str:
return self.success(
data={
'instance_uuid': request_context.instance_uuid,
'workspace_uuid': request_context.workspace_uuid,
'api_key_id': request_context.principal.api_key_uuid,
'permissions': sorted(request_context.workspace.permissions),
}
)
@self.route('/capabilities', methods=['GET'], auth_type=group.AuthType.API_KEY)
async def _() -> str:
return self.success(
data={
'schema_version': 1,
'operations': {operation: {'supported': True} for operation in SYSTEM_CAPABILITY_OPERATIONS},
}
)
@self.route('/info', methods=['GET'], auth_type=group.AuthType.NONE)
async def _() -> str:
# Read wizard_status and wizard_progress from metadata table
@@ -308,24 +206,10 @@ class SystemRouterGroup(group.RouterGroup):
return self.success(data={})
@self.route(
'/wizard/recommended-model',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.RESOURCE_MANAGE,
)
async def _(request_context: RequestContext) -> str:
"""Resolve Space's best available chat model to this Workspace."""
try:
model = await self.ap.space_service.get_recommended_chat_model(request_context)
except ValueError as exc:
return self.http_status(503, -1, str(exc))
return self.success(data=model)
@self.route(
'/tasks',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.RESOURCE_VIEW,
)
async def _(request_context: RequestContext) -> str:
@@ -344,23 +228,18 @@ class SystemRouterGroup(group.RouterGroup):
instance_uuid=request_context.instance_uuid,
workspace_uuid=request_context.workspace_uuid,
placement_generation=request_context.placement_generation,
public=request_context.principal.principal_type == PrincipalType.API_KEY,
)
)
@self.route(
'/tasks/<task_id>',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
auth_type=group.AuthType.USER_TOKEN,
permission=Permission.RESOURCE_VIEW,
)
async def _(task_id: str, request_context: RequestContext) -> str:
try:
task_index = int(task_id)
except (TypeError, ValueError):
return self.http_status(404, 404, 'Task not found')
task = self.ap.task_mgr.get_task_by_id(
task_index,
int(task_id),
instance_uuid=request_context.instance_uuid,
workspace_uuid=request_context.workspace_uuid,
placement_generation=request_context.placement_generation,
@@ -369,8 +248,6 @@ class SystemRouterGroup(group.RouterGroup):
if task is None:
return self.http_status(404, 404, 'Task not found')
if request_context.principal.principal_type == PrincipalType.API_KEY:
return self.success(data=task.to_public_dict())
return self.success(data=task.to_dict())
@self.route(
@@ -1,12 +1,7 @@
from __future__ import annotations
import quart
import argon2
import asyncio
import datetime
import hmac
import time
import typing
import uuid
from urllib.parse import parse_qs, urlsplit
@@ -16,36 +11,16 @@ from ...context import RequestContext
from .....cloud.launch import SpaceLaunchError
from ...service.user import ControlPlaneDirectoryRequiredError, PublicRegistrationClosedError
# Fixed-window admission quota for the unauthenticated reset-password endpoint (#2392).
# The admission check and slot bump share ONE synchronous critical section with no await
# points, so concurrent bursts within a single event loop cannot slip past accounting.
# Every admitted attempt consumes quota (regardless of success), which throttles both the
# legacy 24-bit keyspace exhaustion and brute-force on modern high-entropy keys.
# NOTE: this state is process-local; multi-worker deployments need a shared limiter upstream.
_MAX_RESET_ATTEMPTS_PER_WINDOW = 5
_RESET_WINDOW_SECONDS = 15 * 60
_reset_password_state: dict = {'window_started_at': 0.0, 'attempts': 0}
def _admit_reset_attempt(now: float) -> bool:
"""Atomically reserve one reset-password admission slot.
Must stay await-free: running to completion without suspension makes the
check-and-increment atomic under the single-threaded event loop.
"""
st = _reset_password_state
if now - st['window_started_at'] >= _RESET_WINDOW_SECONDS:
st['window_started_at'] = now
st['attempts'] = 0
if st['attempts'] >= _MAX_RESET_ATTEMPTS_PER_WINDOW:
return False
st['attempts'] += 1
return True
@group.group_class('user', '/api/v1/user')
class UserRouterGroup(group.RouterGroup):
@staticmethod
def _origin(value: str) -> tuple[str, str, int | None] | None:
parsed = urlsplit(value)
if parsed.scheme not in {'http', 'https'} or not parsed.hostname:
return None
return parsed.scheme, parsed.hostname.casefold(), parsed.port
def _validate_space_redirect_uri(self, redirect_uri: str, *, bind: bool) -> str:
parsed = urlsplit(redirect_uri)
if (
@@ -63,26 +38,19 @@ class UserRouterGroup(group.RouterGroup):
if query != {'mode': ['bind']}:
raise ValueError('Invalid Space binding redirect_uri')
elif query:
raise ValueError('Invalid LangBot Account login redirect_uri')
raise ValueError('Invalid Space login redirect_uri')
redirect_origin = self._origin(redirect_uri)
api_config = self.ap.instance_config.data.get('api', {})
trusted_origins = {
self._origin(str(api_config.get(config_key, '') or '').strip())
for config_key in ('webui_url', 'webhook_prefix')
}
trusted_origins.discard(None)
if redirect_origin not in trusted_origins:
raise ValueError('Untrusted redirect_uri origin')
return redirect_uri
def _extract_origin_and_rp_id(self, json_data: dict[str, typing.Any] | None = None) -> tuple[str, str]:
origin = ''
if json_data and isinstance(json_data, dict):
origin = json_data.get('origin', '')
if not origin:
origin = quart.request.headers.get('Origin', '')
if not origin:
origin = quart.request.headers.get('Referer', '')
if not origin:
origin = quart.request.url_root.rstrip('/')
parsed = urlsplit(origin)
rp_id = parsed.hostname or 'localhost'
clean_origin = f'{parsed.scheme}://{parsed.netloc}' if parsed.scheme and parsed.netloc else origin.rstrip('/')
return clean_origin, rp_id
async def initialize(self) -> None:
@self.route('/init', methods=['GET', 'POST'], auth_type=group.AuthType.NONE)
async def _() -> str:
@@ -129,12 +97,6 @@ class UserRouterGroup(group.RouterGroup):
@self.route('/reset-password', methods=['POST'], auth_type=group.AuthType.NONE)
async def _() -> str:
# Admit (or reject) BEFORE touching the body or any service call (#2392):
# rejecting requests never reach the slow path, and quota accounting happens
# synchronously at entry, closing the post-await race of burst requests.
if not _admit_reset_attempt(time.monotonic()):
return self.http_status(429, -1, 'Too many attempts, try again later')
json_data = await quart.request.json
user_email = json_data['user']
@@ -152,18 +114,7 @@ class UserRouterGroup(group.RouterGroup):
if user_obj is None:
return self.http_status(400, -1, 'User not found')
stored_key = self.ap.instance_config.data['system']['recovery_key']
try:
key_matches = (
isinstance(recovery_key, str)
and isinstance(stored_key, str)
and hmac.compare_digest(recovery_key.encode(), stored_key.encode())
)
except UnicodeEncodeError:
# JSON can contain lone surrogates, which are not valid UTF-8.
key_matches = False
if not key_matches:
if recovery_key != self.ap.instance_config.data['system']['recovery_key']:
return self.http_status(403, -1, 'Invalid recovery key')
await self.ap.user_service.reset_password(user_email, new_password)
@@ -251,9 +202,6 @@ class UserRouterGroup(group.RouterGroup):
json_data = await quart.request.json
code = json_data.get('code')
state = json_data.get('state')
redirect_uri = json_data.get('redirect_uri') or (
quart.request.url_root.rstrip('/') + '/auth/space/callback'
)
launch_assertion = json_data.get('launch_assertion')
workspace_uuid = json_data.get('workspace_uuid')
@@ -267,11 +215,8 @@ class UserRouterGroup(group.RouterGroup):
return self.fail(1, 'Missing authorization code')
if not state:
return self.fail(1, 'Missing state parameter')
if not str(code).startswith('v4_'):
return self.fail(1, 'Unsupported Space OAuth code contract')
try:
redirect_uri = self._validate_space_redirect_uri(str(redirect_uri), bind=False)
consumed_state = await self.ap.user_service.consume_space_oauth_state_details(state, 'login')
# Exchange code for tokens
launch_workspace_uuid = consumed_state.launch_workspace_uuid
@@ -289,36 +234,24 @@ class UserRouterGroup(group.RouterGroup):
code,
workspace_uuids,
workspace_created_ats,
redirect_uri=redirect_uri,
)
access_token = token_data.get('access_token')
refresh_token = token_data.get('refresh_token')
expires_in = token_data.get('expires_in', 0)
cloud_workspace_uuid = token_data.get('cloud_workspace_uuid')
if not access_token:
return self.fail(1, 'Failed to get access token from Space')
cloud_mode = getattr(getattr(self.ap, 'deployment', None), 'mode', 'oss') == 'cloud'
if cloud_mode and launch_workspace_uuid and launch_workspace_uuid != cloud_workspace_uuid:
return self.fail(1, 'Space OAuth Workspace binding mismatch')
target_workspace_uuid = launch_workspace_uuid or cloud_workspace_uuid
if cloud_mode:
if not target_workspace_uuid:
return self.fail(1, 'Space OAuth response is missing the Cloud Workspace binding')
await self.ap.directory_projection_service.reconcile_workspaces((target_workspace_uuid,))
# Authenticate only after the signed, exact Workspace delta has
# established the Account and membership runtime shadow rows.
# Authenticate and create/update local user
jwt_token, user_obj = await self.ap.user_service.authenticate_space_user(
access_token, refresh_token, expires_in
)
if target_workspace_uuid:
if launch_workspace_uuid:
try:
access = await self.ap.workspace_collaboration_service.resolve_account_workspace(
user_obj.uuid,
target_workspace_uuid,
launch_workspace_uuid,
)
except Exception:
self.ap.logger.warning('Rejected Space OAuth launch for unauthorized Workspace')
@@ -405,9 +338,6 @@ class UserRouterGroup(group.RouterGroup):
if cloud_mode:
capabilities['password_login_enabled'] = False
capabilities['authenticated_invitation_acceptance_enabled'] = cloud_mode
capabilities['invitation_registration_enabled'] = not cloud_mode
capabilities['passkey_login_enabled'] = True
capabilities['passkey_supported'] = True
return self.success(data={'initialized': True, **capabilities})
@self.route('/set-password', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
@@ -452,17 +382,12 @@ class UserRouterGroup(group.RouterGroup):
json_data = await quart.request.json
code = json_data.get('code')
state = json_data.get('state')
redirect_uri = json_data.get('redirect_uri') or (
quart.request.url_root.rstrip('/') + '/auth/space/callback?mode=bind'
)
if not code:
return self.http_status(400, -1, 'Missing authorization code')
if not state:
return self.http_status(400, -1, 'Missing state parameter')
if not str(code).startswith('v4_'):
return self.http_status(400, -1, 'Unsupported Space OAuth code contract')
try:
user_obj = await self.ap.user_service.consume_space_oauth_state(state, 'bind')
@@ -475,10 +400,7 @@ class UserRouterGroup(group.RouterGroup):
return self.http_status(400, -1, 'Only local accounts can bind to Space')
try:
redirect_uri = self._validate_space_redirect_uri(str(redirect_uri), bind=True)
updated_user = await self.ap.user_service.bind_space_account(
user_obj.user, code, redirect_uri=redirect_uri
)
updated_user = await self.ap.user_service.bind_space_account(user_obj.user, code)
jwt_token = await self.ap.user_service.generate_jwt_token(updated_user)
return self.success(
data={
@@ -494,186 +416,10 @@ class UserRouterGroup(group.RouterGroup):
'Bind the LangBot Account with the same email as this local Account',
)
except ValueError:
return self.http_status(400, -1, 'LangBot Account binding failed')
return self.http_status(400, -1, 'Space account binding failed')
except Exception:
raise
@self.route('/passkey/register/options', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
async def _(user_email: str) -> str:
"""Generate WebAuthn registration options for current account."""
allow_modify_login_info = self.ap.instance_config.data.get('system', {}).get(
'allow_modify_login_info', True
)
if not allow_modify_login_info:
return self.http_status(403, -1, 'Modifying login info is disabled')
user_obj = await self.ap.user_service.get_user_by_email(user_email)
if user_obj is None:
return self.http_status(404, -1, 'User not found')
json_data = (await quart.request.json) or {}
origin, rp_id = self._extract_origin_and_rp_id(json_data)
try:
options, challenge_token = await self.ap.user_service.generate_passkey_registration_options(
account_uuid=user_obj.uuid,
rp_id=rp_id,
origin=origin,
rp_name='LangBot',
)
return self.success(data={'options': options, 'challenge_token': challenge_token})
except Exception as e:
return self.fail(1, str(e))
@self.route('/passkey/register/verify', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
async def _(user_email: str) -> str:
"""Verify WebAuthn registration response and save credential."""
allow_modify_login_info = self.ap.instance_config.data.get('system', {}).get(
'allow_modify_login_info', True
)
if not allow_modify_login_info:
return self.http_status(403, -1, 'Modifying login info is disabled')
user_obj = await self.ap.user_service.get_user_by_email(user_email)
if user_obj is None:
return self.http_status(404, -1, 'User not found')
json_data = await quart.request.json
challenge_token = json_data.get('challenge_token')
credential = json_data.get('credential') or json_data.get('response')
name = json_data.get('name')
if not challenge_token or not credential:
return self.fail(1, 'Missing challenge_token or credential')
try:
cred = await self.ap.user_service.verify_and_save_passkey_registration(
challenge_token=challenge_token,
credential_data=credential,
name=name,
)
return self.success(
data={
'uuid': cred.uuid,
'name': cred.name,
'created_at': cred.created_at.isoformat() if cred.created_at else None,
}
)
except Exception as e:
return self.fail(1, str(e))
@self.route('/passkey/auth/options', methods=['POST'], auth_type=group.AuthType.NONE)
async def _() -> str:
"""Generate WebAuthn authentication options for passkey login."""
json_data = (await quart.request.json) or {}
email = json_data.get('email')
origin, rp_id = self._extract_origin_and_rp_id(json_data)
try:
options, challenge_token = await self.ap.user_service.generate_passkey_authentication_options(
rp_id=rp_id,
origin=origin,
email=email,
)
return self.success(data={'options': options, 'challenge_token': challenge_token})
except Exception as e:
return self.fail(1, str(e))
@self.route('/passkey/auth/verify', methods=['POST'], auth_type=group.AuthType.NONE)
async def _() -> str:
"""Verify WebAuthn authentication response and log in."""
json_data = await quart.request.json
challenge_token = json_data.get('challenge_token')
credential = json_data.get('credential') or json_data.get('response')
if not challenge_token or not credential:
return self.fail(1, 'Missing challenge_token or credential')
try:
token, user_obj = await self.ap.user_service.verify_passkey_authentication(
challenge_token=challenge_token,
credential_data=credential,
)
return self.success(
data={
'token': token,
'user': user_obj.user,
}
)
except Exception as e:
return self.fail(1, str(e))
@self.route('/passkeys', methods=['GET'], auth_type=group.AuthType.USER_TOKEN)
async def _(user_email: str) -> str:
"""List registered passkeys for the current user."""
user_obj = await self.ap.user_service.get_user_by_email(user_email)
if user_obj is None:
return self.http_status(404, -1, 'User not found')
passkeys = await self.ap.user_service.get_user_passkeys(user_obj.uuid)
return self.success(
data=[
{
'uuid': pk.uuid,
'name': pk.name,
'aaguid': pk.aaguid,
'transports': pk.transports,
'backed_up': pk.backed_up,
'created_at': pk.created_at.isoformat() if pk.created_at else None,
'last_used_at': pk.last_used_at.isoformat() if pk.last_used_at else None,
}
for pk in passkeys
]
)
@self.route('/passkey/<passkey_uuid>', methods=['PATCH'], auth_type=group.AuthType.USER_TOKEN)
async def _(user_email: str, passkey_uuid: str) -> str:
"""Rename a registered passkey."""
allow_modify_login_info = self.ap.instance_config.data.get('system', {}).get(
'allow_modify_login_info', True
)
if not allow_modify_login_info:
return self.http_status(403, -1, 'Modifying login info is disabled')
user_obj = await self.ap.user_service.get_user_by_email(user_email)
if user_obj is None:
return self.http_status(404, -1, 'User not found')
json_data = await quart.request.json
name = (json_data.get('name') or '').strip()
if not name:
return self.fail(1, 'Passkey name cannot be empty')
updated = await self.ap.user_service.rename_user_passkey(
account_uuid=user_obj.uuid,
passkey_uuid=passkey_uuid,
new_name=name,
)
if not updated:
return self.http_status(404, -1, 'Passkey not found')
return self.success(data={'uuid': updated.uuid, 'name': updated.name})
@self.route('/passkey/<passkey_uuid>', methods=['DELETE'], auth_type=group.AuthType.USER_TOKEN)
async def _(user_email: str, passkey_uuid: str) -> str:
"""Delete/revoke a registered passkey."""
allow_modify_login_info = self.ap.instance_config.data.get('system', {}).get(
'allow_modify_login_info', True
)
if not allow_modify_login_info:
return self.http_status(403, -1, 'Modifying login info is disabled')
user_obj = await self.ap.user_service.get_user_by_email(user_email)
if user_obj is None:
return self.http_status(404, -1, 'User not found')
deleted = await self.ap.user_service.delete_user_passkey(
account_uuid=user_obj.uuid,
passkey_uuid=passkey_uuid,
)
if not deleted:
return self.http_status(404, -1, 'Passkey not found')
return self.success()
async def _handle_space_direct_launch(
self,
launch_assertion: str,
@@ -697,10 +443,6 @@ class UserRouterGroup(group.RouterGroup):
}
)
projection_service = self.ap.directory_projection_service
if projection_service is None:
raise SpaceLaunchError('Cloud directory projection is unavailable')
await projection_service.reconcile_workspaces((launch['workspace_uuid'],))
account = await self.ap.user_service.get_user_by_uuid(launch['account_uuid'])
if account is None:
raise SpaceLaunchError('Launch Account is not projected into Core')
+3 -61
View File
@@ -1,7 +1,6 @@
from __future__ import annotations
import uuid
import json
import sqlalchemy
from ....core import app
@@ -9,8 +8,6 @@ from ....entity.persistence import bot as persistence_bot
from ....entity.persistence import pipeline as persistence_pipeline
from ....workspace.errors import WorkspaceNotFoundError
from .tenant import TenantContext, require_workspace_uuid, scope_statement
from ....utils import httpclient
from ....platform.sources import http_bot_signing
class BotService:
@@ -72,6 +69,8 @@ class BotService:
runtime_bot = await self.ap.platform_mgr.get_bot_by_uuid(context, bot_uuid)
if runtime_bot is not None:
adapter_runtime_values['bot_account_id'] = runtime_bot.adapter.bot_account_id
if hasattr(runtime_bot.adapter, 'get_runtime_status'):
adapter_runtime_values['runtime_status'] = runtime_bot.adapter.get_runtime_status()
# Webhook URL for unified webhook adapters (independent of bot running state)
if persistence_bot['adapter'] in [
@@ -83,7 +82,6 @@ class BotService:
'wecomcs',
'LINE',
'lark',
'http_bot',
]:
webhook_prefix = self.ap.instance_config.data['api'].get('webhook_prefix', 'http://127.0.0.1:5300')
extra_webhook_prefix = self.ap.instance_config.data['api'].get('extra_webhook_prefix', '')
@@ -137,16 +135,7 @@ class BotService:
bot = await self.get_bot(context, bot_data['uuid'], include_secret=True)
try:
await self.ap.platform_mgr.load_bot(context, bot)
except Exception:
# The bot row was already inserted above; without this rollback a
# failing adapter constructor (e.g. a missing optional credential
# key) would leave a permanently disabled orphan bot in the DB.
await self.ap.persistence_mgr.execute_async(
sqlalchemy.delete(persistence_bot.Bot).where(persistence_bot.Bot.uuid == bot_data['uuid'])
)
raise
await self.ap.platform_mgr.load_bot(context, bot)
return bot_data['uuid']
@@ -229,53 +218,6 @@ class BotService:
return [log.to_json() for log in logs], total_count
async def send_http_bot_test_message(
self,
context: TenantContext,
bot_uuid: str,
message: str,
) -> dict:
"""Send a signed test message through the HTTP Bot public ingress."""
bot = await self.get_bot(context, bot_uuid, include_secret=True)
if bot is None:
raise WorkspaceNotFoundError('Bot not found')
if bot.get('adapter') != 'http_bot':
raise ValueError('Inbound test is only available for HTTP Bot')
if not bot.get('enable'):
raise ValueError('Bot must be enabled before sending a test message')
text = message.strip()
if not text or len(text) > 2000:
raise ValueError('Test message must contain 1 to 2000 characters')
payload = {
'session_id': f'wizard-{uuid.uuid4().hex}',
'sender': {'id': 'wizard-user', 'name': 'Wizard Test'},
'message': [{'type': 'Plain', 'text': text}],
}
body = json.dumps(payload, ensure_ascii=False, separators=(',', ':')).encode()
config = bot.get('adapter_config') or {}
headers = {'Content-Type': 'application/json'}
if config.get('signature_required', True):
secret = str(config.get('inbound_secret') or '')
if not secret:
raise ValueError('HTTP Bot inbound signing secret is required')
timestamp, signature = http_bot_signing.sign(secret, body)
headers[http_bot_signing.HEADER_TIMESTAMP] = timestamp
headers[http_bot_signing.HEADER_SIGNATURE] = signature
port = int(self.ap.instance_config.data.get('api', {}).get('port', 5300))
session = httpclient.get_session()
async with session.post(
f'http://127.0.0.1:{port}/bots/{bot_uuid}',
data=body,
headers=headers,
) as response:
result = await httpclient.read_json_limited(response)
if response.status not in {200, 202}:
raise ValueError(result.get('msg') or f'HTTP Bot test failed with status {response.status}')
return result.get('data') or {}
async def send_message(
self,
context: TenantContext,
+9 -16
View File
@@ -446,19 +446,15 @@ class MCPService:
persisted_session = runtime_mcp_session
async def _refresh_and_report() -> None:
try:
needs_start = (
persisted_session.status == MCPSessionStatus.ERROR or persisted_session.session is None
)
if needs_start:
needs_start = persisted_session.status == MCPSessionStatus.ERROR or persisted_session.session is None
if needs_start:
await persisted_session.start()
else:
try:
await persisted_session.refresh()
except Exception:
await persisted_session.start()
else:
try:
await persisted_session.refresh()
except Exception:
await persisted_session.start()
finally:
ctx.metadata['runtime_info'] = persisted_session.get_runtime_info_dict()
ctx.metadata['runtime_info'] = persisted_session.get_runtime_info_dict()
coroutine = _refresh_and_report()
else:
@@ -475,11 +471,8 @@ class MCPService:
async def _run_and_cleanup() -> None:
try:
await test_session.start()
finally:
# start() raises for a failed connection. Preserve the
# terminal runtime state so the UI can render actionable
# failure phases such as OAuth-required.
ctx.metadata['runtime_info'] = test_session.get_runtime_info_dict()
finally:
try:
await test_session.shutdown()
except Exception as exc:
+23 -84
View File
@@ -29,19 +29,6 @@ _DEFAULT_CLEANUP_BATCHES_PER_TABLE = 4
_HARD_MAX_CLEANUP_BATCHES_PER_TABLE = 100
def _normalize_user_id(value: str | int | None) -> str | None:
"""Convert numeric platform IDs before binding a VARCHAR with asyncpg.
Opaque string IDs (including whitespace and leading zeros) and missing
IDs must remain unchanged. Do not silently stringify unsupported objects.
"""
if value is None or isinstance(value, str):
return value
if isinstance(value, int) and not isinstance(value, bool):
return str(value)
raise TypeError('user_id must be a string, integer, or None')
def _workspace_transaction(method):
"""Run an explicit service entrypoint in one Workspace transaction."""
@@ -294,21 +281,19 @@ class MonitoringService:
for _batch_number in range(max_batches):
async def delete_batch() -> tuple[int, int]:
key_columns = list(model_cls.__table__.primary_key.columns)
select_result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(*key_columns)
sqlalchemy.select(pk_column)
.where(model_cls.workspace_uuid == workspace_uuid, ts_column < cutoff)
.limit(batch_size)
)
pk_values = [tuple(row) for row in select_result.all()]
pk_values = list(select_result.scalars().all())
if not pk_values:
return 0, 0
delete_result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.delete(model_cls).where(
model_cls.workspace_uuid == workspace_uuid,
sqlalchemy.tuple_(*key_columns).in_(pk_values),
ts_column < cutoff,
pk_column.in_(pk_values),
)
)
return len(pk_values), int(delete_result.rowcount or 0)
@@ -430,7 +415,7 @@ class MonitoringService:
status: str = 'success',
level: str = 'info',
platform: str | None = None,
user_id: str | int | None = None,
user_id: str | None = None,
user_name: str | None = None,
runner_name: str | None = None,
variables: str | None = None,
@@ -452,7 +437,7 @@ class MonitoringService:
'status': status,
'level': level,
'platform': platform,
'user_id': _normalize_user_id(user_id),
'user_id': user_id,
'user_name': user_name,
'runner_name': runner_name,
'variables': variables,
@@ -625,7 +610,7 @@ class MonitoringService:
pipeline_id: str,
pipeline_name: str,
platform: str | None = None,
user_id: str | int | None = None,
user_id: str | None = None,
user_name: str | None = None,
) -> None:
"""Record a new session"""
@@ -637,29 +622,17 @@ class MonitoringService:
'bot_name': bot_name,
'pipeline_id': pipeline_id,
'pipeline_name': pipeline_name,
'message_count': 1,
'message_count': 0,
'start_time': datetime.datetime.now(datetime.timezone.utc).replace(tzinfo=None),
'last_activity': datetime.datetime.now(datetime.timezone.utc).replace(tzinfo=None),
'is_active': True,
'platform': platform,
'user_id': _normalize_user_id(user_id),
'user_id': user_id,
'user_name': user_name,
}
model = persistence_monitoring.MonitoringSession
dialect = self.ap.persistence_mgr.get_db_engine().dialect.name
insert = postgresql_dialect.insert if dialect == 'postgresql' else sqlite_dialect.insert
statement = insert(model).values(session_data)
await self.ap.persistence_mgr.execute_async(
statement.on_conflict_do_update(
index_elements=['workspace_uuid', 'bot_id', 'session_id'],
set_={
'message_count': model.message_count + 1,
'last_activity': statement.excluded.last_activity,
'pipeline_id': statement.excluded.pipeline_id,
'pipeline_name': statement.excluded.pipeline_name,
},
)
sqlalchemy.insert(persistence_monitoring.MonitoringSession).values(session_data)
)
@_workspace_transaction
@@ -669,7 +642,6 @@ class MonitoringService:
session_id: str,
pipeline_id: str | None = None,
pipeline_name: str | None = None,
bot_id: str | None = None,
) -> bool:
"""Update session last activity time and increment message count.
@@ -679,9 +651,6 @@ class MonitoringService:
True if session was found and updated, False if session doesn't exist.
"""
workspace_uuid = self._require_write_context(context)
bot_id = bot_id if bot_id is not None else context.bot_uuid
if not bot_id:
raise ValueError('Session activity requires a bot_id')
update_values = {
'last_activity': datetime.datetime.now(datetime.timezone.utc).replace(tzinfo=None),
'message_count': persistence_monitoring.MonitoringSession.message_count + 1,
@@ -698,7 +667,6 @@ class MonitoringService:
.where(
persistence_monitoring.MonitoringSession.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringSession.session_id == session_id,
persistence_monitoring.MonitoringSession.bot_id == bot_id,
)
.values(update_values)
)
@@ -801,13 +769,13 @@ class MonitoringService:
message_conditions.append(persistence_monitoring.MonitoringMessage.timestamp >= start_time)
llm_conditions.append(persistence_monitoring.MonitoringLLMCall.timestamp >= start_time)
embedding_conditions.append(persistence_monitoring.MonitoringEmbeddingCall.timestamp >= start_time)
session_conditions.append(persistence_monitoring.MonitoringSession.last_activity >= start_time)
session_conditions.append(persistence_monitoring.MonitoringSession.start_time >= start_time)
if end_time:
message_conditions.append(persistence_monitoring.MonitoringMessage.timestamp <= end_time)
llm_conditions.append(persistence_monitoring.MonitoringLLMCall.timestamp <= end_time)
embedding_conditions.append(persistence_monitoring.MonitoringEmbeddingCall.timestamp <= end_time)
session_conditions.append(persistence_monitoring.MonitoringSession.last_activity <= end_time)
session_conditions.append(persistence_monitoring.MonitoringSession.start_time <= end_time)
# Total messages
message_query = sqlalchemy.select(sqlalchemy.func.count(persistence_monitoring.MonitoringMessage.id))
@@ -1289,7 +1257,6 @@ class MonitoringService:
pipeline_ids: list[str] | None = None,
start_time: datetime.datetime | None = None,
end_time: datetime.datetime | None = None,
user_query: str | None = None,
is_active: bool | None = None,
limit: int = 100,
offset: int = 0,
@@ -1304,17 +1271,9 @@ class MonitoringService:
if pipeline_ids:
conditions.append(persistence_monitoring.MonitoringSession.pipeline_id.in_(pipeline_ids))
if start_time:
conditions.append(persistence_monitoring.MonitoringSession.last_activity >= start_time)
conditions.append(persistence_monitoring.MonitoringSession.start_time >= start_time)
if end_time:
conditions.append(persistence_monitoring.MonitoringSession.last_activity <= end_time)
if user_query and user_query.strip():
user_pattern = f'%{user_query.strip()}%'
conditions.append(
sqlalchemy.or_(
persistence_monitoring.MonitoringSession.user_id.ilike(user_pattern),
persistence_monitoring.MonitoringSession.user_name.ilike(user_pattern),
)
)
conditions.append(persistence_monitoring.MonitoringSession.start_time <= end_time)
if is_active is not None:
conditions.append(persistence_monitoring.MonitoringSession.is_active == is_active)
@@ -1406,9 +1365,6 @@ class MonitoringService:
self,
context: TenantContext,
session_id: str,
start_time: datetime.datetime | None = None,
end_time: datetime.datetime | None = None,
bot_id: str | None = None,
) -> dict:
"""Get bounded session details with full statistics computed in SQL."""
workspace_uuid = require_workspace_uuid(context)
@@ -1418,13 +1374,8 @@ class MonitoringService:
persistence_monitoring.MonitoringSession.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringSession.session_id == session_id,
)
if bot_id is not None:
session_query = session_query.where(persistence_monitoring.MonitoringSession.bot_id == bot_id)
session_result = await self.ap.persistence_mgr.execute_async(session_query.limit(2))
session_rows = session_result.all()
if len(session_rows) > 1:
return {'session_id': session_id, 'found': False, 'ambiguous': True}
session_row = session_rows[0] if session_rows else None
session_result = await self.ap.persistence_mgr.execute_async(session_query)
session_row = session_result.first()
if not session_row:
return {
@@ -1433,7 +1384,6 @@ class MonitoringService:
}
session = session_row[0] if isinstance(session_row, tuple) else session_row
bot_id = session.bot_id
message_stats_result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(
@@ -1461,7 +1411,6 @@ class MonitoringService:
).where(
persistence_monitoring.MonitoringMessage.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringMessage.session_id == session_id,
persistence_monitoring.MonitoringMessage.bot_id == bot_id,
)
)
message_stats = message_stats_result.one()
@@ -1500,7 +1449,6 @@ class MonitoringService:
).where(
persistence_monitoring.MonitoringLLMCall.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringLLMCall.session_id == session_id,
persistence_monitoring.MonitoringLLMCall.bot_id == bot_id,
)
)
llm_stats = llm_stats_result.one()
@@ -1527,22 +1475,15 @@ class MonitoringService:
).where(
persistence_monitoring.MonitoringToolCall.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringToolCall.session_id == session_id,
persistence_monitoring.MonitoringToolCall.bot_id == bot_id,
)
)
tool_stats = tool_stats_result.one()
tool_conditions = [
persistence_monitoring.MonitoringToolCall.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringToolCall.session_id == session_id,
persistence_monitoring.MonitoringToolCall.bot_id == bot_id,
]
if start_time is not None:
tool_conditions.append(persistence_monitoring.MonitoringToolCall.timestamp >= start_time)
if end_time is not None:
tool_conditions.append(persistence_monitoring.MonitoringToolCall.timestamp <= end_time)
tool_query = (
sqlalchemy.select(persistence_monitoring.MonitoringToolCall)
.where(*tool_conditions)
.where(
persistence_monitoring.MonitoringToolCall.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringToolCall.session_id == session_id,
)
.order_by(persistence_monitoring.MonitoringToolCall.timestamp.asc())
.limit(detail_limit + 1)
)
@@ -1563,7 +1504,6 @@ class MonitoringService:
.where(
persistence_monitoring.MonitoringError.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringError.session_id == session_id,
persistence_monitoring.MonitoringError.bot_id == bot_id,
)
.order_by(persistence_monitoring.MonitoringError.timestamp.desc())
.limit(detail_limit + 1)
@@ -2048,9 +1988,9 @@ class MonitoringService:
if pipeline_ids:
conditions.append(persistence_monitoring.MonitoringSession.pipeline_id.in_(pipeline_ids))
if start_time:
conditions.append(persistence_monitoring.MonitoringSession.last_activity >= start_time)
conditions.append(persistence_monitoring.MonitoringSession.start_time >= start_time)
if end_time:
conditions.append(persistence_monitoring.MonitoringSession.last_activity <= end_time)
conditions.append(persistence_monitoring.MonitoringSession.start_time <= end_time)
query = sqlalchemy.select(persistence_monitoring.MonitoringSession).order_by(
persistence_monitoring.MonitoringSession.last_activity.desc()
@@ -2084,7 +2024,6 @@ class MonitoringService:
# ========== Feedback Methods ==========
@_workspace_transaction
async def record_feedback(
self,
context: ExecutionContext,
@@ -2099,7 +2038,7 @@ class MonitoringService:
session_id: str | None = None,
message_id: str | None = None,
stream_id: str | None = None,
user_id: str | int | None = None,
user_id: str | None = None,
platform: str | None = None,
) -> str | None:
"""Record user feedback (like/dislike) from AI Bot conversation.
@@ -2155,7 +2094,7 @@ class MonitoringService:
'session_id': session_id,
'message_id': message_id,
'stream_id': stream_id,
'user_id': _normalize_user_id(user_id),
'user_id': user_id,
'platform': platform,
}
dialect_name = self.ap.persistence_mgr.get_db_engine().dialect.name
@@ -1,83 +0,0 @@
"""Bounded traffic aggregation, independent of record-list pagination."""
from __future__ import annotations
import datetime
import typing
import sqlalchemy
from ....entity.persistence.monitoring import MonitoringLLMCall, MonitoringMessage
from .tenant import TenantContext, require_workspace_uuid
if typing.TYPE_CHECKING:
from ....core.app import Application
MAX_TRAFFIC_POINTS = 1000
async def get_traffic_series(
ap: Application,
context: TenantContext,
*,
bot_ids: list[str] | None = None,
pipeline_ids: list[str] | None = None,
start_time: datetime.datetime | None = None,
end_time: datetime.datetime | None = None,
) -> dict:
"""Count all matching records in UTC buckets, returning at most 1000 points."""
workspace_uuid = require_workspace_uuid(context)
bucket = 'hour' if start_time and end_time and end_time - start_time <= datetime.timedelta(days=7) else 'day'
step = datetime.timedelta(hours=1) if bucket == 'hour' else datetime.timedelta(days=1)
postgres = ap.persistence_mgr.get_db_engine().dialect.name == 'postgresql'
points: dict[datetime.datetime, dict[str, int]] = {}
truncated = False
for model, field in ((MonitoringMessage, 'messages'), (MonitoringLLMCall, 'llm_calls')):
timestamp = model.timestamp
if postgres:
time_bucket = sqlalchemy.func.date_trunc(bucket, timestamp)
else:
pattern = '%Y-%m-%dT%H:00:00' if bucket == 'hour' else '%Y-%m-%dT00:00:00'
time_bucket = sqlalchemy.func.strftime(pattern, timestamp)
conditions = [model.workspace_uuid == workspace_uuid]
if bot_ids:
conditions.append(model.bot_id.in_(bot_ids))
if pipeline_ids:
conditions.append(model.pipeline_id.in_(pipeline_ids))
if start_time is not None:
conditions.append(timestamp >= start_time)
if end_time is not None:
conditions.append(timestamp <= end_time)
statement = (
sqlalchemy.select(time_bucket.label('bucket'), sqlalchemy.func.count(model.id).label('count'))
.where(*conditions)
.group_by(time_bucket)
.order_by(time_bucket)
.limit(MAX_TRAFFIC_POINTS + 1)
)
result = await ap.persistence_mgr.execute_async(statement)
rows = result.all()
truncated = truncated or len(rows) > MAX_TRAFFIC_POINTS
for timestamp_value, count in rows[:MAX_TRAFFIC_POINTS]:
key = (
datetime.datetime.fromisoformat(timestamp_value)
if isinstance(timestamp_value, str)
else timestamp_value
)
points.setdefault(key, {'messages': 0, 'llm_calls': 0})[field] = int(count)
def floor(value: datetime.datetime) -> datetime.datetime:
return value.replace(minute=0, second=0, microsecond=0, **({'hour': 0} if bucket == 'day' else {}))
first = floor(start_time) if start_time is not None else min(points, default=None)
last = floor(end_time) if end_time is not None else max(points, default=None)
series = []
if first is not None and last is not None:
cursor = first
while cursor <= last and len(series) < MAX_TRAFFIC_POINTS:
series.append(
{'timestamp': cursor.isoformat() + 'Z', **points.get(cursor, {'messages': 0, 'llm_calls': 0})}
)
cursor += step
truncated = truncated or cursor <= last
return {'bucket': bucket, 'points': series, 'truncated': truncated}
+50 -129
View File
@@ -1,6 +1,5 @@
from __future__ import annotations
import asyncio
import uuid
import traceback
@@ -8,10 +7,8 @@ import sqlalchemy
from ....cloud.model_catalog import LANGBOT_MODELS_PROVIDER_REQUESTER
from ....core import app
from ....core.task_boundary import create_detached_task
from ....entity.persistence import model as persistence_model
from ....workspace.errors import WorkspaceNotFoundError
from ....provider.modelmgr.codex_auth import CodexAuth, REQUESTER as CODEX_REQUESTER, validate_config
from .secrets import contains_secret_placeholder, redact_secrets, restore_secret_placeholders
from .tenant import TenantContext, require_workspace_uuid, scope_statement
@@ -23,8 +20,6 @@ class ModelProviderService:
def __init__(self, ap: app.Application) -> None:
self.ap = ap
self.codex_auth = CodexAuth(ap)
self._deletion_tasks: set[asyncio.Task[None]] = set()
def _is_cloud_runtime(self) -> bool:
mode = getattr(self.ap.persistence_mgr, 'mode', None)
@@ -121,30 +116,14 @@ class ModelProviderService:
provider_data = provider_data.copy()
if self._system_requester_is_reserved(provider_data.get('requester')):
raise ValueError('space-chat-completions is reserved for the Cloud-managed LangBot Models provider')
validate_config(provider_data)
provider_data['uuid'] = str(uuid.uuid4())
provider_data['workspace_uuid'] = require_workspace_uuid(context)
provider_data['api_keys'] = self._normalize_api_keys(
restore_secret_placeholders(provider_data.get('api_keys'), sensitive=True)
)
if provider_data.get('requester') == CODEX_REQUESTER:
async with self.ap.persistence_mgr.tenant_uow(provider_data['workspace_uuid']):
await self.ap.persistence_mgr.execute_async(
sqlalchemy.insert(persistence_model.ModelProvider).values(**provider_data)
)
await self.ap.persistence_mgr.execute_async(
sqlalchemy.insert(persistence_model.CodexCredential).values(
workspace_uuid=provider_data['workspace_uuid'],
provider_uuid=provider_data['uuid'],
payload={},
version=0,
lease_until=0,
)
)
else:
await self.ap.persistence_mgr.execute_async(
sqlalchemy.insert(persistence_model.ModelProvider).values(**provider_data)
)
await self.ap.persistence_mgr.execute_async(
sqlalchemy.insert(persistence_model.ModelProvider).values(**provider_data)
)
# load to runtime
runtime_provider = await self.ap.model_mgr.load_provider(context, provider_data)
@@ -159,17 +138,6 @@ class ModelProviderService:
raise ValueError('space-chat-completions is reserved for the Cloud-managed LangBot Models provider')
provider_data.pop('uuid', None)
provider_data.pop('workspace_uuid', None)
if {'requester', 'base_url', 'api_keys'} & provider_data.keys():
current = await self.get_provider(context, provider_uuid, include_secret=True)
if current is None:
raise WorkspaceNotFoundError('Provider not found')
if CODEX_REQUESTER in (current.get('requester'), provider_data.get('requester')):
if provider_data.get('requester', current.get('requester')) != current.get('requester'):
raise ValueError('Create a separate provider to change the ChatGPT authentication type')
merged = {**current, **provider_data}
validate_config(merged)
provider_data['base_url'] = merged['base_url']
provider_data['api_keys'] = []
if 'api_keys' in provider_data:
submitted_keys = provider_data.get('api_keys')
if contains_secret_placeholder(submitted_keys, sensitive=True):
@@ -195,107 +163,60 @@ class ModelProviderService:
raise WorkspaceNotFoundError('Provider not found')
await self.ap.model_mgr.reload_provider(context, provider_uuid)
async def delete_provider(self, context: TenantContext, provider_uuid: str, cascade: bool = False) -> None:
"""Delete a provider, optionally deleting all its Workspace-scoped models."""
async def delete_provider(self, context: TenantContext, provider_uuid: str) -> None:
"""Delete a provider (only if no models reference it)"""
await self._assert_provider_mutable(context, provider_uuid)
workspace_uuid = require_workspace_uuid(context)
persistence = self.ap.persistence_mgr
model_types = (
(persistence_model.LLMModel, 'LLM', 'remove_llm_model'),
(persistence_model.EmbeddingModel, 'Embedding', 'remove_embedding_model'),
(persistence_model.RerankModel, 'Rerank', 'remove_rerank_model'),
# Check if any models use this provider
llm_result = await self.ap.persistence_mgr.execute_async(
scope_statement(
sqlalchemy.select(persistence_model.LLMModel).where(
persistence_model.LLMModel.provider_uuid == provider_uuid
),
persistence_model.LLMModel,
workspace_uuid,
)
)
deleted_models: list[tuple[str, list[str]]] = []
async with persistence.tenant_uow(workspace_uuid):
# Check ownership before touching children. Lock the provider on PostgreSQL
# so concurrent model inserts cannot race the reference check/deletion.
provider_result = await persistence.execute_async(
scope_statement(
sqlalchemy.select(persistence_model.ModelProvider.requester)
.where(persistence_model.ModelProvider.uuid == provider_uuid)
.with_for_update(),
persistence_model.ModelProvider,
workspace_uuid,
)
if llm_result.first() is not None:
raise ValueError('Cannot delete provider: LLM models still reference it')
embedding_result = await self.ap.persistence_mgr.execute_async(
scope_statement(
sqlalchemy.select(persistence_model.EmbeddingModel).where(
persistence_model.EmbeddingModel.provider_uuid == provider_uuid
),
persistence_model.EmbeddingModel,
workspace_uuid,
)
provider = provider_result.first()
if provider is None:
raise WorkspaceNotFoundError('Provider not found')
if self._system_requester_is_reserved(provider.requester):
raise ValueError('LangBot Models is managed by Cloud and cannot be modified')
)
if embedding_result.first() is not None:
raise ValueError('Cannot delete provider: Embedding models still reference it')
for model_type, label, remover in model_types:
result = await persistence.execute_async(
scope_statement(
sqlalchemy.select(model_type.uuid).where(model_type.provider_uuid == provider_uuid),
model_type,
workspace_uuid,
)
)
model_uuids = list(result.scalars())
if model_uuids and not cascade:
raise ValueError(f'Cannot delete provider: {label} models still reference it')
if model_uuids:
# Model services have no pipeline/KB deletion side effects: they
# delete the scoped row and evict its runtime cache. Defer eviction
# here rather than calling those services before our commit.
await persistence.execute_async(
scope_statement(
sqlalchemy.delete(model_type).where(model_type.provider_uuid == provider_uuid),
model_type,
workspace_uuid,
)
)
deleted_models.append((remover, model_uuids))
# Explicit cleanup also works on legacy SQLite connections without FK
# enforcement; never load or serialize the private credential payload.
await persistence.execute_async(
scope_statement(
sqlalchemy.delete(persistence_model.CodexCredential).where(
persistence_model.CodexCredential.provider_uuid == provider_uuid
),
persistence_model.CodexCredential,
workspace_uuid,
)
rerank_result = await self.ap.persistence_mgr.execute_async(
scope_statement(
sqlalchemy.select(persistence_model.RerankModel).where(
persistence_model.RerankModel.provider_uuid == provider_uuid
),
persistence_model.RerankModel,
workspace_uuid,
)
result = await persistence.execute_async(
scope_statement(
sqlalchemy.delete(persistence_model.ModelProvider).where(
persistence_model.ModelProvider.uuid == provider_uuid
),
persistence_model.ModelProvider,
workspace_uuid,
)
)
if rerank_result.first() is not None:
raise ValueError('Cannot delete provider: Rerank models still reference it')
result = await self.ap.persistence_mgr.execute_async(
scope_statement(
sqlalchemy.delete(persistence_model.ModelProvider).where(
persistence_model.ModelProvider.uuid == provider_uuid
),
persistence_model.ModelProvider,
workspace_uuid,
)
if result.rowcount == 0:
raise WorkspaceNotFoundError('Provider not found')
)
if getattr(result, 'rowcount', None) == 0:
raise WorkspaceNotFoundError('Provider not found')
async def remove_runtime() -> None:
async with persistence.tenant_scope(workspace_uuid):
for remover, model_uuids in deleted_models:
for model_uuid in model_uuids:
await getattr(self.ap.model_mgr, remover)(context, model_uuid)
# This also closes the requester's HTTP client; models go first.
await self.ap.model_mgr.remove_provider(context, provider_uuid)
if persistence.current_session() is None:
await remove_runtime()
else:
# A nested UoW has not committed yet. Reuse the rollback-cancelled gate
# and detached context boundary instead of evicting uncommitted data.
task = create_detached_task(
remove_runtime(),
after_commit_manager=persistence,
workspace_uuid=workspace_uuid,
)
self._deletion_tasks.add(task)
def completed(task: asyncio.Task[None]) -> None:
self._deletion_tasks.discard(task)
if not task.cancelled() and task.exception() is not None:
self.ap.logger.error('Failed to remove deleted provider runtime', exc_info=task.exception())
task.add_done_callback(completed)
await self.ap.model_mgr.remove_provider(context, provider_uuid)
async def get_provider_model_counts(self, context: TenantContext, provider_uuid: str) -> dict:
"""Get count of models using this provider"""
+1 -80
View File
@@ -11,9 +11,6 @@ import sqlalchemy
from ....core import app
from ....entity.persistence import user
from ....entity.dto.space_model import SpaceModel
from ....entity.dto.space_model import SpaceModelSelection
from ....entity.persistence import model as persistence_model
from ....cloud.model_catalog import LANGBOT_MODELS_PROVIDER_REQUESTER
_CREDITS_CACHE_TTL_SECONDS = 60
@@ -119,7 +116,7 @@ class SpaceService:
space_config = self._get_space_config()
authorize_url = space_config['oauth_authorize_url']
params = {'redirect_uri': redirect_uri, 'code_contract': 'redirect-v1'}
params = {'redirect_uri': redirect_uri}
if state:
params['state'] = state
return f'{authorize_url}?{urlencode(params)}'
@@ -129,8 +126,6 @@ class SpaceService:
code: str,
workspace_uuids: list[str] | None = None,
workspace_created_ats: dict[str, int] | None = None,
*,
redirect_uri: str = '',
) -> typing.Dict:
"""Exchange OAuth authorization code for tokens"""
from langbot.pkg.utils import constants
@@ -143,7 +138,6 @@ class SpaceService:
f'{space_url}/api/v1/accounts/oauth/token',
json={
'code': code,
'redirect_uri': redirect_uri,
'instance_id': constants.instance_id,
# Sending an explicit empty list tells new Space servers not to
# synthesize a legacy instance-derived Workspace binding.
@@ -244,76 +238,3 @@ class SpaceService:
raise ValueError(f'Failed to get models: {data.get("msg")}')
models_data = data.get('data', {}).get('models', [])
return [SpaceModel.model_validate(model_dict) for model_dict in models_data]
async def get_model_selection(self, category: str) -> typing.List[SpaceModelSelection]:
"""Return Space models in the availability-ranked selection order."""
space_url = self._get_space_config()['url']
session = httpclient.get_session()
async with session.get(
f'{space_url}/api/v1/models/selection',
params={'category': category},
) as response:
if response.status != 200:
error = await httpclient.read_text_limited(response)
raise ValueError(f'Failed to get model selection: {error}')
payload = await httpclient.read_json_limited(response)
if payload.get('code') != 0:
raise ValueError(f'Failed to get model selection: {payload.get("msg")}')
data = payload.get('data', [])
if isinstance(data, dict):
data = data.get('models', data.get('items', []))
if not isinstance(data, list):
raise ValueError('Failed to get model selection: invalid response')
models = []
for selection in data:
if isinstance(selection, dict) and isinstance(selection.get('model'), dict):
models.append(selection['model'])
else:
models.append(selection)
return [SpaceModelSelection.model_validate(model) for model in models]
async def get_recommended_chat_model(self, context: typing.Any) -> dict:
"""Resolve Space's first ranked chat model to a local Workspace model."""
selection = await self.get_model_selection('chat')
if not selection:
raise ValueError('No recommended chat model is available')
recommended = selection[0]
async def find_local_model():
result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(persistence_model.LLMModel)
.join(
persistence_model.ModelProvider,
sqlalchemy.and_(
persistence_model.ModelProvider.workspace_uuid == persistence_model.LLMModel.workspace_uuid,
persistence_model.ModelProvider.uuid == persistence_model.LLMModel.provider_uuid,
),
)
.where(
persistence_model.LLMModel.workspace_uuid == context.workspace_uuid,
persistence_model.ModelProvider.requester == LANGBOT_MODELS_PROVIDER_REQUESTER,
sqlalchemy.or_(
persistence_model.LLMModel.uuid == recommended.uuid,
persistence_model.LLMModel.name == recommended.model_id,
),
)
)
return result.first()
local_model = await find_local_model()
if local_model is None:
# OSS synchronizes the public catalog locally. Refresh once in case
# the recommendation was published after this process started.
from ..context import ExecutionContext
try:
await self.ap.model_mgr.sync_new_models_from_space(ExecutionContext.from_request(context))
except Exception:
pass
local_model = await find_local_model()
if local_model is None:
raise ValueError('Recommended chat model is not available in this Workspace')
return {'uuid': local_model.uuid, 'name': local_model.name}
+6 -340
View File
@@ -4,7 +4,6 @@ import sqlalchemy
import argon2
import jwt
import datetime
import json
import typing
import asyncio
import dataclasses
@@ -13,19 +12,10 @@ import hashlib
import secrets
import time
import uuid
import webauthn
from webauthn.helpers import bytes_to_base64url, base64url_to_bytes
from webauthn.helpers.structs import (
AuthenticatorSelectionCriteria,
PublicKeyCredentialDescriptor,
ResidentKeyRequirement,
UserVerificationRequirement,
)
from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker
from ....entity.persistence import user
from ....entity.persistence import passkey
from ....entity.persistence.workspace import MembershipRole, MembershipStatus, WorkspaceMembership
from ....utils import constants
from ....entity.errors import account as account_errors
@@ -39,9 +29,6 @@ if typing.TYPE_CHECKING:
_SPACE_OAUTH_STATE_MAX_ENTRIES = 4096
_SPACE_OAUTH_STATE_HEAP_COMPACT_FLOOR = 64
_SPACE_OAUTH_STATE_HEAP_MAX_MULTIPLIER = 4
_PASSKEY_CHALLENGE_MAX_ENTRIES = 4096
_PASSKEY_CHALLENGE_HEAP_COMPACT_FLOOR = 64
_PASSKEY_CHALLENGE_HEAP_MAX_MULTIPLIER = 4
class AccountExistsLoginRequiredError(ValueError):
@@ -67,17 +54,6 @@ class SpaceOAuthStateConsumption:
launch_workspace_uuid: str | None = None
@dataclasses.dataclass(frozen=True, slots=True)
class PasskeyChallengeData:
challenge: bytes
purpose: typing.Literal['register', 'auth']
rp_id: str
origin: str
expires_at: float
account_uuid: str | None = None
user_email: str | None = None
class UserService:
ap: Application
_create_user_lock: asyncio.Lock
@@ -89,9 +65,6 @@ class UserService:
self._space_oauth_state_lock = asyncio.Lock()
self._space_oauth_states: dict[str, tuple[str, str | None, float, str | None]] = {}
self._space_oauth_state_expiry_heap: list[tuple[float, str]] = []
self._passkey_challenge_lock = asyncio.Lock()
self._passkey_challenges: dict[str, PasskeyChallengeData] = {}
self._passkey_challenge_expiry_heap: list[tuple[float, str]] = []
@staticmethod
def _space_oauth_state_digest(state: str) -> str:
@@ -141,7 +114,7 @@ class UserService:
if purpose == 'login' and account_uuid is not None:
raise ValueError('Login state cannot be bound to an Account')
if purpose != 'login' and launch_workspace_uuid is not None:
raise ValueError('Launch Workspace state is only valid for LangBot Account login')
raise ValueError('Launch Workspace state is only valid for Space login')
if ttl_seconds <= 0:
raise ValueError('OAuth state lifetime must be positive')
@@ -354,7 +327,7 @@ class UserService:
normalized_email = normalize_email(user_email)
if self._uses_control_plane_directory():
raise ControlPlaneDirectoryRequiredError(
'Cloud invitation registration must use a LangBot Account to preserve control-plane identity'
'Cloud invitation registration must use a Space account to preserve control-plane identity'
)
invitation, _ = await self.ap.workspace_collaboration_service.inspect_invitation(invitation_token)
if invitation.normalized_email != normalized_email:
@@ -421,7 +394,7 @@ class UserService:
# Check if this user has a local password set
if not user_obj.password:
raise ValueError('请使用 LangBot登录')
raise ValueError('请使用 Space登录')
await self._verify_password(user_obj.password, password)
@@ -801,7 +774,7 @@ class UserService:
f'email:{normalized_email}',
)
async def bind_space_account(self, user_email: str, code: str, *, redirect_uri: str = '') -> user.User:
async def bind_space_account(self, user_email: str, code: str) -> user.User:
"""Bind Space account to existing local account"""
local_account = await self.get_user_by_email(user_email)
if local_account is None:
@@ -821,13 +794,12 @@ class UserService:
code,
[binding.workspace_uuid],
{binding.workspace_uuid: created_ts},
redirect_uri=redirect_uri,
)
else:
# Compatibility for early/bootstrap call sites that have not wired
# WorkspaceService yet; old Space servers still derive the legacy
# Workspace identity from instance_id when the field is omitted.
token_data = await self.ap.space_service.exchange_oauth_code(code, redirect_uri=redirect_uri)
token_data = await self.ap.space_service.exchange_oauth_code(code)
access_token = token_data.get('access_token')
refresh_token = token_data.get('refresh_token')
expires_in = token_data.get('expires_in', 0)
@@ -853,7 +825,7 @@ class UserService:
# Check if this Space account is already bound to another user
existing_space_user = await self.get_user_by_space_account_uuid(space_account_uuid)
if existing_space_user and existing_space_user.normalized_email != normalize_email(user_email):
raise ValueError('This LangBot Account is already bound to another user')
raise ValueError('This Space account is already bound to another user')
# Update local account to Space account
normalized_email = normalize_email(user_email)
@@ -877,309 +849,3 @@ class UserService:
await self._update_space_provider_for_account(local_account, api_key)
return await self.get_user_by_email(space_email)
def _prune_passkey_challenges(self, now: float) -> None:
while self._passkey_challenge_expiry_heap:
expires_at, token = self._passkey_challenge_expiry_heap[0]
entry = self._passkey_challenges.get(token)
if entry is None or entry.expires_at != expires_at:
heapq.heappop(self._passkey_challenge_expiry_heap)
continue
if expires_at > now:
break
heapq.heappop(self._passkey_challenge_expiry_heap)
self._passkey_challenges.pop(token, None)
max_heap_entries = max(
_PASSKEY_CHALLENGE_HEAP_COMPACT_FLOOR,
len(self._passkey_challenges) * _PASSKEY_CHALLENGE_HEAP_MAX_MULTIPLIER,
)
if len(self._passkey_challenge_expiry_heap) > max_heap_entries:
self._passkey_challenge_expiry_heap[:] = [
(entry.expires_at, token) for token, entry in self._passkey_challenges.items()
]
heapq.heapify(self._passkey_challenge_expiry_heap)
async def issue_passkey_challenge(
self,
purpose: typing.Literal['register', 'auth'],
rp_id: str,
origin: str,
*,
account_uuid: str | None = None,
user_email: str | None = None,
ttl_seconds: int = 300,
) -> tuple[str, bytes]:
now = time.monotonic()
challenge_bytes = secrets.token_bytes(32)
challenge_token = secrets.token_urlsafe(32)
expires_at = now + ttl_seconds
async with self._passkey_challenge_lock:
self._prune_passkey_challenges(now)
while len(self._passkey_challenges) >= _PASSKEY_CHALLENGE_MAX_ENTRIES:
if not self._passkey_challenge_expiry_heap:
break
_, oldest_token = heapq.heappop(self._passkey_challenge_expiry_heap)
self._passkey_challenges.pop(oldest_token, None)
self._passkey_challenges[challenge_token] = PasskeyChallengeData(
challenge=challenge_bytes,
purpose=purpose,
rp_id=rp_id,
origin=origin,
expires_at=expires_at,
account_uuid=account_uuid,
user_email=user_email,
)
heapq.heappush(self._passkey_challenge_expiry_heap, (expires_at, challenge_token))
return challenge_token, challenge_bytes
async def consume_passkey_challenge(
self,
challenge_token: str,
purpose: typing.Literal['register', 'auth'],
) -> PasskeyChallengeData:
now = time.monotonic()
async with self._passkey_challenge_lock:
self._prune_passkey_challenges(now)
data = self._passkey_challenges.pop(challenge_token, None)
if data is None or data.expires_at < now:
raise ValueError('Invalid or expired passkey challenge')
if data.purpose != purpose:
raise ValueError('Passkey challenge purpose mismatch')
return data
async def get_user_passkeys(self, account_uuid: str) -> list[passkey.PasskeyCredential]:
statement = (
sqlalchemy.select(passkey.PasskeyCredential)
.where(passkey.PasskeyCredential.account_uuid == account_uuid)
.order_by(passkey.PasskeyCredential.created_at.desc())
)
async with self._session_factory()() as session:
result = await session.scalars(statement)
return list(result.all())
async def get_passkey_by_credential_id(self, credential_id: str) -> passkey.PasskeyCredential | None:
statement = sqlalchemy.select(passkey.PasskeyCredential).where(
passkey.PasskeyCredential.credential_id == credential_id
)
async with self._session_factory()() as session:
return await session.scalar(statement)
async def get_passkey_by_uuid(self, passkey_uuid: str) -> passkey.PasskeyCredential | None:
statement = sqlalchemy.select(passkey.PasskeyCredential).where(passkey.PasskeyCredential.uuid == passkey_uuid)
async with self._session_factory()() as session:
return await session.scalar(statement)
async def generate_passkey_registration_options(
self,
account_uuid: str,
rp_id: str,
origin: str,
rp_name: str = 'LangBot',
) -> tuple[dict[str, typing.Any], str]:
account = await self.get_user_by_uuid(account_uuid)
if account is None:
raise ValueError('User not found')
self._require_active_account(account)
challenge_token, challenge_bytes = await self.issue_passkey_challenge(
purpose='register',
rp_id=rp_id,
origin=origin,
account_uuid=account_uuid,
user_email=account.user,
)
existing_passkeys = await self.get_user_passkeys(account_uuid)
exclude_credentials = [
PublicKeyCredentialDescriptor(id=base64url_to_bytes(pk.credential_id)) for pk in existing_passkeys
]
options = webauthn.generate_registration_options(
rp_id=rp_id,
rp_name=rp_name,
user_name=account.user,
user_id=account.uuid.encode('utf-8'),
user_display_name=account.user,
challenge=challenge_bytes,
exclude_credentials=exclude_credentials or None,
authenticator_selection=AuthenticatorSelectionCriteria(
resident_key=ResidentKeyRequirement.PREFERRED,
),
)
options_dict = json.loads(webauthn.options_to_json(options))
return options_dict, challenge_token
async def verify_and_save_passkey_registration(
self,
challenge_token: str,
credential_data: dict[str, typing.Any] | str,
name: str | None = None,
) -> passkey.PasskeyCredential:
challenge_data = await self.consume_passkey_challenge(challenge_token, 'register')
if not challenge_data.account_uuid:
raise ValueError('Registration challenge must be bound to an account')
verification = webauthn.verify_registration_response(
credential=credential_data,
expected_challenge=challenge_data.challenge,
expected_rp_id=challenge_data.rp_id,
expected_origin=challenge_data.origin,
require_user_verification=False,
)
cred_id_str = bytes_to_base64url(verification.credential_id)
pub_key_str = bytes_to_base64url(verification.credential_public_key)
transports = None
if isinstance(credential_data, dict):
resp = credential_data.get('response', {})
if isinstance(resp, dict) and 'transports' in resp:
t_list = resp.get('transports')
if isinstance(t_list, list):
transports = ','.join(str(x) for x in t_list)
credential_name = (name or '').strip()
if not credential_name:
credential_name = f'Passkey ({datetime.datetime.now().strftime("%Y-%m-%d %H:%M")})'
record = passkey.PasskeyCredential(
uuid=str(uuid.uuid4()),
account_uuid=challenge_data.account_uuid,
name=credential_name,
credential_id=cred_id_str,
public_key=pub_key_str,
sign_count=verification.sign_count,
aaguid=verification.aaguid,
transports=transports,
backed_up=verification.credential_backed_up,
)
async with self._session_factory()() as session:
async with session.begin():
session.add(record)
await session.flush()
await session.refresh(record)
return record
async def generate_passkey_authentication_options(
self,
rp_id: str,
origin: str,
email: str | None = None,
) -> tuple[dict[str, typing.Any], str]:
challenge_token, challenge_bytes = await self.issue_passkey_challenge(
purpose='auth',
rp_id=rp_id,
origin=origin,
user_email=email,
)
allow_credentials: list[PublicKeyCredentialDescriptor] | None = None
if email:
user_obj = await self.get_user_by_email(email)
if user_obj:
user_passkeys = await self.get_user_passkeys(user_obj.uuid)
if user_passkeys:
allow_credentials = [
PublicKeyCredentialDescriptor(id=base64url_to_bytes(pk.credential_id)) for pk in user_passkeys
]
options = webauthn.generate_authentication_options(
rp_id=rp_id,
challenge=challenge_bytes,
allow_credentials=allow_credentials or None,
user_verification=UserVerificationRequirement.PREFERRED,
)
options_dict = json.loads(webauthn.options_to_json(options))
return options_dict, challenge_token
async def verify_passkey_authentication(
self,
challenge_token: str,
credential_data: dict[str, typing.Any] | str,
) -> tuple[str, user.User]:
challenge_data = await self.consume_passkey_challenge(challenge_token, 'auth')
raw_id = credential_data.get('id') if isinstance(credential_data, dict) else None
if not raw_id:
raise ValueError('Missing credential id')
stored_credential = await self.get_passkey_by_credential_id(raw_id)
if stored_credential is None:
raise ValueError('Passkey credential not recognized')
user_obj = await self.get_user_by_uuid(stored_credential.account_uuid)
if user_obj is None:
raise ValueError('Associated user not found')
self._require_active_account(user_obj)
verification = webauthn.verify_authentication_response(
credential=credential_data,
expected_challenge=challenge_data.challenge,
expected_rp_id=challenge_data.rp_id,
expected_origin=challenge_data.origin,
credential_public_key=base64url_to_bytes(stored_credential.public_key),
credential_current_sign_count=stored_credential.sign_count,
require_user_verification=False,
)
async with self._session_factory()() as session:
async with session.begin():
record = await session.scalar(
sqlalchemy.select(passkey.PasskeyCredential).where(
passkey.PasskeyCredential.id == stored_credential.id
)
)
if record:
record.sign_count = verification.new_sign_count
record.last_used_at = datetime.datetime.now()
record.backed_up = verification.credential_backed_up
token = await self.generate_jwt_token(user_obj)
return token, user_obj
async def rename_user_passkey(
self,
account_uuid: str,
passkey_uuid: str,
new_name: str,
) -> passkey.PasskeyCredential | None:
async with self._session_factory()() as session:
async with session.begin():
record = await session.scalar(
sqlalchemy.select(passkey.PasskeyCredential).where(
passkey.PasskeyCredential.uuid == passkey_uuid,
passkey.PasskeyCredential.account_uuid == account_uuid,
)
)
if record is None:
return None
record.name = new_name
await session.flush()
await session.refresh(record)
return record
async def delete_user_passkey(
self,
account_uuid: str,
passkey_uuid: str,
) -> bool:
async with self._session_factory()() as session:
async with session.begin():
record = await session.scalar(
sqlalchemy.select(passkey.PasskeyCredential).where(
passkey.PasskeyCredential.uuid == passkey_uuid,
passkey.PasskeyCredential.account_uuid == account_uuid,
)
)
if record is None:
return False
await session.delete(record)
return True
+1 -10
View File
@@ -147,16 +147,7 @@ class LangBotMCPServer:
)
async def create_pipeline(pipeline_data: dict) -> str:
context = _authorized(Permission.RESOURCE_MANAGE)
create_as_default = pipeline_data.get('is_default') is True
return _dump(
{
'uuid': await ap.pipeline_service.create_pipeline(
context,
pipeline_data,
default=create_as_default,
)
}
)
return _dump({'uuid': await ap.pipeline_service.create_pipeline(context, pipeline_data)})
@mcp.tool(description='Update a pipeline by UUID. `pipeline_data` matches the PUT body.')
async def update_pipeline(pipeline_uuid: str, pipeline_data: dict) -> str:
-4
View File
@@ -368,10 +368,6 @@ class BoxRuntimeConnector(ManagedRuntimeConnector):
if not self._control_token and allow_generate:
self._control_token = secrets.token_urlsafe(48)
if not self._control_token:
if getattr(getattr(self.ap, 'deployment', None), 'mode', 'oss') == 'cloud':
raise BoxRuntimeUnavailableError(
f'{BOX_CONTROL_TOKEN_ENV} must be configured with a strong shared secret for a Cloud Box runtime'
)
return ''
try:
self._control_token = validate_control_token(self._control_token)
+11 -27
View File
@@ -455,9 +455,7 @@ class BoxService:
async def _require_validated_workspace_sandbox(self, execution_context: ExecutionContext) -> None:
if not self._available:
raise BoxError(
'Box runtime is not available. Configure an available Box backend before using Box features.'
)
raise BoxError('Box runtime is not available. Install and start Docker to use sandbox features.')
if self._cloud_managed:
if self._admission is None:
raise BoxAdmissionError('Cloud Box sandbox admission is unavailable')
@@ -567,9 +565,7 @@ class BoxService:
skip_host_mount_validation: bool = False,
) -> dict:
if not self._available:
raise BoxError(
'Box runtime is not available. Configure an available Box backend before using Box features.'
)
raise BoxError('Box runtime is not available. Install and start Docker to use sandbox features.')
execution_context = await self._validated_execution_context(self._query_execution_context(query))
spec_payload = self._managed_policy_payload(execution_context, spec_payload)
await self._require_validated_workspace_sandbox(execution_context)
@@ -1214,9 +1210,8 @@ class BoxService:
async def _read_outbox_via_exec(self, query: pipeline_query.Query) -> list[dict]:
"""Fallback: read the outbox over the exec channel (E2B / remote).
Uses ``client.execute`` directly (bypassing ``_serialize_result``)
so stdout is NOT truncated by ``output_limit_chars`` - the raw
base64 payload can be far larger than the 4000-char display limit.
Note: exec stdout is truncated by ``output_limit_chars``, so this path
only reliably transfers small files. The host path is preferred.
"""
import json as _json
@@ -1270,22 +1265,14 @@ class BoxService:
' break\n'
'print(json.dumps(out))\n'
)
spec_payload: dict = {
'cmd': f"python3 - <<'LBPY'\n{script}\nLBPY",
'timeout_sec': 120,
'session_id': self.resolve_box_session_id(query),
}
if 'extra_mounts' not in spec_payload:
spec_payload['extra_mounts'] = self.build_skill_extra_mounts(query)
try:
spec = self.build_spec(spec_payload)
result = await self.client.execute(spec)
except Exception:
return []
if not result.ok:
result = await self.execute_tool(
{'command': f"python3 - <<'LBPY'\n{script}\nLBPY", 'timeout_sec': 120},
query,
)
if not result.get('ok'):
return []
try:
return _json.loads(str(result.stdout or '').strip().splitlines()[-1])
return _json.loads(str(result.get('stdout') or '').strip().splitlines()[-1])
except Exception:
return []
@@ -2146,8 +2133,5 @@ class BoxService:
if backend_name:
payload['connector_error'] = f'Configured sandbox backend "{backend_name}" is unavailable'
else:
payload['connector_error'] = (
'No supported sandbox backend (Docker / nsjail / E2B) is available. '
'Trusted local development may explicitly select the unsafe host backend.'
)
payload['connector_error'] = 'No supported sandbox backend (Docker / nsjail / E2B) is available'
return payload
+2 -100
View File
@@ -125,21 +125,10 @@ class DirectoryProjectionService:
# The database cursor remains the shared projection high-water mark,
# while this cursor tracks what this process has actually observed.
self._consumer_cursor: int | None = None
self._sync_lock = asyncio.Lock()
async def initialize(self) -> None:
"""Block Cloud startup until one full signed snapshot is committed."""
async with self._sync_lock:
await self._refresh_snapshot()
async def refresh_snapshot(self) -> None:
"""Refresh from one full signed snapshot within the sync single-flight."""
async with self._sync_lock:
await self._refresh_snapshot()
async def _refresh_snapshot(self) -> None:
last_superseded: _DirectorySnapshotSuperseded | None = None
for _attempt in range(5):
snapshot = await self.provider.fetch_snapshot(self.instance_uuid)
@@ -170,84 +159,9 @@ class DirectoryProjectionService:
delay = min(max(delay * 2, self.sync_interval_seconds), self.max_staleness_seconds / 2)
async def sync_once(self) -> None:
async with self._sync_lock:
await self._sync_once()
async def reconcile_workspaces(self, workspace_uuids: Iterable[str]) -> None:
"""Synchronously project an exact Workspace set without moving the event cursor."""
requested = tuple(sorted({str(value).strip() for value in workspace_uuids if str(value).strip()}))
if not requested:
raise DirectoryProjectionUnavailableError('Targeted directory reconciliation requires a Workspace')
if len(requested) > self.event_limit:
raise DirectoryProjectionUnavailableError('Targeted directory reconciliation exceeds the batch limit')
async with self._sync_lock:
delta = await self.provider.fetch_workspaces(self.instance_uuid, requested)
await self._apply_targeted_delta(delta, requested)
async def _apply_targeted_delta(
self,
delta: DirectoryDelta,
requested_workspace_uuids: tuple[str, ...],
) -> None:
if not isinstance(delta, DirectoryDelta):
raise DirectoryProjectionUnavailableError('Directory provider returned an invalid delta')
workspace_count, membership_count = self._validate_batch_capacity(
delta.workspaces,
full_snapshot=False,
)
delta = DirectoryDelta.model_validate(delta.model_dump())
if delta.instance_uuid != self.instance_uuid:
raise DirectoryProjectionUnavailableError('Directory delta targets another LangBot instance')
requested = set(requested_workspace_uuids)
if set(delta.requested_workspace_uuids) != requested:
raise DirectoryProjectionUnavailableError('Directory delta does not match the requested Workspaces')
if {workspace.uuid for workspace in delta.workspaces} != requested:
raise DirectoryProjectionUnavailableError('Directory delta omitted a requested Workspace')
directory_uow = getattr(self.ap.persistence_mgr, 'directory_projection_uow', None)
if not callable(directory_uow):
raise DirectoryProjectionUnavailableError('Directory projection persistence scope is unavailable')
async with directory_uow(self.instance_uuid) as uow:
session = uow.session
state = await session.scalar(
sqlalchemy.select(DirectoryProjectionState)
.where(DirectoryProjectionState.instance_uuid == self.instance_uuid)
.with_for_update()
)
if state is None:
raise DirectoryProjectionUnavailableError('Directory projection is not initialized')
snapshot = DirectorySnapshot(
instance_uuid=self.instance_uuid,
cursor=state.cursor,
generated_at=delta.generated_at,
workspaces=delta.workspaces,
)
accounts_by_uuid = await self._apply_accounts(session, snapshot, preserve_existing=True)
await self._apply_workspaces(session, snapshot, accounts_by_uuid=accounts_by_uuid)
active_workspace_count = await self._enforce_active_workspace_capacity(session)
await session.flush()
await self._update_entitlement_workspace_activity(
snapshot.workspaces,
requested_workspace_uuids=requested,
)
self._publish_runtime_execution_projection(
snapshot.workspaces,
affected_workspace_uuids=requested,
)
self._request_model_catalog_sync()
self._record_batch_cardinality(
active_workspaces=active_workspace_count,
workspaces=workspace_count,
memberships=membership_count,
)
async def _sync_once(self) -> None:
cursor = self._consumer_cursor
if cursor is None:
await self._refresh_snapshot()
await self.initialize()
return
batch = await self.provider.fetch_events(
self.instance_uuid,
@@ -794,13 +708,7 @@ class DirectoryProjectionService:
for row in inbox_rows:
row.applied_at = now
async def _apply_accounts(
self,
session: Any,
snapshot: DirectorySnapshot,
*,
preserve_existing: bool = False,
) -> dict[str, User]:
async def _apply_accounts(self, session: Any, snapshot: DirectorySnapshot) -> dict[str, User]:
selected: dict[str, DirectoryMember] = {}
emails: dict[str, str] = {}
for workspace in snapshot.workspaces:
@@ -865,12 +773,6 @@ class DirectoryProjectionService:
continue
if account.source != AccountSource.CLOUD_PROJECTION.value:
raise DirectoryProjectionUnavailableError('Directory account UUID collides with a local Core account')
if preserve_existing:
# A targeted Workspace fetch has no independently monotonic
# Account revision. It may create a missing runtime shadow, but
# ordered event/snapshot projection remains the only updater of
# existing Account identity and status fields.
continue
if account.projection_revision > snapshot.cursor:
raise DirectoryProjectionUnavailableError('Directory account revision rolled back')
projected_account = self._account_projection(member)
+1 -9
View File
@@ -3,7 +3,6 @@ from __future__ import annotations
import typing
import inspect
from ..api.http.context import ExecutionContext
from ..core import app
from . import operator
from ..utils import importutil
@@ -67,14 +66,7 @@ class CommandManager:
require_context = getattr(self.ap.plugin_connector, 'require_workspace_context', None)
if require_context is not None:
result = require_context(
ExecutionContext(
instance_uuid=context.instance_uuid,
workspace_uuid=context.workspace_uuid,
placement_generation=context.placement_generation,
query_uuid=context.query_uuid,
)
)
result = require_context(context)
if inspect.isawaitable(result):
await result
+5 -39
View File
@@ -249,10 +249,6 @@ class Application:
{},
)
),
'plugin_runtime_connected': bool(
self.plugin_connector is not None
and getattr(self.plugin_connector, '_runtime_available', lambda: False)()
),
}
mcp_loader = getattr(self.tool_mgr, 'mcp_tool_loader', None)
runtime_stats.update(
@@ -301,36 +297,11 @@ class Application:
async def initialize(self):
pass
async def _initialize_plugin_runtime(self) -> None:
try:
await self.plugin_connector.initialize()
except asyncio.CancelledError:
raise
except Exception as exc:
self.logger.warning(f'Plugin runtime unavailable during startup; reconnecting in background: {exc}')
self.plugin_connector.schedule_reconnect()
def _start_plugin_runtime_initialization(self) -> asyncio.Task | None:
task = getattr(self, '_plugin_runtime_initialization_task', None)
if task is not None and not task.done():
return task
# This is application lifecycle work, not a request side effect. It must
# not wait on PersistenceManager's after-commit gate at boot.
task = asyncio.create_task(
self._initialize_plugin_runtime(),
name='plugin-runtime-initialization',
)
self._plugin_runtime_initialization_task = task
return task
async def run(self):
self.event_loop_monitor.start()
try:
if (
self.directory_projection_service is not None
and getattr(self, 'directory_projection_task', None) is None
):
self.directory_projection_task = self.task_mgr.create_task(
if self.directory_projection_service is not None:
self.task_mgr.create_task(
self.directory_projection_service.run(),
name='cloud-directory-projection',
scopes=[core_entities.LifecycleControlScope.APPLICATION],
@@ -347,6 +318,7 @@ class Application:
name='cloud-manifest-refresh',
scopes=[core_entities.LifecycleControlScope.APPLICATION],
)
await self.plugin_connector.initialize_plugins()
# 后续可能会允许动态重启其他任务
# 故为了防止程序在非 Ctrl-C 情况下退出,这里创建一个不会结束的协程
@@ -372,7 +344,6 @@ class Application:
name='http-api-controller',
scopes=[core_entities.LifecycleControlScope.APPLICATION],
)
self._start_plugin_runtime_initialization()
# Telemetry instance heartbeat (startup + daily); respects
# space.disable_telemetry via TelemetryManager.send().
@@ -554,11 +525,6 @@ class Application:
if self.task_mgr is not None:
self.task_mgr.cancel_by_scope(core_entities.LifecycleControlScope.APPLICATION)
plugin_runtime_task = getattr(self, '_plugin_runtime_initialization_task', None)
if plugin_runtime_task is not None and not plugin_runtime_task.done():
plugin_runtime_task.cancel()
with contextlib.suppress(asyncio.CancelledError):
await plugin_runtime_task
with contextlib.suppress(Exception):
await self.event_loop_monitor.stop()
mcp_mount = getattr(self.http_ctrl, 'mcp_mount', None)
@@ -635,9 +601,9 @@ class Application:
frontend_path = paths.get_frontend_path()
if not os.path.exists(frontend_path):
self.logger.warning('WebUI 文件缺失,请根据文档部署:https://langbot.app/docs/zh')
self.logger.warning('WebUI 文件缺失,请根据文档部署:https://docs.langbot.app/zh')
self.logger.warning(
'WebUI files are missing, please deploy according to the documentation: https://langbot.app/docs/en'
'WebUI files are missing, please deploy according to the documentation: https://docs.langbot.app/en'
)
return
+8 -11
View File
@@ -1,6 +1,6 @@
from __future__ import annotations
from .. import stage, app, entities as core_entities
from .. import stage, app
from ...utils import version, proxy, constants
from ...pipeline import pool, controller, pipelinemgr
from ...pipeline import aggregator as message_aggregator
@@ -292,17 +292,14 @@ class BuildAppStage(stage.BootingStage):
async def runtime_disconnect_callback(connector: plugin_connector.PluginRuntimeConnector) -> None:
connector.schedule_reconnect()
if ap.directory_projection_service is not None:
# Keep the projection fresh while shared Runtime cold restore runs.
# BuildApp initializes the connector before Application.run() starts
# its long-lived tasks, so start the single refresh task here.
ap.directory_projection_task = ap.task_mgr.create_task(
ap.directory_projection_service.run(),
name='cloud-directory-projection',
scopes=[core_entities.LifecycleControlScope.APPLICATION],
)
plugin_connector_inst = plugin_connector.PluginRuntimeConnector(ap, runtime_disconnect_callback)
try:
await plugin_connector_inst.initialize()
except Exception as exc:
# Keep the API/UI available while an external or managed runtime is
# starting, then recover in the background with bounded backoff.
ap.logger.warning(f'Plugin runtime unavailable during startup; reconnecting in background: {exc}')
plugin_connector_inst.schedule_reconnect()
ap.plugin_connector = plugin_connector_inst
workspace_service_inst.release_startup_execution_bindings()
+1 -20
View File
@@ -1,18 +1,9 @@
from __future__ import annotations
import logging
import secrets
from .. import stage, app
# This stage runs before SetupLoggerStage, so ap.logger is still None here;
# the module logger falls back to the stderr lastResort handler.
_logger = logging.getLogger(__name__)
# 32 symbols without 0/O or 1/I; eight independent draws provide 40 random bits.
_RECOVERY_KEY_ALPHABET = '23456789ABCDEFGHJKLMNPQRSTUVWXYZ'
_RECOVERY_KEY_LENGTH = 8
@stage.stage_class('GenKeysStage')
class GenKeysStage(stage.BootingStage):
@@ -29,15 +20,5 @@ class GenKeysStage(stage.BootingStage):
ap.instance_config.data['system']['recovery_key'] = ''
if not ap.instance_config.data['system']['recovery_key']:
# Keep recovery practical to type. Security also requires the reset
# endpoint's concurrency-safe quota (five admissions per 15 minutes).
ap.instance_config.data['system']['recovery_key'] = ''.join(
secrets.choice(_RECOVERY_KEY_ALPHABET) for _ in range(_RECOVERY_KEY_LENGTH)
)
ap.instance_config.data['system']['recovery_key'] = secrets.token_hex(3).upper()
await ap.instance_config.dump_config()
elif len(ap.instance_config.data['system']['recovery_key']) < _RECOVERY_KEY_LENGTH:
_logger.warning(
'Low-entropy legacy recovery key detected (length < 8); '
'regenerate system.recovery_key in the configuration file '
'with a strong random value (#2392)'
)
+12 -49
View File
@@ -1,7 +1,6 @@
from __future__ import annotations
import asyncio
import json
import typing
import datetime
import time
@@ -198,41 +197,6 @@ class TaskWrapper:
},
}
def to_public_dict(self) -> dict:
"""Return the stable task projection exposed to API-key callers."""
if self.task.cancelled():
status = 'cancelled'
error = {'type': 'task_cancelled', 'message': 'Task was cancelled'}
result = None
elif not self.task.done():
status = 'running'
error = None
result = None
else:
exception = self.assume_exception()
if exception is not None:
status = 'failed'
error = {'type': 'task_failed', 'message': 'Task execution failed'}
result = None
else:
status = 'succeeded'
error = None
result = self.assume_result()
try:
json.dumps(result)
except (TypeError, ValueError):
result = None
return {
'id': self.id,
'task_type': self.task_type,
'kind': self.kind,
'status': status,
'error': error,
'result': result,
'created_at': self.created_at,
}
def cancel(self):
self.task.cancel()
@@ -361,20 +325,19 @@ class AsyncTaskManager:
instance_uuid: str | None = None,
workspace_uuid: str | None = None,
placement_generation: int | None = None,
public: bool = False,
) -> dict:
tasks = [
t.to_public_dict() if public else t.to_dict()
for t in self.tasks
if (type is None or t.task_type == type)
and (kind is None or t.kind == kind)
and (instance_uuid is None or t.instance_uuid == instance_uuid)
and (workspace_uuid is None or t.workspace_uuid == workspace_uuid)
and (placement_generation is None or t.placement_generation == placement_generation)
]
if public:
return {'tasks': tasks}
return {'tasks': tasks, 'id_index': TaskWrapper._id_index}
return {
'tasks': [
t.to_dict()
for t in self.tasks
if (type is None or t.task_type == type)
and (kind is None or t.kind == kind)
and (instance_uuid is None or t.instance_uuid == instance_uuid)
and (workspace_uuid is None or t.workspace_uuid == workspace_uuid)
and (placement_generation is None or t.placement_generation == placement_generation)
],
'id_index': TaskWrapper._id_index,
}
def get_stats(self) -> dict:
completed = sum(1 for t in self.tasks if t.task.done())
@@ -47,10 +47,3 @@ class SpaceModel(pydantic.BaseModel):
status: str
created_at: str | None = None
updated_at: str | None = None
class SpaceModelSelection(pydantic.BaseModel):
"""Minimal model identity returned by the ranked selection endpoint."""
uuid: str
model_id: str
+1 -1
View File
@@ -17,4 +17,4 @@ class SpaceAccountBindingRequiredError(AccountEmailMismatchError):
code = 'space_account_binding_required'
def __str__(self) -> str:
return 'This local account must bind a LangBot Account from Account settings before LangBot Account login'
return 'This local Account must bind Space from Account settings before Space login'
@@ -33,28 +33,6 @@ class ModelProvider(Base):
)
class CodexCredential(Base):
"""Server-only OAuth state. Never joined into provider/model serialization."""
__tablename__ = 'codex_credentials'
provider_uuid = sqlalchemy.Column(sqlalchemy.String(255), primary_key=True)
workspace_uuid = sqlalchemy.Column(sqlalchemy.String(36), nullable=False)
payload = sqlalchemy.Column(sqlalchemy.JSON, nullable=False, default=dict)
version = sqlalchemy.Column(sqlalchemy.Integer, nullable=False, default=0)
lease_owner = sqlalchemy.Column(sqlalchemy.String(64), nullable=True)
lease_until = sqlalchemy.Column(sqlalchemy.Float, nullable=False, default=0)
__table_args__ = (
sqlalchemy.ForeignKeyConstraint(
['workspace_uuid', 'provider_uuid'],
['model_providers.workspace_uuid', 'model_providers.uuid'],
name='fk_codex_credentials_workspace_provider',
ondelete='CASCADE',
),
sqlalchemy.Index('ix_codex_credentials_workspace', 'workspace_uuid'),
)
class LLMModel(Base):
"""LLM model"""

Some files were not shown because too many files have changed in this diff Show More