Merge remote-tracking branch 'origin/master' into dev/4.11.x

# Conflicts:
#	src/langbot/pkg/api/http/controller/groups/pipelines/pipelines.py
#	src/langbot/pkg/api/http/service/bot.py
#	src/langbot/pkg/provider/runners/localagent.py
#	src/langbot/templates/metadata/pipeline/ai.yaml
#	tests/unit_tests/api/service/test_bot_service.py
#	tests/unit_tests/provider/runners/test_difysvapi_runner.py
#	tests/unit_tests/utils/test_safe_regex.py
#	web/src/app/infra/entities/adapter-categories.ts
#	web/src/app/wizard/page.tsx
#	web/src/i18n/locales/en-US.ts
#	web/src/i18n/locales/ja-JP.ts
#	web/src/i18n/locales/zh-Hans.ts
#	web/tests/e2e/plugin-page-auth.spec.ts
This commit is contained in:
Hyu
2026-08-31 17:17:47 +08:00
67 changed files with 3604 additions and 255 deletions
+38 -11
View File
@@ -1,13 +1,14 @@
from __future__ import annotations
import asyncio
import httpx
import typing
import json
import os
import typing
from pathlib import Path
import httpx
from .errors import DifyAPIError
from pathlib import Path
import os
_MAX_DIFY_RESPONSE_BYTES = 1024 * 1024
_MAX_DIFY_SSE_LINE_BYTES = 1024 * 1024
@@ -15,6 +16,32 @@ _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,
*,
@@ -56,16 +83,16 @@ async def _iter_sse_json(
line = raw_line.rstrip(b'\r').strip()
if not line or not line.startswith(b'data:'):
continue
payload = json.loads(line[5:].decode('utf-8', errors='replace'))
if isinstance(payload, dict):
payload = _decode_sse_data(line)
if payload is not None:
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 = json.loads(line[5:].decode('utf-8', errors='replace'))
if isinstance(payload, dict):
payload = _decode_sse_data(line)
if payload is not None:
yield payload
@@ -242,7 +269,7 @@ class AsyncDifyServiceClient:
file: httpx._types.FileTypes,
user: str,
timeout: float = 30.0,
) -> str:
) -> dict[str, typing.Any]:
# 处理 Path 对象
if isinstance(file, Path):
if not file.exists():
@@ -271,6 +298,6 @@ class AsyncDifyServiceClient:
timeout=timeout,
) as response:
body = await _read_limited_response(response)
if response.status_code != 201:
if response.status_code not in (200, 201):
raise DifyAPIError(f'{response.status_code} {body.decode(errors="replace")}')
return json.loads(body)
return _decode_upload_response(body)
+63
View File
@@ -422,6 +422,69 @@ 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,
@@ -940,6 +940,13 @@ 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
@@ -321,6 +321,34 @@ class WecomCSClient:
raise Exception('Failed to send image 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)
@@ -43,10 +43,13 @@ 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
try:
pipeline_uuid = await self.ap.pipeline_service.create_pipeline(
request_context,
await quart.request.json,
pipeline_data,
default=create_as_default,
)
except ValueError as exc:
return self.http_status(400, -1, str(exc))
@@ -131,9 +134,7 @@ class PipelinesRouterGroup(group.RouterGroup):
except Exception as exc:
self.ap.logger.warning('Unable to list skills for pipeline extensions: %s', exc)
available_skills = []
extensions_prefs = normalize_extension_preferences(
pipeline.get('extensions_preferences')
)
extensions_prefs = normalize_extension_preferences(pipeline.get('extensions_preferences'))
return self.success(
data={
'enable_all_plugins': extensions_prefs.get('enable_all_plugins', True),
@@ -192,9 +193,7 @@ class PipelinesRouterGroup(group.RouterGroup):
bound_skills=json_data.get('bound_skills', []),
enable_all_skills=json_data.get('enable_all_skills', True),
bound_mcp_resources=json_data.get('bound_mcp_resources'),
mcp_resource_agent_read_enabled=json_data.get(
'mcp_resource_agent_read_enabled'
),
mcp_resource_agent_read_enabled=json_data.get('mcp_resource_agent_read_enabled'),
)
except ValueError as exc:
return self.http_status(400, -1, str(exc))
@@ -152,6 +152,24 @@ 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'],
@@ -206,6 +206,20 @@ 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'],
+50
View File
@@ -2,6 +2,7 @@ from __future__ import annotations
import uuid
import typing
import json
import sqlalchemy
from ....core import app
@@ -11,6 +12,8 @@ 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:
@@ -797,6 +800,53 @@ class BotService:
'stale_routes': stale_routes,
}
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,
+76
View File
@@ -11,6 +11,9 @@ 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
@@ -238,3 +241,76 @@ 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}
+10 -1
View File
@@ -151,7 +151,16 @@ class LangBotMCPServer:
)
async def create_pipeline(pipeline_data: dict) -> str:
context = _authorized(Permission.RESOURCE_MANAGE)
return _dump({'uuid': await ap.pipeline_service.create_pipeline(context, pipeline_data)})
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,
)
}
)
@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:
@@ -47,3 +47,10 @@ 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
@@ -3,6 +3,7 @@
from __future__ import annotations
import asyncio
import contextlib
import dataclasses
import datetime
import json
@@ -82,7 +83,7 @@ def _verify_connection(connection: sqlite3.Connection, expected_revision: str) -
def _verify_file(path: pathlib.Path, expected_revision: str) -> None:
with _open_read_only(path) as connection:
with contextlib.closing(_open_read_only(path)) as connection:
_verify_connection(connection, expected_revision)
@@ -119,12 +120,16 @@ def _write_manifest(backup: SQLiteMigrationBackup, status: str, **extra: typing.
def _fsync_file(path: pathlib.Path, *, reopen_attempts: int = 20) -> None:
"""Sync a file, tolerating delayed visibility after replace on bind mounts."""
"""Sync a file, tolerating delayed visibility after replace on bind mounts.
Uses O_RDWR so os.fsync works on Windows (where _commit requires write
access to the file descriptor).
"""
descriptor: int | None = None
for attempt in range(reopen_attempts):
try:
descriptor = os.open(path, os.O_RDONLY)
descriptor = os.open(path, os.O_RDWR)
break
except FileNotFoundError:
if attempt + 1 >= reopen_attempts:
@@ -138,13 +143,37 @@ def _fsync_file(path: pathlib.Path, *, reopen_attempts: int = 20) -> None:
def _fsync_directory(path: pathlib.Path) -> None:
descriptor = os.open(path, os.O_RDONLY)
if os.name == 'nt':
# Windows cannot fsync directory handles opened through os.open.
return
descriptor = os.open(path, os.O_RDONLY | getattr(os, 'O_DIRECTORY', 0))
try:
os.fsync(descriptor)
finally:
os.close(descriptor)
def _remove_stale_temporary_files(
directory: pathlib.Path,
*,
prefix: str,
suffix: str,
) -> None:
"""Remove temporary files left by an interrupted backup or restore."""
for candidate in directory.iterdir():
if candidate.is_dir() or not candidate.name.startswith(prefix) or not candidate.name.endswith(suffix):
continue
try:
candidate.unlink()
except FileNotFoundError:
continue
except PermissionError:
# Another process may still own this file. Do not turn harmless
# cleanup into a migration failure; its unique name cannot collide.
continue
def _create_backup(
database_path: pathlib.Path,
source_revision: str,
@@ -153,6 +182,11 @@ def _create_backup(
backup_directory = database_path.parent / 'migration-backups'
backup_directory.mkdir(mode=0o700, parents=True, exist_ok=True)
os.chmod(backup_directory, 0o700)
_remove_stale_temporary_files(
backup_directory,
prefix=f'.{database_path.stem}-pre-',
suffix='.creating',
)
created_at = datetime.datetime.now(datetime.UTC).strftime('%Y-%m-%dT%H-%M-%S.%fZ')
stem = (
f'{database_path.stem}-pre-{_safe_label(target_revision)}-'
@@ -169,11 +203,8 @@ def _create_backup(
temporary_path = pathlib.Path(temporary_name)
try:
with (
_open_read_only(database_path) as source,
sqlite3.connect(
temporary_path,
timeout=30,
) as destination,
contextlib.closing(_open_read_only(database_path)) as source,
contextlib.closing(sqlite3.connect(temporary_path, timeout=30)) as destination,
):
source.execute('PRAGMA busy_timeout = 30000')
source.backup(destination)
@@ -221,6 +252,11 @@ async def create_verified_backup(
def _restore_backup(backup: SQLiteMigrationBackup) -> None:
_verify_file(backup.backup_path, backup.source_revision)
_remove_stale_temporary_files(
backup.database_path.parent,
prefix=f'.{backup.database_path.name}.',
suffix='.restoring',
)
descriptor, temporary_name = tempfile.mkstemp(
prefix=f'.{backup.database_path.name}.',
suffix='.restoring',
@@ -230,11 +266,8 @@ def _restore_backup(backup: SQLiteMigrationBackup) -> None:
temporary_path = pathlib.Path(temporary_name)
try:
with (
_open_read_only(backup.backup_path) as source,
sqlite3.connect(
temporary_path,
timeout=30,
) as destination,
contextlib.closing(_open_read_only(backup.backup_path)) as source,
contextlib.closing(sqlite3.connect(temporary_path, timeout=30)) as destination,
):
source.backup(destination)
destination.commit()
+1 -1
View File
@@ -210,7 +210,7 @@ _ALLOWED_SCOPED_BUILTIN_FUNCTION_TYPES = {
'now': sqlalchemy.sql.functions.now,
'sum': sqlalchemy.sql.functions.sum,
}
_ALLOWED_SCOPED_GENERIC_FUNCTIONS = frozenset({'date_trunc', 'length', 'nullif'})
_ALLOWED_SCOPED_GENERIC_FUNCTIONS = frozenset({'date_trunc', 'length', 'nullif', 'strftime'})
_ALLOWED_SCOPED_CUSTOM_OPERATORS = frozenset({'<=>'})
_ALLOWED_SCOPED_STATEMENT_TYPES = (
sqlalchemy.sql.dml.UpdateBase,
@@ -5,6 +5,11 @@ from .. import entities
import langbot_plugin.api.entities.builtin.pipeline.query as pipeline_query
from ....utils.safe_regex import SafeRegexError, mask_patterns
# Legacy sensitive-words.json files shipped ~70 rules, which exceeds the
# default safe_regex per-call cap of 64 and used to fail-close every message.
# Keep one 50ms CPU budget for the whole list; only raise the pattern cap.
_MAX_SENSITIVE_WORD_PATTERNS = 256
@filter_model.filter_class('ban-word-filter')
class BanWordFilter(filter_model.ContentFilter):
@@ -14,24 +19,29 @@ class BanWordFilter(filter_model.ContentFilter):
pass
async def process(self, query: pipeline_query.Query, message: str) -> entities.FilterResult:
words = self.ap.sensitive_meta.data.get('words') or []
mask = self.ap.sensitive_meta.data['mask']
mask_word = self.ap.sensitive_meta.data['mask_word']
try:
found, message = await mask_patterns(
self.ap.sensitive_meta.data['words'],
found, current = await mask_patterns(
words,
message,
mask=self.ap.sensitive_meta.data['mask'],
mask_word=self.ap.sensitive_meta.data['mask_word'],
mask=mask,
mask_word=mask_word,
max_pattern_count=_MAX_SENSITIVE_WORD_PATTERNS,
)
except SafeRegexError as exc:
return entities.FilterResult(
level=entities.ResultLevel.BLOCK,
replacement='',
user_notice='内容安全检查配置有误,请检查敏感词设置',
user_notice='内容检查规则执行失败,请联系管理员',
console_notice=f'Sensitive-word regex rejected: {exc}',
)
return entities.FilterResult(
level=entities.ResultLevel.MASKED if found else entities.ResultLevel.PASS,
replacement=message,
replacement=current,
user_notice='消息中存在不合适的内容, 请修改' if found else '',
console_notice='',
)
+61 -12
View File
@@ -25,6 +25,7 @@ from linebot.v3.webhooks import (
ImageMessageContent,
VideoMessageContent,
AudioMessageContent,
UserMentionee,
)
# from linebot import WebhookParser
@@ -58,15 +59,19 @@ class LINEMessageConverter(abstract_platform_adapter.AbstractMessageConverter):
return content_list
@staticmethod
async def target2yiri(message, bot_client) -> platform_message.MessageChain:
def __init__(self, bot_account_id: str = ''):
self.bot_account_id = bot_account_id
async def target2yiri(self, message, bot_client) -> platform_message.MessageChain:
lb_msg_list = []
msg_create_time = datetime.datetime.fromtimestamp(int(message.timestamp) / 1000)
lb_msg_list.append(platform_message.Source(id=message.webhook_event_id, time=msg_create_time))
if isinstance(message.message, TextMessageContent):
lb_msg_list.append(platform_message.Plain(text=message.message.text))
lb_msg_list.extend(
self._build_text_components(message.message.text, getattr(message.message, 'mention', None))
)
elif isinstance(message.message, AudioMessageContent):
pass
elif isinstance(message.message, VideoMessageContent):
@@ -86,22 +91,60 @@ class LINEMessageConverter(abstract_platform_adapter.AbstractMessageConverter):
lb_msg_list.append(platform_message.Image(base64=data_uri))
return platform_message.MessageChain(lb_msg_list)
def _build_text_components(self, text: str, mention) -> list:
"""Build message components from text, inserting At components for mentions.
LINE provides mention positions (index/length) and is_self per mentionee in the
webhook payload. Mapping the bot mention to At(target=bot_account_id) makes the
'at-bot' group respond rule work for LINE, consistent with other adapters.
"""
components: list = []
if not mention or not mention.mentionees:
if text:
components.append(platform_message.Plain(text=text))
return components
segments: list[tuple[int, int, object]] = sorted((m.index, m.index + m.length, m) for m in mention.mentionees)
cursor = 0
for start, end, mentionee in segments:
if start < cursor:
start, end = cursor, min(end, len(text))
if start < cursor or end <= start or end > len(text):
continue
if start > cursor:
components.append(platform_message.Plain(text=text[cursor:start]))
if isinstance(mentionee, UserMentionee):
target = self.bot_account_id if mentionee.is_self else mentionee.user_id
if not target:
target = text[start:end]
else:
target = text[start:end]
# At.__str__ already prepends '@', so strip one from the LINE text token.
display = text[start:end].lstrip('@')
components.append(platform_message.At(target=str(target), display=display))
cursor = end
if cursor < len(text):
components.append(platform_message.Plain(text=text[cursor:]))
return components
class LINEEventConverter(abstract_platform_adapter.AbstractEventConverter):
def __init__(self, bot_account_id: str = ''):
self.bot_account_id = bot_account_id
self.message_converter = LINEMessageConverter(bot_account_id)
@staticmethod
async def yiri2target(
event: platform_events.MessageEvent,
) -> MessageEvent:
pass
@staticmethod
async def target2yiri(event, bot_client) -> platform_events.Event:
message_chain = await LINEMessageConverter.target2yiri(event, bot_client)
async def target2yiri(self, event, bot_client) -> platform_events.Event:
message_chain = await self.message_converter.target2yiri(event, bot_client)
if event.source.type == 'user':
return platform_events.FriendMessage(
sender=platform_entities.Friend(
id=event.message.id,
id=event.source.user_id,
nickname=event.source.user_id,
remark='',
),
@@ -110,13 +153,19 @@ class LINEEventConverter(abstract_platform_adapter.AbstractEventConverter):
source_platform_object=event,
)
else:
# 'group' and 'room' sources carry the stable chat id under different
# field names; user_id may be absent for some members, so fall back
# to the group/room id rather than the per-message id.
group_id = event.source.group_id if event.source.type == 'group' else event.source.room_id
member_id = event.source.user_id or group_id
return platform_events.GroupMessage(
sender=platform_entities.GroupMember(
id=event.event.sender.sender_id.open_id,
member_name=event.event.sender.sender_id.union_id,
id=member_id,
member_name=member_id,
permission=platform_entities.Permission.Member,
group=platform_entities.Group(
id=event.message.id,
id=group_id,
name='',
permission=platform_entities.Permission.Member,
),
@@ -163,8 +212,8 @@ class LINEAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
listeners={},
card_id_dict={},
seq=1,
event_converter=LINEEventConverter(),
message_converter=LINEMessageConverter(),
event_converter=LINEEventConverter(bot_account_id),
message_converter=LINEMessageConverter(bot_account_id),
line_webhook=line_webhook,
parser=parser,
configuration=configuration,
+50 -29
View File
@@ -329,17 +329,12 @@ class QQOfficialAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter
content_type = content.get('type', 'text')
if content_type == 'text':
if target_type == 'c2c':
await self.bot.send_private_text_msg(
if target_type in {'c2c', 'group'}:
await self._send_c2c_or_group_text_reply(
target_type,
target_id,
content['content'],
qq_official_event.d_id,
)
elif target_type == 'group':
await self.bot.send_group_text_msg(
target_id,
content['content'],
qq_official_event.d_id,
msg_id=qq_official_event.d_id,
)
elif content_type == 'image':
@@ -383,6 +378,39 @@ class QQOfficialAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter
async def send_message(self, target_type: str, target_id: str, message: platform_message.MessageChain):
pass
async def _send_c2c_or_group_text_reply(
self,
target_type: str,
target_id: str,
content: str,
*,
msg_id: typing.Optional[str] = None,
event_id: typing.Optional[str] = None,
msg_seq: int = 1,
) -> None:
"""Send a text reply using the configured C2C/group render mode."""
use_markdown = self.config.get('enable-markdown-rendering', False)
if target_type == 'c2c':
send = self.bot.send_private_markdown_msg if use_markdown else self.bot.send_private_text_msg
await send(
user_openid=target_id,
content=content,
msg_id=msg_id,
event_id=event_id,
msg_seq=msg_seq,
)
elif target_type == 'group':
send = self.bot.send_group_markdown_msg if use_markdown else self.bot.send_group_text_msg
await send(
group_openid=target_id,
content=content,
msg_id=msg_id,
event_id=event_id,
msg_seq=msg_seq,
)
else:
raise ValueError(f'Unsupported QQ Official text reply target: {target_type}')
def register_listener(
self,
event_type: typing.Type[platform_events.Event],
@@ -650,13 +678,13 @@ class QQOfficialAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter
# 用第一个 chunk 的文本建立会话(不发 "..." 避免污染前缀)
ctx['session_started'] = True
# 发送内容 = 全量累积文本
# QQ API 的 replace 模式不允许修改已下发前缀,所以:
# - 首次:发送全部文本,建立会话
# - 后续:只能发送新增部分(append 行为)
content_to_send = ctx['accumulated_text'][ctx['sent_length'] :]
if not content_to_send and not is_final:
# `replace` mode requires every update to contain the previously
# delivered content as its prefix. `sent_length` only tells us whether
# a non-final snapshot has new content; it must not truncate the
# content sent to QQ.
if len(ctx['accumulated_text']) <= ctx['sent_length'] and not is_final:
return
content_to_send = ctx['accumulated_text']
input_state = 10 if is_final else 1
@@ -778,20 +806,13 @@ class QQOfficialAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter
return
try:
if target_type == 'c2c':
await self.bot.send_private_text_msg(
user_openid=target_id,
content=text,
event_id=event_id,
msg_seq=msg_seq,
)
elif target_type == 'group':
await self.bot.send_group_text_msg(
group_openid=target_id,
content=text,
event_id=event_id,
msg_seq=msg_seq,
)
await self._send_c2c_or_group_text_reply(
target_type,
target_id,
text,
event_id=event_id,
msg_seq=msg_seq,
)
except Exception:
await self.logger.error(f'QQ Official: synthetic reply delivery failed: {traceback.format_exc()}')
@@ -96,6 +96,18 @@ spec:
type: boolean
required: true
default: false
- name: enable-markdown-rendering
label:
en_US: Enable Markdown Rendering
zh_Hans: 启用 Markdown 渲染
zh_Hant: 啟用 Markdown 渲染
description:
en_US: Render non-stream C2C and QQ group text replies as Markdown. Channel messages always use plain text and are not affected by this setting.
zh_Hans: 将非流式 C2C 私聊和 QQ 群聊文本回复渲染为 Markdown。频道消息始终以纯文本发送,不受此设置影响。
zh_Hant: 將非串流 C2C 私聊與 QQ 群聊文字回覆渲染為 Markdown。頻道訊息一律以純文字傳送,不受此設定影響。
type: boolean
required: true
default: false
- name: webhook_url
label:
en_US: Webhook Callback URL
+10 -3
View File
@@ -107,7 +107,7 @@ class WecomEventConverter(abstract_platform_adapter.AbstractEventConverter):
if event.type == 'text':
yiri_chain = await WecomMessageConverter.target2yiri(event.message, event.message_id)
friend = platform_entities.Friend(
id=f'u{event.user_id}',
id=f'{event.receiver_id}|u{event.user_id}',
nickname=nickname,
remark='',
)
@@ -117,7 +117,7 @@ class WecomEventConverter(abstract_platform_adapter.AbstractEventConverter):
)
elif event.type == 'image':
friend = platform_entities.Friend(
id=f'u{event.user_id}',
id=f'{event.receiver_id}|u{event.user_id}',
nickname=nickname,
remark='',
)
@@ -197,7 +197,7 @@ class WecomCSAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
content_list = await WecomMessageConverter.yiri2target(message, self.bot)
for content in content_list:
msgid = f'langbot_{uuid.uuid4().hex}'
msgid = f'{uuid.uuid4().hex}'
if content['type'] == 'text':
await self.bot.send_text_msg(
open_kfid=open_kfid,
@@ -205,6 +205,13 @@ class WecomCSAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
msgid=msgid,
content=content['content'],
)
elif content['type'] == 'image':
await self.bot.send_image_msg(
open_kfid=open_kfid,
external_userid=external_userid,
msgid=msgid,
media_id=content['media_id'],
)
def set_bot_uuid(self, bot_uuid: str):
"""设置 bot UUID(用于生成 webhook URL"""
+7 -2
View File
@@ -1934,9 +1934,14 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
return plugins
async def get_plugin_info(self, author: str, plugin_name: str) -> dict[str, Any]:
async def get_plugin_info(self, author: str, plugin_name: str) -> dict[str, Any] | None:
runtime_handler = self._runtime_handler()
binding = await self._target_binding(author, plugin_name)
try:
binding = await self._target_binding(author, plugin_name)
except ValueError as exc:
if str(exc) == f'Plugin {author}/{plugin_name} is not installed in this Workspace':
return None
raise
with runtime_handler.installation_scope(binding):
return await runtime_handler.get_plugin_info(author, plugin_name)
@@ -747,9 +747,24 @@ class LiteLLMRequester(requester.ProviderAPIRequester):
converted_parts = []
for part in content:
if isinstance(part, dict) and part.get('type') == 'image_base64':
part['image_url'] = {'url': part['image_base64']}
part['type'] = 'image_url'
del part['image_base64']
# History trimming (SessionManager) clears image_base64
# on past turns and exclude_none serialization drops
# the key entirely, so the replayed part may carry no
# payload. Prefer the base64 payload; fall back to an
# image_url that survived on the same element; drop
# hollow parts instead of raising KeyError (#2469).
image_b64 = part.get('image_base64')
fallback_url = None
if not image_b64:
raw_image_url = part.get('image_url')
if isinstance(raw_image_url, dict):
fallback_url = raw_image_url.get('url')
if image_b64 or fallback_url:
part['image_url'] = {'url': image_b64 or fallback_url}
part['type'] = 'image_url'
part.pop('image_base64', None)
else:
continue
# OpenAI-compatible chat models reject non-image file parts
# (audio/document base64 or url). These originate from Voice /
# File attachments — including ones replayed from conversation
+13 -4
View File
@@ -30,10 +30,16 @@ class SafeRegexTimeoutError(SafeRegexError):
"""Raised when the regex engine exhausts the operation CPU budget."""
def _validate_patterns(patterns: Sequence[str]) -> tuple[str, ...]:
def _validate_patterns(
patterns: Sequence[str],
*,
max_pattern_count: int = MAX_PATTERN_COUNT,
) -> tuple[str, ...]:
if max_pattern_count < 1:
raise ValueError('max_pattern_count must be positive')
if len(patterns) > max_pattern_count:
raise SafeRegexLimitError(f'At most {max_pattern_count} regex patterns are allowed')
normalized = tuple(patterns)
if len(normalized) > MAX_PATTERN_COUNT:
raise SafeRegexLimitError(f'At most {MAX_PATTERN_COUNT} regex patterns are allowed')
for pattern in normalized:
if not isinstance(pattern, str):
raise SafeRegexError('Regex patterns must be strings')
@@ -118,8 +124,9 @@ def _mask_patterns_sync(
mask: str,
mask_word: str,
timeout_seconds: float,
max_pattern_count: int,
) -> tuple[bool, str]:
normalized_patterns = _validate_patterns(patterns)
normalized_patterns = _validate_patterns(patterns, max_pattern_count=max_pattern_count)
_validate_input(value)
if len(mask) > MAX_REPLACEMENT_CHARS or len(mask_word) > MAX_REPLACEMENT_CHARS:
raise SafeRegexLimitError(f'Regex replacements may contain at most {MAX_REPLACEMENT_CHARS} characters')
@@ -165,6 +172,7 @@ async def mask_patterns(
mask: str,
mask_word: str,
timeout_seconds: float = DEFAULT_OPERATION_TIMEOUT_SECONDS,
max_pattern_count: int = MAX_PATTERN_COUNT,
) -> tuple[bool, str]:
"""Apply untrusted masking patterns with bounded CPU and output growth."""
@@ -177,4 +185,5 @@ async def mask_patterns(
mask=mask,
mask_word=mask_word,
timeout_seconds=timeout_seconds,
max_pattern_count=max_pattern_count,
)
@@ -240,7 +240,7 @@ class InvitationDeliveryService:
@staticmethod
def _plain_text(workspace_name: str, invitation_link: str) -> str:
return (
'You have been invited to LangBot Cloud\n\n'
'You have been invited to join a Workspace in LangBot\n\n'
f'Join the Workspace “{workspace_name}” to collaborate with your team.\n\n'
f'Accept invitation: {invitation_link}\n\n'
'This secure invitation expires in 7 days and can only be accepted by the email address '
@@ -258,30 +258,77 @@ class InvitationDeliveryService:
<head>
<meta charset="utf-8">
<meta name="viewport" content="width=device-width,initial-scale=1">
<title>Join {escaped_workspace} on LangBot Cloud</title>
<meta http-equiv="X-UA-Compatible" content="IE=edge">
<title>Join {escaped_workspace} in LangBot</title>
</head>
<body style="margin:0;background:#f4f7fb;color:#152033;font-family:Inter,-apple-system,BlinkMacSystemFont,'Segoe UI',sans-serif;">
<div style="display:none;max-height:0;overflow:hidden;opacity:0;">You have been invited to join {escaped_workspace} on LangBot Cloud.</div>
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" style="background:#f4f7fb;padding:40px 16px;">
<tr><td align="center">
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" style="max-width:600px;background:#ffffff;border:1px solid #e5eaf2;border-radius:16px;overflow:hidden;box-shadow:0 12px 32px rgba(20,49,93,.08);">
<tr><td style="padding:28px 36px;background:linear-gradient(135deg,#0f172a,#1d4ed8);color:#ffffff;">
<div style="font-size:14px;font-weight:700;letter-spacing:.08em;text-transform:uppercase;opacity:.78;">LangBot Cloud</div>
<div style="font-size:26px;font-weight:700;margin-top:8px;line-height:1.25;">Youre invited</div>
</td></tr>
<tr><td style="padding:36px;">
<p style="margin:0 0 18px;font-size:16px;line-height:1.65;color:#475569;">You have been invited to collaborate in this Workspace:</p>
<div style="margin:0 0 26px;padding:18px 20px;background:#f8fafc;border:1px solid #e2e8f0;border-radius:12px;font-size:18px;font-weight:700;color:#0f172a;">{escaped_workspace}</div>
<table role="presentation" cellspacing="0" cellpadding="0"><tr><td style="border-radius:9px;background:#2563eb;">
<a href="{escaped_link}" style="display:inline-block;padding:13px 22px;color:#ffffff;text-decoration:none;font-size:15px;font-weight:700;">Accept invitation</a>
</td></tr></table>
<p style="margin:26px 0 8px;font-size:14px;line-height:1.6;color:#64748b;">This invitation expires in 7 days and is bound to the email address that received it.</p>
<p style="margin:0 0 8px;font-size:13px;line-height:1.6;color:#94a3b8;">If the button does not work, copy and paste this URL into your browser:</p>
<p style="margin:0;padding:12px;background:#f8fafc;border-radius:8px;word-break:break-all;font-size:12px;line-height:1.55;color:#475569;">{escaped_link}</p>
</td></tr>
<tr><td style="padding:20px 36px;border-top:1px solid #eef2f7;font-size:12px;line-height:1.6;color:#94a3b8;">If you were not expecting this invitation, you can safely ignore this email.</td></tr>
</table>
</td></tr>
<body style="margin:0;padding:0;background:#f4f7fb;color:#111827;font-family:Arial,'Helvetica Neue',sans-serif;">
<div style="display:none;max-height:0;overflow:hidden;opacity:0;">You have been invited to join {escaped_workspace} in LangBot.</div>
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" border="0" style="width:100%;background:#f4f7fb;">
<tr>
<td align="center" style="padding:48px 16px;">
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" border="0" style="width:100%;max-width:600px;">
<tr>
<td style="padding:0 4px 20px;">
<img src="https://docs.langbot.app/langbot-logo.png" alt="LangBot" width="34" height="34" style="display:inline-block;width:34px;height:34px;border:0;vertical-align:middle;">
<span style="display:inline-block;margin-left:10px;vertical-align:middle;font-size:18px;font-weight:700;letter-spacing:-.01em;">LangBot</span>
</td>
</tr>
<tr>
<td style="background:#ffffff;border-radius:10px;overflow:hidden;">
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" border="0">
<tr>
<td style="padding:42px 42px 38px;">
<div style="margin:0 0 12px;font-size:13px;line-height:1.4;font-weight:600;color:#5f6f84;">Workspace invitation</div>
<h1 style="margin:0 0 16px;font-size:28px;line-height:1.25;font-weight:700;letter-spacing:-.025em;color:#111827;">Youre invited to collaborate</h1>
<p style="margin:0 0 28px;font-size:15px;line-height:1.7;color:#526173;">Join your team in LangBot and start building together in this Workspace.</p>
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" border="0" style="background:#f6f8fb;border-radius:8px;">
<tr>
<td style="padding:16px 18px;">
<div style="margin:0 0 4px;font-size:11px;line-height:1.4;font-weight:700;letter-spacing:.08em;text-transform:uppercase;color:#5f6f84;">Workspace</div>
<div style="font-size:18px;line-height:1.4;font-weight:700;color:#111827;">{escaped_workspace}</div>
</td>
</tr>
</table>
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" border="0">
<tr><td height="28" style="height:28px;font-size:0;line-height:0;">&nbsp;</td></tr>
</table>
<table role="presentation" cellspacing="0" cellpadding="0" border="0">
<tr>
<td style="background:#2563eb;border-radius:8px;">
<a href="{escaped_link}" target="_blank" style="display:inline-block;padding:13px 22px;font-size:15px;line-height:1.2;font-weight:700;color:#ffffff;text-decoration:none;border-radius:8px;">Accept invitation</a>
</td>
</tr>
</table>
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" border="0">
<tr><td height="32" style="height:32px;font-size:0;line-height:0;">&nbsp;</td></tr>
</table>
<table role="presentation" width="100%" cellspacing="0" cellpadding="0" border="0" style="border-top:1px solid #e8edf4;">
<tr>
<td style="padding-top:22px;">
<p style="margin:0 0 10px;font-size:13px;line-height:1.6;color:#5f6f84;">For your security, this invitation expires in 7 days and only works for the email address that received it.</p>
<a href="{escaped_link}" target="_blank" style="font-size:13px;line-height:1.6;font-weight:600;color:#2563eb;text-decoration:none;">Open invitation link&nbsp;&rarr;</a>
</td>
</tr>
</table>
</td>
</tr>
</table>
</td>
</tr>
<tr>
<td align="center" style="padding:20px 24px 0;font-size:12px;line-height:1.6;color:#5f6f84;">
Sent by LangBot<br>
If you were not expecting this invitation, you can safely ignore this email.
</td>
</tr>
</table>
</td>
</tr>
</table>
</body>
</html>'''
+20
View File
@@ -7,6 +7,9 @@
// Read config from script tag data attributes
var scriptEl = document.currentScript;
var scriptTitle = scriptEl ? scriptEl.getAttribute("data-title") : null;
var scriptTestNotice = scriptEl
? scriptEl.getAttribute("data-test-notice")
: null;
// ========== i18n ==========
var I18N = {
@@ -192,6 +195,7 @@
.lb-header-btn { background: none; border: none; color: #fff; cursor: pointer; padding: 4px; border-radius: 6px; display: flex; align-items: center; justify-content: center; opacity: 0.8; transition: opacity 0.15s; }\
.lb-header-btn:hover { opacity: 1; }\
.lb-header-btn svg { width: 18px; height: 18px; fill: currentColor; }\
.lb-test-notice { padding: 8px 16px; border-bottom: 1px solid #fde68a; background: #fffbeb; color: #92400e; font-size: 12px; line-height: 1.5; text-align: center; flex-shrink: 0; }\
.lb-messages { flex: 1; overflow-y: auto; padding: 16px; display: flex; flex-direction: column; gap: 16px; scroll-behavior: smooth; }\
.lb-messages::-webkit-scrollbar { width: 6px; }\
.lb-messages::-webkit-scrollbar-track { background: transparent; }\
@@ -1240,6 +1244,14 @@
// Root container
var root = document.createElement("div");
root.id = "langbot-widget-root";
root.langbotDestroy = function () {
wsDisconnect();
if (state.historyReloadTimer) {
clearTimeout(state.historyReloadTimer);
state.historyReloadTimer = null;
}
root.remove();
};
document.body.appendChild(root);
var shadow = root.attachShadow({ mode: "open" });
@@ -1328,6 +1340,14 @@
header.appendChild(headerActions);
panel.appendChild(header);
if (scriptTestNotice) {
var testNotice = document.createElement("div");
testNotice.className = "lb-test-notice";
testNotice.setAttribute("role", "note");
testNotice.textContent = scriptTestNotice;
panel.appendChild(testNotice);
}
// Messages area
var messages = document.createElement("div");
messages.className = "lb-messages";