Compare commits

..

27 Commits

Author SHA1 Message Date
leonoxo 855ae2bdba fix(line): map LINE mentions to At elements so the at-bot rule works (#2478)
The LINE adapter passed text through as a single Plain component,
ignoring the mention payload (mentions[].index/length/isSelf) that the
Line Messaging API includes in the webhook. As a result:

- At(target=bot_account_id) never appeared in the message chain, so the
  'at-bot' group respond rule silently dropped every @bot mention.
- The bot only replied when the message happened to match the prefix
  rule (e.g. starting with 'ai').

Now LINEMessageConverter reads message.message.mention and builds the
chain per mention position:

- Bot mention (isSelf) -> At(target=bot_account_id) so AtBotRule matches
  the same way as other adapters (dingtalk/lark etc.).
- Other mentions -> At(target=<line user id>, display=<mention text>).
  At.__str__ already prepends '@', so the display text carries no
  double '@' and the rendered text (prefix/regexp rules, quotes,
  session context) is byte-identical to before.
- Missing/out-of-bounds mentions are skipped defensively.

target2yiri becomes an instance method (like wechatpad/aiocqhttp) so
the converters can hold bot_account_id; LINEAdapter passes it in from
its own config.
2026-08-27 18:38:27 +08:00
fishzjp b66db86bff fix(provider): stringify MCP tool results for OpenAI-compatible APIs (#2476)
execute_func_call returns list[ContentElement] for MCP tools, but the
runner assigned that list directly to the tool-message content. The
OpenAI chat-completions spec requires tool-message content to be a
string, so OpenAI-compatible endpoints return HTTP 500 when the raw
list is sent.

Serialize the list to a string before building the tool message, using
ContentElement.__str__ which returns the text payload for text elements
and a human-readable placeholder for images and files. Fixes #2457.
2026-08-27 18:19:04 +08:00
fishzjp 95b8736e93 fix(provider): tolerate trimmed image parts in litellm message conversion (#2475)
SessionManager clears image_base64 on past turns to save memory, and
exclude_none serialization drops the hollowed field entirely, so a
replayed history part can arrive as {'type': 'image_base64'} with no
payload. The converter accessed the missing key unconditionally and
raised KeyError on every turn after an image was sent.

Prefer the base64 payload when present, fall back to an image_url that
survived on the same element, and drop hollow parts otherwise (same
strategy as the existing file-part handling). Fixes #2469.
2026-08-27 18:07:43 +08:00
Yang cabde423a1 Fix wecomcs open_kfid msgid (#2449)
* Update wecomcs.py

fix bug wecomcs send_message open_kfid
event.receiver_id  is open_kfid

* Update wecomcs.py

fix bug msgid exceeds the 32-byte limit

* test(wecomcs): cover bounded message IDs and images

---------

Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-27 17:30:49 +08:00
ciri667 08307790e5 fix(cntfilter): allow legacy sensitive-word lists over 64 patterns (#2467)
* fix(cntfilter): allow legacy sensitive-word lists over 64 patterns

Legacy sensitive-words.json files shipped ~70 rules. After v4.10.7,
BanWordFilter treated the 64-pattern safe_regex cap as a hard failure
and blocked every message. Raise the cap only on the sensitive-word
path, keep the 50ms CPU budget, and truncate oversized lists with a
one-time warning.

Fixes #2443

* fix(cntfilter): reject oversized sensitive-word lists

---------

Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-25 23:53:41 +08:00
leonoxo 777fe1f20b fix(line): use stable source id for session identity (#2398)
LINEEventConverter.target2yiri() built Friend.id/Group.id from
event.message.id, which is unique per message. Every incoming message
therefore mapped to a new session key, so LINE users and groups lost
conversation context on every turn.

Use event.source.user_id/group_id/room_id instead, matching the stable
identifiers other adapters (e.g. Telegram) use for session identity.
Falls back to the group/room id when user_id is absent, per LINE's
documented behavior for some group/room members.
2026-08-25 12:30:29 +08:00
Hyu f0ee57c1e0 style(email): use generic LangBot invitation branding (#2466)
Co-authored-by: Junyan Qin <rockchinq@gmail.com>
2026-08-24 14:53:59 +08:00
Hyu a45e27e76e style(cloud): redesign workspace invitation email (#2464)
* style(cloud): redesign workspace invitation email

* fix(email): harden Outlook spacing and text contrast

---------

Co-authored-by: Junyan Qin <rockchinq@gmail.com>
2026-08-24 14:14:56 +08:00
Hyu 536fcdf29f fix(web): restore plugin page SDK loading (#2435)
Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-24 12:40:50 +08:00
QuasarRyan c87548c0b9 feat(qqofficial): add markdown reply rendering (#2459) 2026-08-24 11:59:23 +08:00
QuasarRyan 79634772da fix(qqofficial): send complete stream snapshots (#2458) 2026-08-24 11:49:22 +08:00
Hyu bb366779af fix(auth): align invitation password fields (#2461)
Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-24 10:53:05 +08:00
Hyu 1336f47cb4 fix(auth): enable local registration for OSS invitations (#2460)
Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-24 10:28:45 +08:00
RockChinQ 962366c507 fix(deps): make SeekDB optional for native installs 2026-08-21 12:07:36 +08:00
Hyu 23875b240f Merge pull request #2452 from langbot-app/ci/immutable-core-images
ci: publish immutable Core image tags
2026-08-21 10:19:58 +08:00
dadachann e699358a5a ci: publish immutable Core image tags 2026-08-21 02:16:34 +00:00
Hyu 14277d129c Merge pull request #2451 from langbot-app/fix/release-4.10.8-cloud
fix(cloud): carry runtime readiness fixes into 4.10.8
2026-08-21 03:12:12 +08:00
dadachann 6bf1546df2 style: format Cloud runtime readiness fixes 2026-08-20 19:08:41 +00:00
dadachann 0bec72a3f9 fix(cloud): run runtime initialization outside transaction gate 2026-08-20 18:55:24 +00:00
dadachann f36542135a fix(cloud): start API before runtime reconcile 2026-08-20 18:55:24 +00:00
dadachann 693c59b726 fix(cloud): import lifecycle task scope 2026-08-20 18:55:24 +00:00
dadachann c3fe312a43 fix(cloud): refresh directory during plugin startup 2026-08-20 18:55:24 +00:00
dadachann c4bad508d2 fix(runtime): honor configured cold reconcile timeout 2026-08-20 18:55:24 +00:00
RockChinQ 3d4a726cd8 chore: release v4.10.8 (#2450) 2026-08-21 01:47:39 +08:00
DongXiaoming e934f08adf fix(wecombot): deliver sandbox outbox media through full pipeline chain (#2328)
* fix(wecombot): align media upload protocol

* fix(wecombot): deliver outbox media in reply and fix tool call recording

- Integrate _send_media into reply_message and reply_message_chunk so
  sandbox outbox images/voices/files are uploaded and sent instead of
  being silently dropped.
- Add missing import base64 that caused _send_media to fail with a
  NameError swallowed by its except clause.
- Change yiri2target to return component dicts (text/image/voice/file)
  so callers can distinguish text from media.
- Fix _get_message_for_tool_context using result.first()/row[0] which
  returned a raw string instead of the ORM object, causing
  "'str' object has no attribute 'pipeline_id'" in tool call recording.
  Use result.scalars().first() per SQLAlchemy 2.0 convention.

* fix(pipeline): collect outbox attachments on final chunk with empty content

When the last streaming chunk has is_final=True but empty content
(e.g. the LLM sends all text in earlier chunks), the 'if result.content'
branch is skipped entirely, so _append_outbound_attachments never runs
and sandbox outbox images are silently dropped.

Add an elif branch for _is_final_assistant_message that creates an
empty MessageChain and still collects outbox attachments, so images
are delivered even when the final chunk carries no text.

* fix(box): bypass stdout truncation when reading outbox via exec

_read_outbox_via_exec used execute_tool which returns _serialize_result
where stdout is truncated to output_limit_chars (4000). A 7KB JPEG
encodes to ~9400 base64 chars, so the JSON payload was truncated and
json.loads failed silently, returning an empty list.

Call client.execute directly to get the raw BoxExecutionResult with
untruncated stdout, so base64 file data is preserved.

* fix(tests): adapt box and wrapper tests for client.execute and strict is_final check

- wrapper.py: restrict outbox collection on empty-content chunks to
  actual MessageChunk instances with is_final=True, not generic Mock
  objects that happen to have role='assistant'
- test_box_service.py: update _read_outbox_via_exec tests to mock
  client.execute (returning BoxExecutionResult) instead of
  execute_tool, matching the implementation change

* chore(wecombot): remove temporary upload log

* test(box): preserve direct outbox read and cleanup coverage

---------

Co-authored-by: fdc310 <2213070223@qq.com>
Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-20 16:56:20 +08:00
Hyu 7803d56254 fix(plugin): keep pre-tenancy plugin storage readable (#2446)
* fix(plugin): adopt legacy scoped storage rows

* fix(plugin): delete adopted legacy storage rows

* test(plugin): cover legacy storage deletion

* fix(plugin): avoid mutating legacy storage on reads

* fix(plugin): make legacy storage adoption atomic

* fix(plugin): resolve concurrent legacy adoption

* fix(plugin): upsert concurrent storage writes

* fix(plugin): retry reads after legacy adoption

* fix(plugin): deduplicate migrated storage keys

---------

Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-17 11:49:26 +08:00
Hyu 54c96a18e1 test(migration): preserve legacy plugin storage payloads (#2444)
Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-08-17 10:31:03 +08:00
56 changed files with 2595 additions and 265 deletions
+32 -13
View File
@@ -7,23 +7,42 @@ 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@v2
uses: actions/checkout@v4
with:
persist-credentials: false
- name: Generate Tag
id: generate_tag
- name: Set up Docker Buildx
uses: docker/setup-buildx-action@v3
- name: Generate image metadata
id: image
shell: bash
run: |
# 获取分支名称,把/替换为-
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
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 }}
+1 -1
View File
@@ -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 \
&& uv sync --extra seekdb \
&& apt-get purge -y --auto-remove curl git gnupg \
&& rm -rf /var/lib/apt/lists/* \
&& touch /.dockerenv
+18 -1
View File
@@ -10,6 +10,19 @@ 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:
@@ -20,6 +33,10 @@ 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:
@@ -101,7 +118,7 @@ uvx langbot
## System Requirements
- Python 3.10.1 or higher
- Python 3.11 or higher (lower than Python 4)
- Operating System: Linux, macOS, or Windows
## Differences from Source Installation
+34 -43
View File
@@ -16,12 +16,20 @@ This document describes how to use OceanBase SeekDB as the vector database backe
## Installation
SeekDB support is automatically included when you install LangBot. The required dependency `pyseekdb` is listed in `pyproject.toml`.
SeekDB is an optional LangBot feature. A normal LangBot installation uses
Chroma by default and does not install `pyseekdb` or its native bindings.
If you need to install it manually:
Choose the command that matches how you run LangBot:
```bash
pip install pyseekdb
# PyPI / uvx
uvx --from 'langbot[seekdb]@latest' langbot
# Installed package
pip install 'langbot[seekdb]'
# Source checkout
uv sync --extra seekdb
```
## ⚠️ Platform Compatibility
@@ -30,31 +38,36 @@ pip install pyseekdb
| Platform | Status | Notes |
|----------|--------|-------|
| 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 |
| 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 |
**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.
**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.
### Server Mode (Docker)
| Platform | Status | Notes |
|----------|--------|-------|
| Linux | ✅ Supported | Full Docker support |
| 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
| 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 |
### Server Mode (Remote Connection)
| Platform | Status | Notes |
|----------|--------|-------|
| All Platforms | ✅ Supported | Connect to SeekDB running on a remote Linux server |
| 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 |
**Recommendation for macOS/Windows users**: Deploy SeekDB on a Linux server and connect via server mode configuration.
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.
## Configuration
@@ -170,22 +183,23 @@ Key methods:
### Import Error
If you see: `ImportError: pyseekdb is not installed`
If you see: `SeekDB support is not installed`
Solution:
```bash
pip install pyseekdb
uv sync --extra seekdb
# or: uvx --from 'langbot[seekdb]@latest' langbot
```
### Embedded Mode Error on macOS/Windows
### Embedded Mode Is Unavailable on the Current Platform
**Error**:
```
RuntimeError: Embedded Client is not available because pylibseekdb is not available.
Please install pylibseekdb (Linux only) or use RemoteServerClient (host/port) instead.
```
**Cause**: `pylibseekdb` is only available on Linux platforms.
**Cause**: No compatible `pylibseekdb` wheel is installed for the current OS,
CPU architecture, Python version, and macOS deployment target.
**Solution**: Use server mode instead:
1. Deploy SeekDB on a Linux server or VM
@@ -208,29 +222,6 @@ 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:
+6 -2
View File
@@ -1,6 +1,6 @@
[project]
name = "langbot"
version = "4.10.7"
version = "4.10.8"
description = "Production-grade platform for building agentic IM bots"
readme = "README.md"
license-files = ["LICENSE"]
@@ -70,7 +70,6 @@ dependencies = [
"langchain-text-splitters>=1.1.2",
"chromadb>=1.0.0,<2.0.0",
"qdrant-client (>=1.15.1,<2.0.0)",
"pyseekdb==1.1.0.post3",
"langbot-plugin==0.5.5",
"asyncpg>=0.30.0",
"line-bot-sdk>=3.19.0",
@@ -108,6 +107,11 @@ classifiers = [
"Topic :: Communications :: Chat",
]
[project.optional-dependencies]
seekdb = [
"pyseekdb==1.1.0.post3",
]
[project.urls]
Homepage = "https://langbot.app"
Documentation = "https://docs.langbot.app"
+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,
@@ -46,6 +46,14 @@ 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
@@ -495,6 +503,145 @@ 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.
@@ -295,6 +295,34 @@ 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)
@@ -322,6 +322,7 @@ 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
return self.success(data={'initialized': True, **capabilities})
@self.route('/set-password', methods=['POST'], auth_type=group.AuthType.USER_TOKEN)
+17 -8
View File
@@ -1210,8 +1210,9 @@ 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).
Note: exec stdout is truncated by ``output_limit_chars``, so this path
only reliably transfers small files. The host path is preferred.
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.
"""
import json as _json
@@ -1265,14 +1266,22 @@ class BoxService:
' break\n'
'print(json.dumps(out))\n'
)
result = await self.execute_tool(
{'command': f"python3 - <<'LBPY'\n{script}\nLBPY", 'timeout_sec': 120},
query,
)
if not result.get('ok'):
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:
return []
try:
return _json.loads(str(result.get('stdout') or '').strip().splitlines()[-1])
return _json.loads(str(result.stdout or '').strip().splitlines()[-1])
except Exception:
return []
+33 -3
View File
@@ -301,11 +301,36 @@ 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:
self.task_mgr.create_task(
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(
self.directory_projection_service.run(),
name='cloud-directory-projection',
scopes=[core_entities.LifecycleControlScope.APPLICATION],
@@ -322,7 +347,6 @@ class Application:
name='cloud-manifest-refresh',
scopes=[core_entities.LifecycleControlScope.APPLICATION],
)
await self.plugin_connector.initialize_plugins()
# 后续可能会允许动态重启其他任务
# 故为了防止程序在非 Ctrl-C 情况下退出,这里创建一个不会结束的协程
@@ -348,6 +372,7 @@ 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().
@@ -529,6 +554,11 @@ 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)
+11 -8
View File
@@ -1,6 +1,6 @@
from __future__ import annotations
from .. import stage, app
from .. import stage, app, entities as core_entities
from ...utils import version, proxy, constants
from ...pipeline import pool, controller, pipelinemgr
from ...pipeline import aggregator as message_aggregator
@@ -292,14 +292,17 @@ 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()
@@ -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,12 +19,17 @@ 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(
@@ -31,7 +41,7 @@ class BanWordFilter(filter_model.ContentFilter):
return entities.FilterResult(
level=entities.ResultLevel.MASKED if found else entities.ResultLevel.PASS,
replacement=message,
replacement=current,
user_notice='消息中存在不合适的内容, 请修改' if found else '',
console_notice='',
)
@@ -158,6 +158,18 @@ class ResponseWrapper(stage.PipelineStage):
result_type=entities.ResultType.CONTINUE,
new_query=query,
)
elif (
isinstance(result, provider_message.MessageChunk) and result.is_final and not result.tool_calls
):
# Final streaming chunk with no text content but
# possibly carrying sandbox outbox attachments.
reply_chain = platform_message.MessageChain([])
await self._append_outbound_attachments(query, reply_chain)
query.resp_message_chain.append(reply_chain)
yield entities.StageProcessResult(
result_type=entities.ResultType.CONTINUE,
new_query=query,
)
if result.tool_calls is not None and len(result.tool_calls) > 0: # 有函数调用
function_names = [tc.function.name for tc in result.tool_calls]
+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()}')
@@ -95,6 +95,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
+107 -19
View File
@@ -3,8 +3,10 @@ import typing
import asyncio
import time
import traceback
import base64
import datetime
import langbot_plugin.api.definition.abstract.platform.adapter as abstract_platform_adapter
import langbot_plugin.api.entities.builtin.platform.message as platform_message
import langbot_plugin.api.entities.builtin.platform.events as platform_events
@@ -24,11 +26,24 @@ from langbot.libs.wecom_ai_bot_api.ws_client import WecomBotWsClient
class WecomBotMessageConverter(abstract_platform_adapter.AbstractMessageConverter):
@staticmethod
async def yiri2target(message_chain: platform_message.MessageChain):
content = ''
"""Convert a MessageChain into a list of component dicts.
Each dict has a ``type`` key (``'text'``, ``'image'``,
``'voice'``, ``'file'``). Text items carry ``text``; media
items carry ``base64`` (may include a ``data:...;base64,``
prefix) and optionally ``name``.
"""
items: list[dict] = []
for msg in message_chain:
if type(msg) is platform_message.Plain:
content += msg.text
return content
items.append({'type': 'text', 'text': msg.text})
elif type(msg) is platform_message.Image:
items.append({'type': 'image', 'base64': msg.base64 or ''})
elif type(msg) is platform_message.Voice:
items.append({'type': 'voice', 'base64': msg.base64 or ''})
elif type(msg) is platform_message.File:
items.append({'type': 'file', 'base64': msg.base64 or '', 'name': msg.name or ''})
return items
@staticmethod
async def target2yiri(event: WecomBotEvent, bot_name: str = ''):
@@ -362,13 +377,76 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
}
)
@staticmethod
def _join_text_components(items: list[dict]) -> str:
"""Concatenate ``text`` items in order, leaving media items alone."""
return ''.join(item['text'] for item in items if item.get('type') == 'text')
@staticmethod
def _iter_media_components(items: list[dict]):
"""Yield non-text items in order."""
for item in items:
if item.get('type') in {'image', 'voice', 'file'}:
yield item
@staticmethod
async def _send_media(
bot,
req_id: str,
item: dict,
) -> bool:
"""Upload *item* to the WeCom AI Bot CDN and send it as a media reply.
Returns True on success. Falls back to a no-op (with a warning log)
if the SDK does not yet implement ``upload_media`` /
``reply_image`` / ``reply_file`` / ``reply_voice`` the framework
will keep working, just without image delivery.
"""
kind = item.get('type')
upload = getattr(bot, 'upload_media', None)
if upload is None:
return False
b64_text = item.get('base64') or ''
if not b64_text:
return False
if b64_text.startswith('data:') and ',' in b64_text:
b64_text = b64_text.split(',', 1)[1]
try:
data = base64.b64decode(b64_text, validate=False)
except Exception:
return False
if not data:
return False
try:
upload_result = await upload(data, item.get('name') or f'attachment.{kind}', media_type=kind)
except Exception:
return False
media_id = getattr(upload_result, 'media_id', None) or (
isinstance(upload_result, dict) and upload_result.get('media_id')
)
if not media_id:
return False
reply_fn = {
'image': getattr(bot, 'reply_image', None),
'file': getattr(bot, 'reply_file', None),
'voice': getattr(bot, 'reply_voice', None),
}.get(kind)
if reply_fn is None:
return False
try:
await reply_fn(req_id, media_id)
return True
except Exception:
return False
async def reply_message(
self,
message_source: platform_events.MessageEvent,
message: platform_message.MessageChain,
quote_origin: bool = False,
):
content = await self.message_converter.yiri2target(message)
items = await self.message_converter.yiri2target(message)
text = self._join_text_components(items)
_ws_mode = not self.config.get('enable-webhook', False)
event = message_source.source_platform_object
@@ -382,7 +460,7 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
else:
chat_id = str(message_source.sender.id)
try:
await self.bot.send_message(chat_id, content)
await self.bot.send_message(chat_id, text)
except Exception:
await self.logger.error(
f'WeComBot: proactive reply for synthetic event failed: {traceback.format_exc()}'
@@ -396,12 +474,15 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
if _ws_mode:
req_id = event.get('req_id', '') if isinstance(event, dict) else getattr(event, 'req_id', '')
if req_id:
await self.bot.reply_text(req_id, content)
else:
await self.bot.set_message(event.message_id, content)
if text:
if req_id:
await self.bot.reply_text(req_id, text)
else:
await self.bot.set_message(event.message_id, text)
for item in self._iter_media_components(items):
await self._send_media(self.bot, req_id, item)
else:
await self.bot.set_message(event.message_id, content)
await self.bot.set_message(event.message_id, text)
async def reply_message_chunk(
self,
@@ -411,7 +492,8 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
quote_origin: bool = False,
is_final: bool = False,
):
content = await self.message_converter.yiri2target(message)
items = await self.message_converter.yiri2target(message)
text = self._join_text_components(items)
_ws_mode = not self.config.get('enable-webhook', False)
# Synthetic events (e.g. button-click triggered form resume) have
@@ -420,7 +502,7 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
# of the stream/reply path.
spo = message_source.source_platform_object
if spo is None:
return await self._handle_synthetic_chunk(message_source, bot_message, content, is_final, _ws_mode)
return await self._handle_synthetic_chunk(message_source, bot_message, text, is_final, _ws_mode)
msg_id = spo.message_id
@@ -452,7 +534,7 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
form_data.get('actions', []) or [],
)
except Exception:
fallback = content or '(人工输入)'
fallback = text or '(人工输入)'
if _ws_mode:
event = message_source.source_platform_object
req_id = event.get('req_id', '') if isinstance(event, dict) else getattr(event, 'req_id', '')
@@ -463,17 +545,22 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
return {'stream': False, 'form': True, 'fallback': True}
if _ws_mode:
success = await self.bot.push_stream_chunk(msg_id, content, is_final=is_final)
success = await self.bot.push_stream_chunk(msg_id, text, is_final=is_final)
if not success and is_final:
event = message_source.source_platform_object
req_id = event.get('req_id', '')
if req_id:
await self.bot.reply_text(req_id, content)
await self.bot.reply_text(req_id, text)
if is_final:
event = message_source.source_platform_object
req_id = event.get('req_id', '')
for item in self._iter_media_components(items):
await self._send_media(self.bot, req_id, item)
return {'stream': success}
else:
success = await self.bot.push_stream_chunk(msg_id, content, is_final=is_final)
success = await self.bot.push_stream_chunk(msg_id, text, is_final=is_final)
if not success and is_final:
await self.bot.set_message(msg_id, content)
await self.bot.set_message(msg_id, text)
return {'stream': success}
async def is_stream_output_supported(self) -> bool:
@@ -627,8 +714,9 @@ class WecomBotAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
async def send_message(self, target_type, target_id, message):
_ws_mode = not self.config.get('enable-webhook', False)
if _ws_mode:
content = await self.message_converter.yiri2target(message)
await self.bot.send_message(target_id, content)
items = await self.message_converter.yiri2target(message)
text = self._join_text_components(items)
await self.bot.send_message(target_id, text)
else:
pass
+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"""
+14 -2
View File
@@ -701,7 +701,13 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
}
self._known_desired_states.update({state.binding.installation_uuid: state for state in desired_states})
result = await runtime_handler.reconcile_plugin_installations(tuple(self._known_desired_states.values()))
reconcile_timeout_seconds = max(
300.0, self._runtime_connect_timeout(self.ap.instance_config.data.get('plugin', {}))
)
result = await runtime_handler.reconcile_plugin_installations(
tuple(self._known_desired_states.values()),
timeout=reconcile_timeout_seconds,
)
await self._repair_reconcile_missing_artifacts(self._known_desired_states, result)
self._record_reconcile_failures(self._known_desired_states, result)
@@ -736,7 +742,13 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
if state.binding.installation_uuid in all_states:
raise ValueError('Duplicate plugin installation UUID across projected Workspaces')
all_states[state.binding.installation_uuid] = state
result = await runtime_handler.reconcile_plugin_installations(tuple(all_states.values()))
reconcile_timeout_seconds = max(
300.0, self._runtime_connect_timeout(self.ap.instance_config.data.get('plugin', {}))
)
result = await runtime_handler.reconcile_plugin_installations(
tuple(all_states.values()),
timeout=reconcile_timeout_seconds,
)
await self._repair_reconcile_missing_artifacts(all_states, result)
self._record_reconcile_failures(all_states, result)
for installation_uuid, previous in tuple(self._known_desired_states.items()):
+120 -14
View File
@@ -11,6 +11,8 @@ import traceback
from dataclasses import dataclass
import sqlalchemy
import sqlalchemy.dialects.postgresql
import sqlalchemy.dialects.sqlite
from langbot_plugin.runtime.io import handler
from langbot_plugin.runtime.io.connection import Connection
@@ -431,6 +433,19 @@ class RuntimeConnectionHandler(handler.Handler):
return f'{identity.plugin_author}/{identity.plugin_name}'
raise ValueError(f'Unsupported binary storage owner_type {owner_type!r}')
@staticmethod
def _legacy_binary_storage_key(
action_context: ActionContext,
*,
owner_type: str,
owner: str,
key: str,
) -> str:
"""Return the pre-tenancy key shape for a row already scoped to this Workspace."""
legacy_owner = action_context.workspace_uuid if owner_type == 'workspace' else owner
return f'{owner_type}:{legacy_owner}:{key}'
@classmethod
def _binary_storage_key(
cls,
@@ -896,25 +911,82 @@ class RuntimeConnectionHandler(handler.Handler):
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == unique_key)
)
storage = result.first()
if storage is None:
legacy_key = self._legacy_binary_storage_key(
action_context,
owner_type=owner_type,
owner=owner,
key=key,
)
result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(persistence_bstorage.BinaryStorage)
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == legacy_key)
.where(persistence_bstorage.BinaryStorage.key == key)
.where(persistence_bstorage.BinaryStorage.owner_type == owner_type)
.where(persistence_bstorage.BinaryStorage.owner == owner)
)
storage = result.first()
if storage is not None:
update_result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.update(persistence_bstorage.BinaryStorage)
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == legacy_key)
.where(persistence_bstorage.BinaryStorage.key == key)
.where(persistence_bstorage.BinaryStorage.owner_type == owner_type)
.where(persistence_bstorage.BinaryStorage.owner == owner)
.values(unique_key=unique_key, value=value)
)
if update_result.rowcount:
return handler.ActionResponse.success(data={})
canonical_update = await self.ap.persistence_mgr.execute_async(
sqlalchemy.update(persistence_bstorage.BinaryStorage)
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == unique_key)
.where(persistence_bstorage.BinaryStorage.key == key)
.where(persistence_bstorage.BinaryStorage.owner_type == owner_type)
.where(persistence_bstorage.BinaryStorage.owner == owner)
.values(value=value)
)
if canonical_update.rowcount:
return handler.ActionResponse.success(data={})
storage = None
if result.first() is not None:
if storage is not None:
await self.ap.persistence_mgr.execute_async(
sqlalchemy.update(persistence_bstorage.BinaryStorage)
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == unique_key)
.where(persistence_bstorage.BinaryStorage.key == key)
.where(persistence_bstorage.BinaryStorage.owner_type == owner_type)
.where(persistence_bstorage.BinaryStorage.owner == owner)
.values(value=value)
)
else:
await self.ap.persistence_mgr.execute_async(
sqlalchemy.insert(persistence_bstorage.BinaryStorage).values(
workspace_uuid=action_context.workspace_uuid,
unique_key=unique_key,
key=key,
owner_type=owner_type,
owner=owner,
value=value,
)
return handler.ActionResponse.success(data={})
dialect_name = self.ap.persistence_mgr.get_db_engine().dialect.name
insert = {
'postgresql': sqlalchemy.dialects.postgresql.insert,
'sqlite': sqlalchemy.dialects.sqlite.insert,
}.get(dialect_name)
if insert is None:
return handler.ActionResponse.error(message=f'Unsupported storage database dialect: {dialect_name}')
await self.ap.persistence_mgr.execute_async(
insert(persistence_bstorage.BinaryStorage)
.values(
workspace_uuid=action_context.workspace_uuid,
unique_key=unique_key,
key=key,
owner_type=owner_type,
owner=owner,
value=value,
)
.on_conflict_do_update(
index_elements=['workspace_uuid', 'unique_key'],
set_={'value': value},
)
)
return handler.ActionResponse.success(
data={},
@@ -946,6 +1018,29 @@ class RuntimeConnectionHandler(handler.Handler):
)
storage = result.first()
if storage is None:
legacy_key = self._legacy_binary_storage_key(
action_context,
owner_type=owner_type,
owner=owner,
key=key,
)
result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(persistence_bstorage.BinaryStorage)
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == legacy_key)
.where(persistence_bstorage.BinaryStorage.key == key)
.where(persistence_bstorage.BinaryStorage.owner_type == owner_type)
.where(persistence_bstorage.BinaryStorage.owner == owner)
)
storage = result.first()
if storage is None:
retry_result = await self.ap.persistence_mgr.execute_async(
sqlalchemy.select(persistence_bstorage.BinaryStorage)
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == unique_key)
)
storage = retry_result.first()
if storage is None:
return handler.ActionResponse.error(
message=f'Storage with key {key} not found',
@@ -981,10 +1076,19 @@ class RuntimeConnectionHandler(handler.Handler):
message=str(e),
)
legacy_key = self._legacy_binary_storage_key(
action_context,
owner_type=owner_type,
owner=owner,
key=key,
)
await self.ap.persistence_mgr.execute_async(
sqlalchemy.delete(persistence_bstorage.BinaryStorage)
.where(persistence_bstorage.BinaryStorage.workspace_uuid == action_context.workspace_uuid)
.where(persistence_bstorage.BinaryStorage.unique_key == unique_key)
.where(persistence_bstorage.BinaryStorage.unique_key.in_((unique_key, legacy_key)))
.where(persistence_bstorage.BinaryStorage.key == key)
.where(persistence_bstorage.BinaryStorage.owner_type == owner_type)
.where(persistence_bstorage.BinaryStorage.owner == owner)
)
return handler.ActionResponse.success(
@@ -1012,7 +1116,7 @@ class RuntimeConnectionHandler(handler.Handler):
return handler.ActionResponse.success(
data={
'keys': result.scalars().all(),
'keys': list(dict.fromkeys(result.scalars().all())),
},
)
@@ -1573,13 +1677,15 @@ class RuntimeConnectionHandler(handler.Handler):
async def reconcile_plugin_installations(
self,
installations: tuple[PluginInstallationDesiredState, ...],
*,
timeout: float = 300,
) -> dict[str, Any]:
request = ReconcilePluginInstallationsRequest(installations=installations)
with self.installation_scope(None):
return await self.call_action(
LangBotToRuntimeAction.RECONCILE_PLUGIN_INSTALLATIONS,
request.model_dump(),
timeout=300,
timeout=timeout,
)
async def apply_plugin_installation(
@@ -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
@@ -24,7 +24,10 @@ class SeekDBEmbedding(requester.ProviderAPIRequester):
try:
import pyseekdb
except ImportError:
raise ImportError('pyseekdb is not installed. Install it with: pip install pyseekdb')
raise ImportError(
"SeekDB support is not installed. Install LangBot with the 'seekdb' extra: "
"uv sync --extra seekdb (source) or uvx --from 'langbot[seekdb]@latest' langbot (PyPI)."
)
self._embedding_function = pyseekdb.get_default_embedding_function()
@@ -619,7 +619,9 @@ class LocalAgentRunner(runner.RequestRunner):
and len(func_ret) > 0
and isinstance(func_ret[0], provider_message.ContentElement)
):
tool_content = func_ret
# OpenAI-compatible APIs require tool-message content to be a
# string; a raw list of ContentElement causes HTTP 500 (#2457).
tool_content = '\n'.join(str(ce) for ce in func_ret)
else:
tool_content = json.dumps(func_ret, ensure_ascii=False)
+13 -4
View File
@@ -27,10 +27,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')
@@ -115,8 +121,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')
@@ -162,6 +169,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."""
@@ -174,4 +182,5 @@ async def mask_patterns(
mask=mask,
mask_word=mask_word,
timeout_seconds=timeout_seconds,
max_pattern_count=max_pattern_count,
)
+4 -1
View File
@@ -42,7 +42,10 @@ class SeekDBVectorDatabase(VectorDatabase):
def __init__(self, ap: app.Application):
if not SEEKDB_AVAILABLE:
raise ImportError('pyseekdb is not installed. Install it with: pip install pyseekdb')
raise ImportError(
"SeekDB support is not installed. Install LangBot with the 'seekdb' extra: "
"uv sync --extra seekdb (source) or uvx --from 'langbot[seekdb]@latest' langbot (PyPI)."
)
self.ap = ap
config = self.ap.instance_config.data['vdb']['seekdb']
@@ -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>'''
+5
View File
@@ -181,6 +181,11 @@ vdb:
host: localhost
port: 6333
api_key: ''
# SeekDB is optional. Native/package installs need the `seekdb` extra:
# `uv sync --extra seekdb` (source) or
# `uvx --from 'langbot[seekdb]@latest' langbot` (PyPI).
# The official Docker image already includes it.
# Embedded-mode platform support depends on the native pylibseekdb wheels.
seekdb:
mode: embedded # 'embedded' or 'server'
# Embedded mode options:
+23
View File
@@ -307,6 +307,7 @@ class TestUserInitEndpoint:
assert data['data'] == {
'initialized': True,
'authenticated_invitation_acceptance_enabled': False,
'invitation_registration_enabled': True,
'password_login_enabled': True,
'space_login_enabled': False,
}
@@ -330,6 +331,28 @@ class TestUserInitEndpoint:
assert data['data'] == {
'initialized': True,
'authenticated_invitation_acceptance_enabled': True,
'invitation_registration_enabled': False,
'password_login_enabled': False,
'space_login_enabled': True,
}
@pytest.mark.asyncio
async def test_account_info_enables_local_invitation_registration_for_oauth_only_oss(
self, quart_test_client, fake_api_app
):
fake_api_app.user_service.is_initialized.return_value = True
fake_api_app.user_service.get_login_capabilities = AsyncMock(
return_value={'password_login_enabled': False, 'space_login_enabled': True}
)
response = await quart_test_client.get('/api/v1/user/account-info')
assert response.status_code == 200
data = await response.get_json()
assert data['data'] == {
'initialized': True,
'authenticated_invitation_acceptance_enabled': False,
'invitation_registration_enabled': True,
'password_login_enabled': False,
'space_login_enabled': True,
}
@@ -312,6 +312,29 @@ async def test_space_credits_are_resolved_from_workspace_owner(space_oauth_api):
application.space_service.get_credits.assert_awaited_once_with('owner@example.com')
@pytest.mark.asyncio
async def test_oss_local_only_owner_requires_space_binding_for_langbot_models(space_oauth_api):
application, client = space_oauth_api
application.user_service.get_workspace_owner = AsyncMock(
return_value=SimpleNamespace(user='owner@example.com', space_account_uuid=None)
)
application.space_service.get_credits = AsyncMock()
response = await client.get(
'/api/v1/user/space-credits',
headers={'Authorization': 'Bearer account-token', 'X-Workspace-Id': WORKSPACE_UUID},
)
payload = await response.get_json()
assert response.status_code == 200
assert payload['data'] == {
'credits': None,
'owner_space_bound': False,
'is_workspace_owner': True,
}
application.space_service.get_credits.assert_not_awaited()
@pytest.mark.asyncio
async def test_cloud_workspace_owner_is_always_space_bound_after_login(space_oauth_api):
application, client = space_oauth_api
@@ -81,6 +81,7 @@ async def create_legacy_resource_schema(engine, *, instance_uuid: str) -> None:
sa.Column('key', sa.String(255), nullable=False),
sa.Column('owner_type', sa.String(255), nullable=False),
sa.Column('owner', sa.String(255), nullable=False),
sa.Column('value', sa.LargeBinary, nullable=False),
)
mcp_servers = _uuid_table(
metadata,
@@ -210,7 +211,13 @@ async def create_legacy_resource_schema(engine, *, instance_uuid: str) -> None:
await conn.execute(bots.insert().values(uuid='bot-1', name='bot', updated_at=now))
await conn.execute(bot_admins.insert().values(bot_uuid='bot-1', launcher_type='person', launcher_id='owner'))
await conn.execute(
binary_storages.insert().values(unique_key='plugin:demo:key', key='key', owner_type='plugin', owner='demo')
binary_storages.insert().values(
unique_key='plugin:demo:key',
key='key',
owner_type='plugin',
owner='demo',
value=b'legacy-plugin-value',
)
)
await conn.execute(mcp_servers.insert().values(uuid='mcp-1', name='shared-name', enable=True, updated_at=now))
await conn.execute(model_providers.insert().values(uuid='provider-1', name='provider', requester='openai'))
@@ -76,6 +76,26 @@ async def test_legacy_sqlite_resources_are_backfilled_and_contracted(tmp_path):
)
assert legacy_kb['collection_id'] == 'collection-1'
assert legacy_kb['legacy_vector_collection'] == 1
legacy_binary_storage = (
(
await conn.execute(
sa.text(
'SELECT workspace_uuid, unique_key, key, owner_type, owner, value '
"FROM binary_storages WHERE owner_type = 'plugin' AND owner = 'demo'"
)
)
)
.mappings()
.one()
)
assert legacy_binary_storage == {
'workspace_uuid': workspace_uuid,
'unique_key': 'plugin:demo:key',
'key': 'key',
'owner_type': 'plugin',
'owner': 'demo',
'value': b'legacy-plugin-value',
}
assert (
await conn.scalar(
sa.text(
@@ -209,8 +229,8 @@ async def test_sqlite_scoped_keys_allow_cross_workspace_but_reject_same_workspac
await conn.execute(
sa.text(
'INSERT INTO binary_storages '
'(workspace_uuid, unique_key, key, owner_type, owner) '
"VALUES (:workspace_uuid, 'plugin:demo:key', 'key', 'plugin', 'demo')"
'(workspace_uuid, unique_key, key, owner_type, owner, value) '
"VALUES (:workspace_uuid, 'plugin:demo:key', 'key', 'plugin', 'demo', X'')"
),
{'workspace_uuid': second_workspace_uuid},
)
+48 -18
View File
@@ -2163,25 +2163,38 @@ class TestInboundOutboundRoundTrip:
calls = []
async def fake_execute_tool(parameters, q):
calls.append(parameters['command'])
if 'os.scandir' in parameters['command']:
return {
'ok': True,
'stdout': '[{"name": "out.png", "b64": "QUJD"}]',
'stderr': '',
}
async def fake_client_execute(spec):
cmd = spec.cmd
calls.append(cmd)
if 'os.scandir' in cmd:
return BoxExecutionResult(
session_id='s',
backend_name='test',
status=BoxExecutionStatus.COMPLETED,
exit_code=0,
stdout='[{"name": "out.png", "b64": "QUJD"}]',
duration_ms=10,
)
# the rm -rf cleanup call
return {'ok': True, 'stdout': '', 'stderr': ''}
return BoxExecutionResult(
session_id='s',
backend_name='test',
status=BoxExecutionStatus.COMPLETED,
exit_code=0,
stdout='',
duration_ms=10,
)
service.execute_tool = AsyncMock(side_effect=fake_execute_tool)
service.client.execute = AsyncMock(side_effect=fake_client_execute)
service.execute_tool = AsyncMock(return_value={'ok': True, 'stdout': '', 'stderr': ''})
attachments = await service.collect_outbound_attachments(query)
assert len(attachments) == 1
assert attachments[0]['type'] == 'Image'
assert attachments[0]['name'] == 'out.png'
# cleanup (rm -rf) must have been issued after a successful collection
assert any('rm -rf' in c for c in calls)
service.execute_tool.assert_awaited_once()
assert 'rm -rf' in service.execute_tool.await_args.args[0]['command']
@pytest.mark.asyncio
async def test_collect_outbound_empty_still_clears(self):
@@ -2193,16 +2206,33 @@ class TestInboundOutboundRoundTrip:
calls = []
async def fake_execute_tool(parameters, q):
calls.append(parameters['command'])
if 'os.scandir' in parameters['command']:
return {'ok': True, 'stdout': '[]', 'stderr': ''}
return {'ok': True, 'stdout': '', 'stderr': ''}
async def fake_client_execute(spec):
cmd = spec.cmd
calls.append(cmd)
if 'os.scandir' in cmd:
return BoxExecutionResult(
session_id='s',
backend_name='test',
status=BoxExecutionStatus.COMPLETED,
exit_code=0,
stdout='[]',
duration_ms=10,
)
return BoxExecutionResult(
session_id='s',
backend_name='test',
status=BoxExecutionStatus.COMPLETED,
exit_code=0,
stdout='',
duration_ms=10,
)
service.execute_tool = AsyncMock(side_effect=fake_execute_tool)
service.client.execute = AsyncMock(side_effect=fake_client_execute)
service.execute_tool = AsyncMock(return_value={'ok': True, 'stdout': '', 'stderr': ''})
assert await service.collect_outbound_attachments(query) == []
# cleanup (rm -rf) is issued unconditionally now
assert any('rm -rf' in c for c in calls)
service.execute_tool.assert_awaited_once()
assert 'rm -rf' in service.execute_tool.await_args.args[0]['command']
@pytest.mark.asyncio
async def test_passthrough_noop_when_unavailable(self):
@@ -144,3 +144,39 @@ async def test_runtime_resource_stats_are_aggregate_and_constant_time() -> None:
assert stats['models']['providers'] == 1
assert stats['runtimes']['plugin_installations'] == 1
assert stats['runtimes']['plugin_runtime_connected'] is True
@pytest.mark.asyncio
async def test_start_plugin_runtime_initialization_bypasses_after_commit_gate() -> None:
app = Application()
app.plugin_connector = SimpleNamespace(initialize=AsyncMock())
app.task_mgr = SimpleNamespace(create_task=AsyncMock())
task = app._start_plugin_runtime_initialization()
await task
app.plugin_connector.initialize.assert_awaited_once_with()
app.task_mgr.create_task.assert_not_called()
@pytest.mark.asyncio
async def test_shutdown_cancels_plugin_runtime_initialization_task() -> None:
app = Application()
app._plugin_runtime_initialization_task = asyncio.create_task(asyncio.sleep(60))
app.task_mgr = SimpleNamespace(cancel_by_scope=lambda *_: None, tasks=[])
app.event_loop_monitor = SimpleNamespace(stop=AsyncMock())
app.http_ctrl = SimpleNamespace(mcp_mount=None)
app.platform_mgr = None
app.tool_mgr = None
app.model_mgr = None
app.box_service = None
app.plugin_connector = None
app.telemetry = None
app.vector_db_mgr = None
app.storage_mgr = None
app.persistence_mgr = SimpleNamespace(db=SimpleNamespace(engine=SimpleNamespace(dispose=AsyncMock())))
app.deployment = None
await app.shutdown()
assert app._plugin_runtime_initialization_task.cancelled()
+113
View File
@@ -0,0 +1,113 @@
"""BanWordFilter regression tests for legacy sensitive-word lists.
v4.10.7 introduced a 64-pattern cap in safe_regex. Older installs still carry
the previous default list (~70 patterns). The filter must keep applying those
rules instead of blocking every message.
"""
from __future__ import annotations
from importlib import import_module
from unittest.mock import Mock
import pytest
from tests.factories import FakeApp
def _load_banwords():
import_module('langbot.pkg.pipeline.pipelinemgr')
banwords = import_module('langbot.pkg.pipeline.cntfilter.filters.banwords')
entities = import_module('langbot.pkg.pipeline.cntfilter.entities')
safe_regex = import_module('langbot.pkg.utils.safe_regex')
return banwords, entities, safe_regex
def _filter_with_words(words: list[str], *, mask: str = '*', mask_word: str = ''):
banwords, entities, _ = _load_banwords()
app = FakeApp()
app.sensitive_meta = Mock()
app.sensitive_meta.data = {
'words': words,
'mask': mask,
'mask_word': mask_word,
}
return banwords.BanWordFilter(app), entities, app
@pytest.mark.asyncio
async def test_legacy_word_list_over_pattern_cap_does_not_block_clean_message():
"""A pre-v4.10.7 word list must not fail closed on every message."""
_, _, safe_regex = _load_banwords()
words = [f'word{i}' for i in range(safe_regex.MAX_PATTERN_COUNT + 6)]
filt, entities, _ = _filter_with_words(words)
result = await filt.process(Mock(), 'hello there, nothing banned')
assert result.level == entities.ResultLevel.PASS
assert result.replacement == 'hello there, nothing banned'
assert result.user_notice == ''
@pytest.mark.asyncio
async def test_legacy_word_list_still_masks_match_beyond_first_batch():
"""Words past the first 64-pattern batch must still be applied."""
_, _, safe_regex = _load_banwords()
words = [f'word{i}' for i in range(safe_regex.MAX_PATTERN_COUNT)] + ['secret-token']
filt, entities, _ = _filter_with_words(words, mask_word='[hidden]')
result = await filt.process(Mock(), 'please hide secret-token now')
assert result.level == entities.ResultLevel.MASKED
assert 'secret-token' not in result.replacement
assert '[hidden]' in result.replacement
@pytest.mark.asyncio
async def test_legacy_word_list_masks_match_in_first_batch():
_, _, safe_regex = _load_banwords()
words = ['alpha-secret'] + [f'word{i}' for i in range(safe_regex.MAX_PATTERN_COUNT)]
filt, entities, _ = _filter_with_words(words, mask_word='[hidden]')
result = await filt.process(Mock(), 'alpha-secret is here')
assert result.level == entities.ResultLevel.MASKED
assert result.replacement == '[hidden] is here'
@pytest.mark.asyncio
async def test_invalid_sensitive_word_regex_still_blocks():
filt, entities, _ = _filter_with_words(['(unclosed'])
result = await filt.process(Mock(), 'any message')
assert result.level == entities.ResultLevel.BLOCK
assert result.user_notice == '内容检查规则执行失败,请联系管理员'
assert 'rejected' in result.console_notice.lower() or 'invalid' in result.console_notice.lower()
@pytest.mark.asyncio
async def test_oversized_word_list_is_blocked():
"""Configured rules must never be silently skipped when the list is oversized."""
banwords, _, _ = _load_banwords()
words = [f'word{i}' for i in range(banwords._MAX_SENSITIVE_WORD_PATTERNS + 10)]
filt, entities, _ = _filter_with_words(words)
result = await filt.process(Mock(), 'hello there, nothing banned')
assert result.level == entities.ResultLevel.BLOCK
assert result.replacement == ''
assert result.user_notice == '内容检查规则执行失败,请联系管理员'
assert 'at most 256 regex patterns are allowed' in result.console_notice.lower()
@pytest.mark.asyncio
async def test_match_beyond_total_cap_cannot_bypass_filter():
banwords, _, _ = _load_banwords()
words = [f'word{i}' for i in range(banwords._MAX_SENSITIVE_WORD_PATTERNS)] + ['late-secret']
filt, entities, _ = _filter_with_words(words, mask_word='[hidden]')
result = await filt.process(Mock(), 'please hide late-secret now')
assert result.level == entities.ResultLevel.BLOCK
assert result.replacement == ''
@@ -0,0 +1,259 @@
from __future__ import annotations
import pytest
from unittest.mock import MagicMock
from linebot.v3.webhooks import TextMessageContent, UserMentionee, AllMentionee
from langbot.pkg.platform import botmgr as _botmgr # noqa: F401
from langbot.pkg.platform.sources import line
import langbot_plugin.api.entities.builtin.platform.message as platform_message
BOT_ACCOUNT_ID = 'line-bot-account'
def _make_event(
*, source_type: str, user_id, group_id=None, room_id=None, message_id: str, text: str = 'hi', mention=None
):
event = MagicMock()
event.timestamp = 1700000000000
message = MagicMock(spec=TextMessageContent)
message.id = message_id
message.text = text
message.mention = mention
event.message = message
event.message.webhook_event_id = f'webhook-{message_id}'
event.message.timestamp = event.timestamp
source = MagicMock()
source.type = source_type
source.user_id = user_id
if group_id is not None:
source.group_id = group_id
if room_id is not None:
source.room_id = room_id
event.source = source
return event
def _make_converter(bot_account_id: str = BOT_ACCOUNT_ID) -> line.LINEEventConverter:
return line.LINEEventConverter(bot_account_id=bot_account_id)
@pytest.mark.asyncio
async def test_user_message_launcher_id_stable_across_messages() -> None:
"""Two distinct messages from the same LINE user must resolve to the same
sender id, otherwise every message starts a brand new session (context loss).
"""
converter = _make_converter()
event1 = _make_event(source_type='user', user_id='U-stable-user', message_id='msg-1')
event2 = _make_event(source_type='user', user_id='U-stable-user', message_id='msg-2')
result1 = await converter.target2yiri(event1, bot_client=None)
result2 = await converter.target2yiri(event2, bot_client=None)
assert result1.sender.id == 'U-stable-user'
assert result1.sender.id == result2.sender.id
assert result1.sender.id != event1.message.id
@pytest.mark.asyncio
async def test_group_message_uses_group_id_not_message_id() -> None:
converter = _make_converter()
event1 = _make_event(source_type='group', user_id='U-member', group_id='G-stable-group', message_id='msg-1')
event2 = _make_event(source_type='group', user_id='U-member', group_id='G-stable-group', message_id='msg-2')
result1 = await converter.target2yiri(event1, bot_client=None)
result2 = await converter.target2yiri(event2, bot_client=None)
assert result1.sender.group.id == 'G-stable-group'
assert result1.sender.group.id == result2.sender.group.id
assert result1.sender.id == 'U-member'
@pytest.mark.asyncio
async def test_room_message_uses_room_id_and_falls_back_when_user_id_missing() -> None:
converter = _make_converter()
event = _make_event(source_type='room', user_id=None, room_id='R-stable-room', message_id='msg-1')
result = await converter.target2yiri(event, bot_client=None)
assert result.sender.group.id == 'R-stable-room'
assert result.sender.id == 'R-stable-room'
def _plain_texts(chain: platform_message.MessageChain) -> list[str]:
return [c.text for c in chain if isinstance(c, platform_message.Plain)]
def _ats(chain: platform_message.MessageChain) -> list[platform_message.At]:
return [c for c in chain if isinstance(c, platform_message.At)]
@pytest.mark.asyncio
async def test_no_mention_keeps_plain_text() -> None:
converter = _make_converter()
event = _make_event(source_type='group', user_id='U-member', group_id='G1', message_id='m1', text='hello world')
chain = await converter.message_converter.target2yiri(event, bot_client=None)
assert _plain_texts(chain) == ['hello world']
assert _ats(chain) == []
@pytest.mark.asyncio
async def test_bot_mention_maps_to_at_with_bot_account_id() -> None:
"""A @bot mention must become At(target=bot_account_id) so the 'at-bot'
group respond rule matches (previously the mention was lost and the message
was silently dropped in groups with at-only rules).
"""
mention = MagicMock()
mention.mentionees = [
UserMentionee(type='user', index=0, length=4, userId='U-bot-user-id', isSelf=True),
]
converter = _make_converter()
event = _make_event(
source_type='group',
user_id='U-member',
group_id='G1',
message_id='m1',
text='@BOT hey',
mention=mention,
)
chain = await converter.message_converter.target2yiri(event, bot_client=None)
ats = _ats(chain)
assert len(ats) == 1
assert ats[0].target == BOT_ACCOUNT_ID
assert _plain_texts(chain) == [' hey']
@pytest.mark.asyncio
async def test_other_user_mention_keeps_display_text() -> None:
"""Mentions of other users keep their display text in the message string,
so prefix/regexp rules that match the raw '@Name ...' text still work.
"""
mention = MagicMock()
mention.mentionees = [
UserMentionee(type='user', index=0, length=6, userId='U-other', isSelf=False),
]
converter = _make_converter()
event = _make_event(
source_type='group',
user_id='U-member',
group_id='G1',
message_id='m1',
text='@Alice hello',
mention=mention,
)
chain = await converter.message_converter.target2yiri(event, bot_client=None)
ats = _ats(chain)
assert len(ats) == 1
assert ats[0].target == 'U-other'
# str() of the At component falls back to display when set
assert str(chain) == '@Alice hello'
@pytest.mark.asyncio
async def test_bot_mention_triggers_atbot_rule() -> None:
"""End-to-end: a group message that @mentions the bot must be accepted by
the at-bot respond rule (this is the regression that silently dropped
'@bot' messages in LINE groups).
"""
from langbot.pkg.pipeline.resprule.rules.atbot import AtBotRule
mention = MagicMock()
mention.mentionees = [
UserMentionee(type='user', index=0, length=6, userId='U-bot-user-id', isSelf=True),
]
converter = _make_converter()
event = _make_event(
source_type='group',
user_id='U-member',
group_id='G1',
message_id='m1',
text='@RAIQt hi',
mention=mention,
)
chain = await converter.message_converter.target2yiri(event, bot_client=None)
query = MagicMock()
query.adapter = MagicMock()
query.adapter.bot_account_id = BOT_ACCOUNT_ID
rule = AtBotRule(ap=MagicMock())
result = await rule.match(str(chain), chain, {'at': True}, query)
assert result.matching is True
@pytest.mark.asyncio
async def test_group_without_bot_mention_still_dropped_by_atbot_rule() -> None:
from langbot.pkg.pipeline.resprule.rules.atbot import AtBotRule
converter = _make_converter()
event = _make_event(source_type='group', user_id='U-member', group_id='G1', message_id='m1', text='hello')
chain = await converter.message_converter.target2yiri(event, bot_client=None)
query = MagicMock()
query.adapter = MagicMock()
query.adapter.bot_account_id = BOT_ACCOUNT_ID
rule = AtBotRule(ap=MagicMock())
result = await rule.match(str(chain), chain, {'at': True}, query)
assert result.matching is False
@pytest.mark.asyncio
async def test_at_all_mention_preserved_as_at_component() -> None:
mention = MagicMock()
mention.mentionees = [
AllMentionee(type='all', index=0, length=4),
]
converter = _make_converter()
event = _make_event(
source_type='group',
user_id='U-member',
group_id='G1',
message_id='m1',
text='@All hello',
mention=mention,
)
chain = await converter.message_converter.target2yiri(event, bot_client=None)
ats = _ats(chain)
assert len(ats) == 1
assert str(chain) == '@All hello'
@pytest.mark.asyncio
async def test_multiple_mentions_sorted_by_position() -> None:
mention = MagicMock()
# Intentionally out of order to exercise sorting
mention.mentionees = [
UserMentionee(type='user', index=9, length=4, userId='U-b', isSelf=False),
UserMentionee(type='user', index=0, length=4, userId='U-a', isSelf=False),
]
converter = _make_converter()
event = _make_event(
source_type='group',
user_id='U-member',
group_id='G1',
message_id='m1',
text='@aaa mid @bbb tail',
mention=mention,
)
chain = await converter.message_converter.target2yiri(event, bot_client=None)
ats = _ats(chain)
assert [a.target for a in ats] == ['U-a', 'U-b']
assert str(chain) == '@aaa mid @bbb tail'
@@ -1,9 +1,11 @@
"""Tests for QQ Official keyboard payload helpers."""
"""Tests for QQ Official message and keyboard payload helpers."""
import asyncio
import json
import time
from unittest.mock import AsyncMock, MagicMock, patch
import httpx
import pytest
import langbot_plugin.api.entities.builtin.platform.message as platform_message
@@ -99,6 +101,12 @@ def _stream_test_adapter():
adapter.bot = MagicMock()
adapter.bot.send_stream_msg = AsyncMock(return_value={'id': 'stream-1'})
adapter.bot.send_markdown_keyboard = AsyncMock(return_value={'id': 'message-1'})
adapter.bot.send_private_text_msg = AsyncMock()
adapter.bot.send_group_text_msg = AsyncMock()
adapter.bot.send_private_markdown_msg = AsyncMock()
adapter.bot.send_group_markdown_msg = AsyncMock()
adapter.bot.send_channle_group_text_msg = AsyncMock()
adapter.bot.send_channle_private_text_msg = AsyncMock()
adapter.ap = None
adapter._stream_ctx = {}
adapter._stream_ctx_ts = {}
@@ -108,7 +116,7 @@ def _stream_test_adapter():
@pytest.mark.asyncio
async def test_qq_stream_uses_cumulative_chunks_as_snapshots():
async def test_qq_stream_replace_mode_sends_complete_snapshots():
adapter = _stream_test_adapter()
adapter._stream_ctx['message-1'] = {
'user_openid': 'user-1',
@@ -138,10 +146,109 @@ async def test_qq_stream_uses_cumulative_chunks_as_snapshots():
assert [call.kwargs['content'] for call in adapter.bot.send_stream_msg.await_args_list] == [
'<think>one',
' two',
'<think>one two',
]
@pytest.mark.asyncio
async def test_qq_markdown_messages_use_markdown_payloads():
requests = []
def capture_request(request: httpx.Request) -> httpx.Response:
requests.append((str(request.url), json.loads(request.content)))
return httpx.Response(200, json={})
client = QQOfficialClient('secret', 'token', 'app-id', AsyncMock())
client.access_token = 'access-token'
client.access_token_expiry_time = time.time() + 3600
client._http_clients[None] = httpx.AsyncClient(transport=httpx.MockTransport(capture_request))
try:
await client.send_private_markdown_msg('user-1', '# Hello', msg_id='message-1', msg_seq=2)
await client.send_group_markdown_msg('group-1', '* Hello', event_id='event-1', msg_seq=3)
finally:
await client.close()
assert requests == [
(
'https://api.sgroup.qq.com/v2/users/user-1/messages',
{'msg_type': 2, 'markdown': {'content': '# Hello'}, 'msg_seq': 2, 'msg_id': 'message-1'},
),
(
'https://api.sgroup.qq.com/v2/groups/group-1/messages',
{'msg_type': 2, 'markdown': {'content': '* Hello'}, 'msg_seq': 3, 'event_id': 'event-1'},
),
]
@pytest.mark.asyncio
async def test_qq_markdown_rendering_switches_c2c_and_group_text_replies():
adapter = _stream_test_adapter()
adapter.config = {'enable-markdown-rendering': True}
await adapter._send_c2c_or_group_text_reply('c2c', 'user-1', '# Hello', msg_id='message-1')
await adapter._send_c2c_or_group_text_reply('group', 'group-1', '* Hello', event_id='event-1')
adapter.bot.send_private_markdown_msg.assert_awaited_once_with(
user_openid='user-1',
content='# Hello',
msg_id='message-1',
event_id=None,
msg_seq=1,
)
adapter.bot.send_group_markdown_msg.assert_awaited_once_with(
group_openid='group-1',
content='* Hello',
msg_id=None,
event_id='event-1',
msg_seq=1,
)
adapter.bot.send_private_text_msg.assert_not_awaited()
adapter.bot.send_group_text_msg.assert_not_awaited()
@pytest.mark.asyncio
async def test_qq_markdown_rendering_defaults_to_plain_text_replies():
adapter = _stream_test_adapter()
adapter.config = {}
await adapter._send_c2c_or_group_text_reply('c2c', 'user-1', 'Hello')
await adapter._send_c2c_or_group_text_reply('group', 'group-1', 'Hello')
adapter.bot.send_private_text_msg.assert_awaited_once()
adapter.bot.send_group_text_msg.assert_awaited_once()
adapter.bot.send_private_markdown_msg.assert_not_awaited()
adapter.bot.send_group_markdown_msg.assert_not_awaited()
@pytest.mark.asyncio
async def test_qq_markdown_rendering_does_not_affect_channel_messages():
adapter = _stream_test_adapter()
adapter.config = {'enable-markdown-rendering': True}
message = platform_message.MessageChain([platform_message.Plain(text='# Hello')])
channel_source = MagicMock()
channel_source.t = 'AT_MESSAGE_CREATE'
channel_source.channel_id = 'channel-1'
channel_source.d_id = 'message-1'
channel_event = MagicMock()
channel_event.source_platform_object = channel_source
await adapter.reply_message(channel_event, message)
dm_source = MagicMock()
dm_source.t = 'DIRECT_MESSAGE_CREATE'
dm_source.guild_id = 'guild-1'
dm_source.d_id = 'message-2'
dm_event = MagicMock()
dm_event.source_platform_object = dm_source
await adapter.reply_message(dm_event, message)
adapter.bot.send_channle_group_text_msg.assert_awaited_once_with('channel-1', '# Hello', 'message-1')
adapter.bot.send_channle_private_text_msg.assert_awaited_once_with('guild-1', '# Hello', 'message-2')
adapter.bot.send_private_markdown_msg.assert_not_awaited()
adapter.bot.send_group_markdown_msg.assert_not_awaited()
@pytest.mark.asyncio
async def test_qq_non_streaming_fallback_keeps_latest_snapshot_only():
from langbot.pkg.platform.sources.qqofficial import QQOfficialAdapter
@@ -0,0 +1,127 @@
import base64
import pytest
import langbot.pkg.core.app # noqa: F401
import langbot_plugin.api.entities.builtin.platform.message as platform_message
from langbot.libs.wecom_ai_bot_api.ws_client import _UPLOAD_CHUNK_SIZE, WecomBotWsClient
from langbot.pkg.platform.sources.wecombot import WecomBotAdapter, WecomBotMessageConverter
class Logger:
def __init__(self):
self.warnings = []
self.errors = []
async def warning(self, message):
self.warnings.append(message)
async def error(self, message):
self.errors.append(message)
async def info(self, message):
return None
class UploadClient(WecomBotWsClient):
def __init__(self):
super().__init__(bot_id='bot', secret='secret', logger=Logger())
self.frames = []
async def _send_reply(self, req_id: str, body: dict, cmd: str = 'aibot_respond_msg'):
self.frames.append((cmd, body))
if cmd == 'aibot_upload_media_init':
return {'errcode': 0, 'body': {'upload_id': 'upload-1'}}
if cmd == 'aibot_upload_media_finish':
return {'errcode': 0, 'body': {'media_id': 'media-1'}}
return {'errcode': 0}
class Bot:
def __init__(self):
self.calls = []
async def upload_media(self, data, filename='attachment', media_type='file'):
self.calls.append(('upload_media', media_type, filename, data))
return {'media_id': 'media-1'}
async def reply_text(self, req_id, content):
self.calls.append(('reply_text', req_id, content))
async def reply_image(self, req_id, media_id):
self.calls.append(('reply_image', req_id, media_id))
async def send_message(self, target_id, content):
self.calls.append(('send_message', target_id, content))
def make_adapter(bot):
return WecomBotAdapter.model_construct(
bot=bot,
config={'enable-webhook': False},
logger=Logger(),
message_converter=WecomBotMessageConverter(),
)
@pytest.mark.asyncio
async def test_ws_client_upload_media_uses_chunk_protocol():
client = UploadClient()
data = b'a' * (_UPLOAD_CHUNK_SIZE + 1)
upload_result = await client.upload_media(data, 'image.png', media_type='image')
assert upload_result['media_id'] == 'media-1'
assert [cmd for cmd, _ in client.frames] == [
'aibot_upload_media_init',
'aibot_upload_media_chunk',
'aibot_upload_media_chunk',
'aibot_upload_media_finish',
]
init_body = client.frames[0][1]
assert init_body['type'] == 'image'
assert init_body['filename'] == 'image.png'
assert init_body['total_size'] == len(data)
assert init_body['total_chunks'] == 2
assert client.frames[1][1]['chunk_index'] == 0
assert base64.b64decode(client.frames[1][1]['base64_data']) == b'a' * _UPLOAD_CHUNK_SIZE
assert client.frames[2][1]['chunk_index'] == 1
assert base64.b64decode(client.frames[2][1]['base64_data']) == b'a'
@pytest.mark.asyncio
async def test_reply_message_uploads_and_replies_image_media():
bot = Bot()
adapter = make_adapter(bot)
png_data = b'\x89PNG\r\n\x1a\nimage'
image_b64 = base64.b64encode(png_data).decode('utf-8')
chain = platform_message.MessageChain([platform_message.Image(base64=f'data:image/png;base64,{image_b64}')])
items = await WecomBotMessageConverter.yiri2target(chain)
await adapter._send_media(bot, 'req-1', items[0])
assert bot.calls == [
('upload_media', 'image', 'attachment.image', png_data),
('reply_image', 'req-1', 'media-1'),
]
@pytest.mark.asyncio
async def test_send_message_sends_text_and_skips_proactive_image():
bot = Bot()
adapter = make_adapter(bot)
jpg_data = b'\xff\xd8\xffimage'
image_b64 = base64.b64encode(jpg_data).decode('utf-8')
chain = platform_message.MessageChain(
[
platform_message.Plain(text='before'),
platform_message.Image(base64=f'data:image/jpeg;base64,{image_b64}'),
platform_message.Plain(text='after'),
]
)
await adapter.send_message('group', 'chat-1', chain)
assert bot.calls == [
('send_message', 'chat-1', 'beforeafter'),
]
@@ -1,3 +1,4 @@
import uuid
from types import SimpleNamespace
from unittest.mock import AsyncMock
@@ -49,7 +50,29 @@ async def test_send_message_sends_text_to_customer_service_user():
assert kwargs['open_kfid'] == 'kf-test'
assert kwargs['external_userid'] == 'external-user'
assert kwargs['content'] == 'hello'
assert kwargs['msgid'].startswith('langbot_')
assert len(kwargs['msgid'].encode()) <= 32
assert uuid.UUID(hex=kwargs['msgid']).hex == kwargs['msgid']
@pytest.mark.asyncio
async def test_send_message_sends_image_to_customer_service_user():
adapter = make_adapter()
adapter.bot_account_id = 'kf-test'
adapter.bot = SimpleNamespace(
get_media_id=AsyncMock(return_value='media-id'),
send_image_msg=AsyncMock(),
)
message = platform_message.MessageChain([platform_message.Image(base64='aW1hZ2U=')])
await adapter.send_message('person', 'uexternal-user', message)
adapter.bot.send_image_msg.assert_awaited_once()
kwargs = adapter.bot.send_image_msg.await_args.kwargs
assert kwargs['open_kfid'] == 'kf-test'
assert kwargs['external_userid'] == 'external-user'
assert kwargs['media_id'] == 'media-id'
assert len(kwargs['msgid'].encode()) <= 32
@pytest.mark.asyncio
@@ -0,0 +1,47 @@
from __future__ import annotations
import httpx
import pytest
from langbot.libs.wecom_customer_service_api.api import WecomCSClient
@pytest.mark.asyncio
async def test_send_image_msg_posts_customer_service_image_payload() -> None:
captured_request: httpx.Request | None = None
def handle_request(request: httpx.Request) -> httpx.Response:
nonlocal captured_request
captured_request = request
return httpx.Response(200, json={'errcode': 0})
client = WecomCSClient(
corpid='corp-id',
secret='secret',
token='token',
EncodingAESKey='encoding-key',
logger=None,
unified_mode=True,
)
client.access_token = 'access-token'
client._http_client = httpx.AsyncClient(transport=httpx.MockTransport(handle_request))
try:
await client.send_image_msg(
open_kfid='kf-test',
external_userid='external-user',
msgid='a' * 32,
media_id='media-id',
)
finally:
await client.close()
assert captured_request is not None
assert captured_request.url.path == '/cgi-bin/kf/send_msg'
assert captured_request.url.params['access_token'] == 'access-token'
assert captured_request.method == 'POST'
assert captured_request.read().decode() == (
'{"touser":"external-user","open_kfid":"kf-test","msgid":"'
+ 'a' * 32
+ '","msgtype":"image","image":{"media_id":"media-id"}}'
)
@@ -107,6 +107,19 @@ def shared_connector(
return connector
@pytest.mark.asyncio
async def test_shared_reconcile_uses_configured_cold_start_timeout():
binding = execution_binding("workspace-a")
setting = plugin_setting("01", "a" * 64)
connector = shared_connector([[binding]], {"workspace-a": [setting]})
connector.ap.instance_config.data["plugin"]["connect_timeout_seconds"] = 900
connector.handler = runtime_handler()
await connector._prepare_connected_runtime()
assert connector.handler.reconcile_plugin_installations.await_args.kwargs["timeout"] == 900
@pytest.mark.asyncio
async def test_shared_reconnect_replays_two_workspaces_and_removes_missing_projection():
binding_a = execution_binding('workspace-a')
@@ -150,7 +163,7 @@ async def test_empty_projected_workspaces_do_not_retain_installation_sets():
assert connector._workspace_installations == {}
assert connector._known_desired_states == {}
connector.handler.reconcile_plugin_installations.assert_awaited_once_with(())
connector.handler.reconcile_plugin_installations.assert_awaited_once_with((), timeout=300.0)
@pytest.mark.asyncio
+12
View File
@@ -81,6 +81,18 @@ async def test_reconcile_plugin_installations_allows_cloud_cold_start_to_finish(
assert runtime_handler.call_action.await_args.kwargs['timeout'] == 300
@pytest.mark.asyncio
async def test_reconcile_plugin_installations_accepts_configured_cold_start_timeout():
runtime_handler = make_handler(SimpleNamespace())
runtime_handler.call_action = AsyncMock(return_value={})
binding = next(iter(runtime_handler._installation_bindings.values()))[0]
desired = PluginInstallationDesiredState(binding=binding, enabled=True)
await runtime_handler.reconcile_plugin_installations((desired,), timeout=900)
assert runtime_handler.call_action.await_args.kwargs["timeout"] == 900
class TestHandlerQueryVariables:
"""Tests for handler query variable logic."""
+136 -6
View File
@@ -234,6 +234,7 @@ class TestSetBinaryStorage:
},
}
mock_app.persistence_mgr = Mock()
mock_app.persistence_mgr.get_db_engine.return_value = SimpleNamespace(dialect=SimpleNamespace(name='sqlite'))
mock_app.persistence_mgr.execute_async = AsyncMock(return_value=make_result())
mock_app.logger = Mock()
return mock_app
@@ -270,8 +271,8 @@ class TestSetBinaryStorage:
)
assert response.code == 0
assert app.persistence_mgr.execute_async.await_count == 2
insert_params = compiled_params(app.persistence_mgr.execute_async.await_args_list[1].args[0])
assert app.persistence_mgr.execute_async.await_count == 3
insert_params = compiled_params(app.persistence_mgr.execute_async.await_args_list[2].args[0])
assert insert_params['workspace_uuid'] == 'workspace-a'
assert insert_params['unique_key'] == canonical_binary_key(
'plugin',
@@ -301,6 +302,69 @@ class TestSetBinaryStorage:
assert expected_key in update_params.values()
assert update_params['value'] == b'new'
@pytest.mark.asyncio
async def test_adopts_legacy_storage_before_updating(self, app):
"""A migrated pre-tenancy row is updated in place rather than duplicated."""
runtime_handler = make_handler(app)
legacy_storage = SimpleNamespace(unique_key='plugin:test-author/test-plugin:test-key')
adopted = SimpleNamespace(rowcount=1)
app.persistence_mgr.execute_async.side_effect = [
make_result(),
make_result(legacy_storage),
adopted,
]
response = await runtime_handler.actions[RuntimeToLangBotAction.SET_BINARY_STORAGE.value](self.payload(b'new'))
assert response.code == 0
assert app.persistence_mgr.execute_async.await_count == 3
adoption_params = compiled_params(app.persistence_mgr.execute_async.await_args_list[2].args[0])
expected_key = canonical_binary_key('plugin', 'test-author/test-plugin', 'test-key')
assert expected_key in adoption_params.values()
assert adoption_params['value'] == b'new'
@pytest.mark.asyncio
async def test_legacy_adoption_race_updates_winning_canonical_row(self, app):
runtime_handler = make_handler(app)
legacy_storage = SimpleNamespace(unique_key='plugin:test-author/test-plugin:test-key')
lost_race = SimpleNamespace(rowcount=0)
canonical_winner = SimpleNamespace(rowcount=1)
app.persistence_mgr.execute_async.side_effect = [
make_result(),
make_result(legacy_storage),
lost_race,
canonical_winner,
]
response = await runtime_handler.actions[RuntimeToLangBotAction.SET_BINARY_STORAGE.value](self.payload(b'new'))
assert response.code == 0
assert app.persistence_mgr.execute_async.await_count == 4
winner_update = compiled_params(app.persistence_mgr.execute_async.await_args_list[3].args[0])
assert canonical_binary_key('plugin', 'test-author/test-plugin', 'test-key') in winner_update.values()
assert winner_update['value'] == b'new'
@pytest.mark.asyncio
async def test_legacy_adoption_lost_to_delete_inserts_new_value(self, app):
runtime_handler = make_handler(app)
legacy_storage = SimpleNamespace(unique_key='plugin:test-author/test-plugin:test-key')
lost_race = SimpleNamespace(rowcount=0)
app.persistence_mgr.execute_async.side_effect = [
make_result(),
make_result(legacy_storage),
lost_race,
SimpleNamespace(rowcount=0),
make_result(),
]
response = await runtime_handler.actions[RuntimeToLangBotAction.SET_BINARY_STORAGE.value](self.payload(b'new'))
assert response.code == 0
assert app.persistence_mgr.execute_async.await_count == 5
insert_params = compiled_params(app.persistence_mgr.execute_async.await_args_list[4].args[0])
assert insert_params['unique_key'] == canonical_binary_key('plugin', 'test-author/test-plugin', 'test-key')
assert insert_params['value'] == b'new'
@pytest.mark.asyncio
async def test_invalid_max_value_bytes_falls_back_to_default_limit(self, app):
"""Invalid max_value_bytes uses the 10MB default limit."""
@@ -525,6 +589,46 @@ class TestGetBinaryStorage:
in statement_params.values()
)
@pytest.mark.asyncio
async def test_reads_legacy_storage_without_mutating_key(self, app):
runtime_handler = make_handler(app)
legacy_storage = SimpleNamespace(
unique_key='plugin:test-author/test-plugin:test-key',
value=b'legacy bytes',
)
app.persistence_mgr.execute_async.side_effect = [
make_result(),
make_result(legacy_storage),
]
response = await runtime_handler.actions[RuntimeToLangBotAction.GET_BINARY_STORAGE.value](
{'key': 'test-key', 'owner_type': 'plugin', 'owner': 'ignored'}
)
assert response.code == 0
assert base64.b64decode(response.data['value_base64']) == b'legacy bytes'
assert app.persistence_mgr.execute_async.await_count == 2
@pytest.mark.asyncio
async def test_retries_canonical_after_concurrent_legacy_adoption(self, app):
runtime_handler = make_handler(app)
canonical_storage = SimpleNamespace(value=b'adopted bytes')
app.persistence_mgr.execute_async.side_effect = [
make_result(),
make_result(),
make_result(canonical_storage),
]
response = await runtime_handler.actions[RuntimeToLangBotAction.GET_BINARY_STORAGE.value](
{'key': 'test-key', 'owner_type': 'plugin', 'owner': 'ignored'}
)
assert response.code == 0
assert base64.b64decode(response.data['value_base64']) == b'adopted bytes'
assert app.persistence_mgr.execute_async.await_count == 3
retry_params = compiled_params(app.persistence_mgr.execute_async.await_args_list[2].args[0])
assert canonical_binary_key('plugin', 'test-author/test-plugin', 'test-key') in retry_params.values()
@pytest.mark.asyncio
async def test_returns_error_when_not_found(self, app):
"""Missing binary storage rows return an error response."""
@@ -567,21 +671,47 @@ class TestDeleteAndListBinaryStorage:
assert response.code == 0
statement_params = compiled_params(app.persistence_mgr.execute_async.await_args.args[0])
assert 'workspace-a' in statement_params.values()
flat_values = [
item for value in statement_params.values() for item in (value if isinstance(value, list) else [value])
]
assert 'workspace-a' in flat_values
assert (
canonical_binary_key(
'plugin',
'test-author/test-plugin',
'test-key',
)
in statement_params.values()
in flat_values
)
assert 'forged-owner' not in statement_params.values()
assert 'forged-owner' not in flat_values
@pytest.mark.asyncio
async def test_delete_removes_canonical_and_legacy_scoped_keys(self, app):
runtime_handler = make_handler(app)
response = await runtime_handler.actions[RuntimeToLangBotAction.DELETE_BINARY_STORAGE.value](
{
'key': 'test-key',
'owner_type': 'plugin',
'owner': 'forged-owner',
}
)
assert response.code == 0
statement_params = compiled_params(app.persistence_mgr.execute_async.await_args.args[0])
values = [
item for value in statement_params.values() for item in (value if isinstance(value, list) else [value])
]
assert 'workspace-a' in values
assert canonical_binary_key('plugin', 'test-author/test-plugin', 'test-key') in values
assert 'plugin:test-author/test-plugin:test-key' in values
assert 'test-author/test-plugin' in values
assert 'forged-owner' not in values
@pytest.mark.asyncio
async def test_list_keys_uses_trusted_plugin_owner(self, app):
result = Mock()
result.scalars.return_value.all.return_value = ['first', 'second']
result.scalars.return_value.all.return_value = ['first', 'second', 'first']
app.persistence_mgr.execute_async.return_value = result
runtime_handler = make_handler(app)
@@ -91,3 +91,42 @@ def test_convert_messages_plain_string_content_untouched():
msg = provider_message.Message(role='user', content='just text')
out = req._convert_messages([msg])
assert out[0]['content'] == 'just text'
def test_convert_messages_replayed_image_without_base64_does_not_crash():
"""Replayed image parts hollowed out by history trimming must not raise KeyError (#2469).
SessionManager clears image_base64 on past turns, and URL-less platform
images never had a URL, so the replayed part serializes as
{'type': 'image_base64'} with no payload keys. The hollow part should be
dropped while the sibling text part survives.
"""
req = _make_requester()
image = provider_message.ContentElement.from_image_base64('data:image/jpeg;base64,AAAA')
# Simulate SessionManager.trim_conversation_messages clearing binary payloads.
image.image_base64 = None
msg = provider_message.Message(
role='user',
content=[
provider_message.ContentElement.from_text('describe the photo'),
image,
],
)
out = req._convert_messages([msg])
assert [p.get('type') for p in out[0]['content']] == ['text']
def test_convert_messages_replayed_image_with_url_falls_back_to_url():
"""When base64 was trimmed but image_url survived, rebuild the OpenAI image_url part from the URL."""
req = _make_requester()
image = provider_message.ContentElement(
type='image_base64',
image_base64=None,
image_url=provider_message.ImageURLContentObject(url='https://example.com/pic.jpg'),
)
msg = provider_message.Message(role='user', content=[image])
out = req._convert_messages([msg])
parts = out[0]['content']
assert [p.get('type') for p in parts] == ['image_url']
assert parts[0]['image_url'] == {'url': 'https://example.com/pic.jpg'}
assert 'image_base64' not in parts[0]
@@ -0,0 +1,208 @@
"""Regression tests for tool-message content serialization (#2457).
MCP tools return ``list[ContentElement]`` from ``execute_func_call``.
The runner must serialize that list to a string before placing it in a
``role='tool'`` message, because the OpenAI chat-completions spec
requires tool-message content to be a string. Sending the raw list
causes OpenAI-compatible endpoints to return HTTP 500.
"""
from __future__ import annotations
import json
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock
import pytest
import langbot_plugin.api.entities.builtin.pipeline.query as pipeline_query
import langbot_plugin.api.entities.builtin.provider.message as provider_message
import langbot_plugin.api.entities.builtin.provider.session as provider_session
from langbot.pkg.api.http.context import ExecutionContext, PrincipalContext, PrincipalType
from langbot.pkg.provider.runners.localagent import LocalAgentRunner
class _ToolCallProvider:
"""Non-streaming provider: round 1 issues a tool call, round 2 returns text."""
def __init__(self):
self.requests: list[dict] = []
async def invoke_llm(self, query, model, messages, funcs, extra_args=None, remove_think=None):
self.requests.append({'messages': list(messages)})
if len(self.requests) == 1:
return provider_message.Message(
role='assistant',
content='Let me search that.',
tool_calls=[
provider_message.ToolCall(
id='call-mcp-1',
type='function',
function=provider_message.FunctionCall(
name='duckduckgo_search',
arguments=json.dumps({'query': 'swift'}),
),
)
],
)
return provider_message.Message(role='assistant', content='Done.')
class _ToolCallStreamProvider:
"""Streaming variant of _ToolCallProvider."""
def __init__(self):
self.requests: list[dict] = []
def invoke_llm_stream(self, query, model, messages, funcs, extra_args=None, remove_think=None):
self.requests.append({'messages': list(messages)})
async def _stream():
if len(self.requests) == 1:
yield provider_message.MessageChunk(
role='assistant',
content='Let me search that.',
tool_calls=[
provider_message.ToolCall(
id='call-mcp-1',
type='function',
function=provider_message.FunctionCall(
name='duckduckgo_search',
arguments=json.dumps({'query': 'swift'}),
),
)
],
is_final=True,
)
return
yield provider_message.MessageChunk(
role='assistant',
content='Done.',
is_final=True,
)
return _stream()
def _make_query(stream: bool = False) -> pipeline_query.Query:
adapter = AsyncMock()
adapter.is_stream_output_supported = AsyncMock(return_value=stream)
query = pipeline_query.Query.model_construct(
query_id='mcp-tool-query',
launcher_type=provider_session.LauncherTypes.PERSON,
launcher_id=12345,
sender_id=12345,
message_chain=[],
message_event=None,
adapter=adapter,
pipeline_uuid='pipeline-uuid',
bot_uuid='bot-uuid',
pipeline_config={
'ai': {
'runner': {'runner': 'local-agent'},
'local-agent': {'model': {'primary': 'test-model-uuid', 'fallbacks': []}, 'prompt': 'test-prompt'},
},
'output': {'misc': {'remove-think': False}},
},
prompt=SimpleNamespace(messages=[]),
messages=[],
user_message=provider_message.Message(role='user', content='search swift'),
use_funcs=[SimpleNamespace(name='duckduckgo_search')],
use_llm_model_uuid='test-model-uuid',
variables={},
)
object.__setattr__(
query,
'_execution_context',
ExecutionContext(
instance_uuid='instance-test',
workspace_uuid='workspace-test',
placement_generation=1,
trigger_principal=PrincipalContext(PrincipalType.SYSTEM),
),
)
return query
def _make_app(provider, func_ret) -> SimpleNamespace:
"""Build a minimal app whose tool_mgr returns *func_ret*."""
model = SimpleNamespace(
provider=provider,
model_entity=SimpleNamespace(
uuid='test-model-uuid',
name='test-model',
abilities=['func_call'],
extra_args={},
),
)
return SimpleNamespace(
logger=Mock(),
model_mgr=SimpleNamespace(get_model_by_uuid=AsyncMock(return_value=model)),
tool_mgr=SimpleNamespace(execute_func_call=AsyncMock(return_value=func_ret)),
rag_mgr=SimpleNamespace(),
box_service=SimpleNamespace(get_system_guidance=Mock(return_value='sandbox guidance')),
skill_mgr=SimpleNamespace(
get_skills_for_pipeline=AsyncMock(return_value=[]),
detect_skill_activation=AsyncMock(return_value=None),
build_activation_prompt=Mock(return_value=None),
),
)
# The actual shape returned by MCP tools: a list of ContentElement objects.
_MCP_FUNC_RET = [
provider_message.ContentElement.from_text('Title: Swift - Wikipedia\nURL: https://en.wikipedia.org/wiki/Swift'),
provider_message.ContentElement.from_text('Title: Swift Programming Language\nURL: https://swift.org'),
]
@pytest.mark.asyncio
async def test_tool_message_content_is_string_not_list():
"""Non-streaming: tool message content must be a string (#2457).
Before the fix, ``func_ret`` (a ``list[ContentElement]``) was assigned
to ``tool_content`` as-is, so the tool message carried a list instead
of a string, causing OpenAI-compatible APIs to return 500.
"""
provider = _ToolCallProvider()
app = _make_app(provider, _MCP_FUNC_RET)
runner = LocalAgentRunner(app, pipeline_config={})
query = _make_query(stream=False)
results = [msg async for msg in runner.run(query)]
tool_msgs = [m for m in results if m.role == 'tool']
assert len(tool_msgs) == 1
# The content must be a string, not a list.
assert isinstance(tool_msgs[0].content, str), (
f'tool message content should be str, got {type(tool_msgs[0].content).__name__}'
)
# And it should contain the text of both ContentElements.
assert 'Swift - Wikipedia' in tool_msgs[0].content
assert 'Swift Programming Language' in tool_msgs[0].content
@pytest.mark.asyncio
async def test_tool_message_content_is_string_in_stream():
"""Streaming: same regression check for the streaming path (#2457)."""
provider = _ToolCallStreamProvider()
app = _make_app(provider, _MCP_FUNC_RET)
runner = LocalAgentRunner(app, pipeline_config={})
query = _make_query(stream=True)
results = [msg async for msg in runner.run(query)]
tool_msgs = [m for m in results if m.role == 'tool']
assert len(tool_msgs) == 1
assert isinstance(tool_msgs[0].content, str), (
f'tool message content should be str, got {type(tool_msgs[0].content).__name__}'
)
assert 'Swift - Wikipedia' in tool_msgs[0].content
assert 'Swift Programming Language' in tool_msgs[0].content
@@ -0,0 +1,15 @@
from __future__ import annotations
import tomllib
from pathlib import Path
def test_seekdb_is_only_declared_as_an_optional_dependency() -> None:
project_root = Path(__file__).resolve().parents[2]
with (project_root / 'pyproject.toml').open('rb') as pyproject_file:
pyproject = tomllib.load(pyproject_file)
project = pyproject['project']
base_dependencies = project['dependencies']
assert not any(dependency.lower().startswith('pyseekdb') for dependency in base_dependencies)
assert project['optional-dependencies']['seekdb'] == ['pyseekdb==1.1.0.post3']
+39
View File
@@ -53,6 +53,45 @@ async def test_matches_any_rejects_pattern_and_input_amplification():
)
@pytest.mark.asyncio
async def test_mask_patterns_honors_explicit_pattern_count_cap():
patterns = ['a'] * (safe_regex.MAX_PATTERN_COUNT + 6)
found, masked = await safe_regex.mask_patterns(
patterns,
'hello',
mask='*',
mask_word='',
max_pattern_count=len(patterns),
)
assert found is False
assert masked == 'hello'
with pytest.raises(safe_regex.SafeRegexLimitError):
await safe_regex.mask_patterns(
patterns,
'hello',
mask='*',
mask_word='',
)
@pytest.mark.asyncio
async def test_mask_patterns_rejects_oversized_sequence_before_copying_it():
class OversizedPatterns(list):
def __iter__(self):
raise AssertionError('oversized patterns must not be materialized')
patterns = OversizedPatterns(['a'] * (safe_regex.MAX_PATTERN_COUNT + 1))
with pytest.raises(safe_regex.SafeRegexLimitError):
await safe_regex.mask_patterns(
patterns,
'hello',
mask='*',
mask_word='',
)
@pytest.mark.asyncio
async def test_mask_patterns_bounds_replacement_growth_and_masks_matches():
found, masked = await safe_regex.mask_patterns(
@@ -0,0 +1,34 @@
from __future__ import annotations
import importlib
from unittest.mock import MagicMock
import pytest
from tests.utils.import_isolation import isolated_sys_modules
_INSTALL_HINT = "Install LangBot with the 'seekdb' extra"
def test_seekdb_vector_backend_reports_missing_optional_extra() -> None:
module_name = 'langbot.pkg.vector.vdbs.seekdb'
with isolated_sys_modules({'pyseekdb': None}, clear=[module_name]):
seekdb_module = importlib.import_module(module_name)
assert seekdb_module.SEEKDB_AVAILABLE is False
with pytest.raises(ImportError, match=_INSTALL_HINT):
seekdb_module.SeekDBVectorDatabase(MagicMock())
@pytest.mark.asyncio
async def test_seekdb_embedding_reports_missing_optional_extra() -> None:
module_name = 'langbot.pkg.provider.modelmgr.requesters.seekdbembed'
with isolated_sys_modules({'pyseekdb': None}, clear=[module_name]):
seekdb_embedding_module = importlib.import_module(module_name)
requester = seekdb_embedding_module.SeekDBEmbedding.__new__(seekdb_embedding_module.SeekDBEmbedding)
with pytest.raises(ImportError, match=_INSTALL_HINT):
await requester.initialize()
@@ -88,14 +88,15 @@ async def test_environment_mapping_enables_provider_without_leaking_secret(monke
assert service.capability() == {'enabled': True, 'provider': 'smtp'}
async def test_cloud_invitation_email_has_branded_html_plain_fallback_and_expiry_copy():
async def test_invitation_email_has_generic_langbot_brand_plain_fallback_and_expiry_copy():
service = InvitationDeliveryService(_app({}))
link = 'https://cloud.langbot.app/invitations/accept#token=lbi_secret&next=<unsafe>'
text = service._plain_text('Research & Development', link)
html = service._html('Research & Development', link)
assert 'LangBot Cloud' in text
assert 'LangBot' in text
assert 'LangBot Cloud' not in text
assert 'Research & Development' in text
assert '7 days' in text
assert link in text
@@ -103,3 +104,55 @@ async def test_cloud_invitation_email_has_branded_html_plain_fallback_and_expiry
assert 'Research &amp; Development' in html
assert 'expires in 7 days' in html
assert 'lbi_secret&amp;next=&lt;unsafe&gt;' in html
assert 'LangBot Cloud' not in html
async def test_invitation_email_uses_quiet_brand_lockup_and_compact_fallback_link():
service = InvitationDeliveryService(_app({}))
link = 'https://cloud.langbot.app/invitations/accept#token=lbi_secret'
html = service._html("RockChinQ's Workspace", link)
assert 'https://docs.langbot.app/langbot-logo.png' in html
assert '>LangBot<' in html
assert 'Workspace invitation' in html
assert 'Open invitation link' in html
assert 'linear-gradient' not in html
assert 'box-shadow' not in html
assert 'border-top:4px solid' not in html
assert 'border:1px solid #dfe6f0' not in html
assert 'height="28"' in html
assert 'height="32"' in html
assert 'margin-top:32px' not in html
assert f'>{link}<' not in html
async def test_oss_smtp_configuration_delivers_the_generic_invitation_email():
service = InvitationDeliveryService(
_app(
{
'workspace': {
'invitations': {
'email': {
'provider': 'smtp',
'from': 'LangBot <noreply@example.com>',
'smtp': {'host': 'smtp.example.com'},
}
}
}
}
)
)
service._send_smtp = AsyncMock(return_value=True)
link = 'https://self-hosted.example/invitations/accept#token=lbi_secret'
result = await service.deliver_invitation(
recipient_email='member@example.com',
workspace_name='Self-hosted Workspace',
invitation_link=link,
)
assert result == InvitationDeliveryResult(status='sent', provider='smtp')
service._send_smtp.assert_awaited_once()
assert 'LangBot Cloud' not in service._plain_text('Self-hosted Workspace', link)
assert 'LangBot Cloud' not in service._html('Self-hosted Workspace', link)
Generated
+10 -5
View File
@@ -9,10 +9,10 @@ resolution-markers = [
"python_full_version == '3.13.*' and sys_platform == 'emscripten'",
"python_full_version == '3.13.*' and sys_platform != 'emscripten' and sys_platform != 'win32'",
"python_full_version == '3.12.*' and sys_platform == 'win32'",
"python_full_version < '3.12' and sys_platform == 'win32'",
"python_full_version == '3.12.*' and sys_platform == 'emscripten'",
"python_full_version < '3.12' and sys_platform == 'emscripten'",
"python_full_version == '3.12.*' and sys_platform != 'emscripten' and sys_platform != 'win32'",
"python_full_version < '3.12' and sys_platform == 'win32'",
"python_full_version < '3.12' and sys_platform == 'emscripten'",
"python_full_version < '3.12' and sys_platform != 'emscripten' and sys_platform != 'win32'",
]
@@ -2008,7 +2008,7 @@ wheels = [
[[package]]
name = "langbot"
version = "4.10.7"
version = "4.10.8"
source = { editable = "." }
dependencies = [
{ name = "aiocqhttp" },
@@ -2063,7 +2063,6 @@ dependencies = [
{ name = "pymilvus" },
{ name = "pynacl" },
{ name = "pypdf2" },
{ name = "pyseekdb" },
{ name = "python-docx" },
{ name = "python-multipart" },
{ name = "python-socks" },
@@ -2089,6 +2088,11 @@ dependencies = [
{ name = "websockets" },
]
[package.optional-dependencies]
seekdb = [
{ name = "pyseekdb" },
]
[package.dev-dependencies]
dev = [
{ name = "moto" },
@@ -2153,7 +2157,7 @@ requires-dist = [
{ name = "pymilvus", specifier = ">=2.6.4" },
{ name = "pynacl", specifier = ">=1.5.0" },
{ name = "pypdf2", specifier = ">=3.0.1" },
{ name = "pyseekdb", specifier = "==1.1.0.post3" },
{ name = "pyseekdb", marker = "extra == 'seekdb'", specifier = "==1.1.0.post3" },
{ name = "python-docx", specifier = ">=1.1.0" },
{ name = "python-multipart", specifier = ">=0.0.27" },
{ name = "python-socks", specifier = ">=2.7.1" },
@@ -2178,6 +2182,7 @@ requires-dist = [
{ name = "valkey-glide", marker = "sys_platform != 'win32'", specifier = ">=2.4.1,<3.0.0" },
{ name = "websockets", specifier = ">=15.0.1" },
]
provides-extras = ["seekdb"]
[package.metadata.requires-dev]
dev = [
+25 -2
View File
@@ -710,11 +710,32 @@ export class BackendClient extends BaseHttpClient {
);
}
private async getAuthenticatedObjectURL(path: string): Promise<string> {
private async getAuthenticatedObjectURL(
path: string,
rewritePluginPageSdk = false,
): Promise<string> {
const response = await this.instance.get<Blob>(path, {
responseType: 'blob',
});
return URL.createObjectURL(response.data);
let blob = response.data;
if (rewritePluginPageSdk && blob.type.startsWith('text/html')) {
const apiBase =
this.instance.defaults.baseURL === '/'
? window.location.origin
: this.instance.defaults.baseURL?.replace(/\/$/, '');
const pageSdkUrl = `${apiBase}/api/v1/plugins/_sdk/page-sdk.js`;
const html = await blob.text();
blob = new Blob(
[
html.replace(
/(<script\b[^>]*\bsrc\s*=\s*)(["'])\/api\/v1\/plugins\/_sdk\/page-sdk\.js\2/gi,
`$1$2${pageSdkUrl}$2`,
),
],
{ type: blob.type },
);
}
return URL.createObjectURL(blob);
}
public getAuthenticatedPluginAssetURL(
@@ -724,6 +745,7 @@ export class BackendClient extends BaseHttpClient {
): Promise<string> {
return this.getAuthenticatedObjectURL(
`/api/v1/plugins/${author}/${name}/authenticated-assets/${filepath}`,
true,
);
}
@@ -1181,6 +1203,7 @@ export class BackendClient extends BaseHttpClient {
public getAccountInfo(): Promise<{
initialized: boolean;
authenticated_invitation_acceptance_enabled?: boolean;
invitation_registration_enabled?: boolean;
password_login_enabled?: boolean;
space_login_enabled?: boolean;
}> {
+28 -22
View File
@@ -91,7 +91,9 @@ export default function AcceptInvitationPage() {
const [errorMessage, setErrorMessage] = useState('');
const [password, setPassword] = useState('');
const [confirmPassword, setConfirmPassword] = useState('');
const [passwordRegistrationEnabled, setPasswordRegistrationEnabled] =
const [invitationRegistrationEnabled, setInvitationRegistrationEnabled] =
useState(false);
const [invitationCapabilitiesLoaded, setInvitationCapabilitiesLoaded] =
useState(false);
const [
authenticatedInvitationAcceptanceEnabled,
@@ -116,12 +118,16 @@ export default function AcceptInvitationPage() {
backendClient
.getAccountInfo()
.then((info) => {
setPasswordRegistrationEnabled(info.password_login_enabled !== false);
setInvitationRegistrationEnabled(
info.invitation_registration_enabled ??
info.password_login_enabled !== false,
);
setAuthenticatedInvitationAcceptanceEnabled(
info.authenticated_invitation_acceptance_enabled === true,
);
})
.catch(() => setPasswordRegistrationEnabled(false));
.catch(() => setInvitationRegistrationEnabled(false))
.finally(() => setInvitationCapabilitiesLoaded(true));
if (!invitationToken) {
setErrorMessage(t('workspace.invitationMissing'));
setStatus('error');
@@ -311,7 +317,11 @@ export default function AcceptInvitationPage() {
</div>
)}
{hasLoginToken && authenticatedInvitationAcceptanceEnabled ? (
{!invitationCapabilitiesLoaded ? (
<div className="flex justify-center py-8">
<Loader2 className="size-6 animate-spin" />
</div>
) : hasLoginToken && authenticatedInvitationAcceptanceEnabled ? (
<Button
className="w-full"
disabled={status === 'submitting'}
@@ -331,7 +341,7 @@ export default function AcceptInvitationPage() {
{t('workspace.logoutAndReturn')}
</Button>
</div>
) : passwordRegistrationEnabled ? (
) : invitationRegistrationEnabled ? (
<>
<div className="space-y-2">
<label
@@ -376,15 +386,19 @@ export default function AcceptInvitationPage() {
>
{t('workspace.confirmPassword')}
</label>
<Input
id="invite-password-confirm"
type="password"
value={confirmPassword}
onChange={(event) =>
setConfirmPassword(event.target.value)
}
autoComplete="new-password"
/>
<div className="relative">
<Lock className="absolute left-3 top-3 size-4 text-muted-foreground" />
<Input
id="invite-password-confirm"
type="password"
value={confirmPassword}
onChange={(event) =>
setConfirmPassword(event.target.value)
}
className="pl-10"
autoComplete="new-password"
/>
</div>
</div>
<Button
className="w-full"
@@ -396,14 +410,6 @@ export default function AcceptInvitationPage() {
)}
{t('workspace.registerAndAccept')}
</Button>
<Button
variant="ghost"
className="w-full"
disabled={status === 'submitting'}
onClick={() => navigate('/login?invitation=1')}
>
{t('workspace.alreadyHaveAccount')}
</Button>
</>
) : (
<Button
+2
View File
@@ -503,6 +503,8 @@ async function handleBackendApi(route: Route, state: LangBotApiMockState) {
if (path === '/api/v1/user/account-info') {
return fulfillJson(route, {
initialized: true,
authenticated_invitation_acceptance_enabled: false,
invitation_registration_enabled: true,
password_login_enabled: true,
space_login_enabled: false,
});
+107 -1
View File
@@ -93,7 +93,7 @@ test('login preserves an explicit invitation email mismatch error', async ({
});
await page.goto('/invitations/accept#token=mismatch-invitation');
await page.getByRole('button', { name: 'I already have an account' }).click();
await page.goto('/login?invitation=1');
await page.getByPlaceholder('Enter email address').fill('other@example.com');
await page.getByPlaceholder('Enter password').fill('password');
await page.getByRole('button', { name: 'Login with password' }).click();
@@ -107,6 +107,111 @@ test('login preserves an explicit invitation email mismatch error', async ({
await expect(page.getByText('Login successful')).toHaveCount(0);
});
test('an OAuth-only OSS instance registers the invited email with a local password', async ({
page,
}) => {
let registration: { email?: string; password?: string } | undefined;
await installLangBotApiMocks(page, { authenticated: false });
await page.route('**/api/v1/user/account-info', async (route) => {
await new Promise((resolve) => setTimeout(resolve, 800));
await route.fulfill({
status: 200,
contentType: 'application/json',
body: JSON.stringify({
code: 0,
data: {
initialized: true,
authenticated_invitation_acceptance_enabled: false,
invitation_registration_enabled: true,
password_login_enabled: false,
space_login_enabled: true,
},
msg: 'ok',
}),
});
});
await page.route('**/api/v1/invitations/inspect', async (route) => {
await route.fulfill({
status: 200,
contentType: 'application/json',
body: JSON.stringify({
code: 0,
data: {
invitation: {
uuid: 'oss-local-registration',
workspace_uuid: 'workspace-playwright',
normalized_email: 'invited@example.com',
role: 'viewer',
status: 'pending',
},
workspace: {
uuid: 'workspace-playwright',
name: 'Playwright Workspace',
},
},
msg: 'ok',
}),
});
});
await page.route('**/api/v1/invitations/accept', async (route) => {
const body = JSON.parse(route.request().postData() || '{}') as {
registration?: { email?: string; password?: string };
};
registration = body.registration;
await route.fulfill({
status: 200,
contentType: 'application/json',
body: JSON.stringify({
code: 0,
data: {
login_required: true,
workspace_uuid: 'workspace-playwright',
},
msg: 'ok',
}),
});
});
await page.goto('/invitations/accept#token=oss-local-registration');
await expect(page.getByText('Playwright Workspace')).toBeVisible();
await expect(
page.getByRole('button', { name: 'Login with LangBot Account' }),
).toHaveCount(0);
await expect(page.locator('#invite-email')).toHaveValue(
'invited@example.com',
);
await expect(page.locator('#invite-email')).toHaveAttribute('readonly', '');
await expect(page.locator('#invite-password')).toBeVisible();
await expect(page.locator('#invite-password-confirm')).toBeVisible();
for (const inputId of ['invite-password', 'invite-password-confirm']) {
const input = page.locator(`#${inputId}`);
const field = input.locator('xpath=..');
await expect(field).toHaveClass(/relative/);
await expect(field.locator('svg')).toBeVisible();
await expect(input).toHaveClass(/pl-10/);
}
await expect(
page.getByRole('button', { name: 'Create account and accept' }),
).toBeVisible();
await expect(
page.getByRole('button', { name: 'I already have an account' }),
).toHaveCount(0);
await expect(
page.getByRole('button', { name: 'Login with LangBot Account' }),
).toHaveCount(0);
await page.locator('#invite-password').fill('invite-password-123');
await page.locator('#invite-password-confirm').fill('invite-password-123');
await page.getByRole('button', { name: 'Create account and accept' }).click();
await expect(page).toHaveURL(/\/login\?invitation=1$/);
expect(registration).toEqual({
email: 'invited@example.com',
password: 'invite-password-123',
});
});
test('an authenticated OSS invitation requires logout before registration', async ({
page,
}) => {
@@ -185,6 +290,7 @@ test('an authenticated Cloud Account can accept its invitation directly', async
data: {
initialized: true,
authenticated_invitation_acceptance_enabled: true,
invitation_registration_enabled: false,
password_login_enabled: false,
space_login_enabled: true,
},
+63
View File
@@ -0,0 +1,63 @@
import { expect, test } from '@playwright/test';
import { installLangBotApiMocks } from './fixtures/langbot-api';
test('an OSS local-only owner is prompted to bind before using LangBot Models', async ({
page,
}) => {
await installLangBotApiMocks(page, { authenticated: true });
await page.route('**/api/v1/user/space-credits', (route) =>
route.fulfill({
contentType: 'application/json',
body: JSON.stringify({
code: 0,
data: {
credits: null,
owner_space_bound: false,
is_workspace_owner: true,
},
msg: 'ok',
}),
}),
);
await page.route('**/api/v1/provider/providers', (route) =>
route.fulfill({
contentType: 'application/json',
body: JSON.stringify({
code: 0,
data: {
providers: [
{
uuid: 'langbot-models-provider',
name: 'LangBot Models',
requester: 'space-chat-completions',
base_url: '',
api_keys: [],
},
],
},
msg: 'ok',
}),
}),
);
await page.route('**/api/v1/provider/requesters**', (route) =>
route.fulfill({
contentType: 'application/json',
body: JSON.stringify({ code: 0, data: { requesters: [] }, msg: 'ok' }),
}),
);
await page.route('**/api/v1/provider/models/**', (route) =>
route.fulfill({
contentType: 'application/json',
body: JSON.stringify({ code: 0, data: { models: [] }, msg: 'ok' }),
}),
);
await page.goto('/home?action=showModelSettings');
await expect(
page.getByRole('button', {
name: 'The Workspace owner must connect a LangBot Account for LangBot Models.',
}),
).toBeVisible();
});
+58 -4
View File
@@ -69,7 +69,55 @@ test('loads a Cloud plugin page through the authenticated asset route', async ({
await route.fulfill({
status: 200,
contentType: 'text/html',
body: '<!doctype html><html><body><h1>LangRAG Observability</h1></body></html>',
body: `<!doctype html>
<html>
<body>
<h1>LangRAG Observability</h1>
<button id="save">Save</button>
<script src="/api/v1/plugins/_sdk/page-sdk.js"></script>
<script>
document.querySelector('#save').addEventListener('click', async () => {
await window.langbot.api('/settings', { enabled: true }, 'POST');
document.body.dataset.saved = 'true';
});
</script>
</body>
</html>`,
});
},
);
let pageSdkRequests = 0;
await page.route('**/api/v1/plugins/_sdk/page-sdk.js', async (route) => {
pageSdkRequests += 1;
await route.fulfill({
status: 200,
contentType: 'application/javascript',
body: `window.langbot = {
api(endpoint, body, method) {
return new Promise((resolve) => {
const requestId = 'request-' + Date.now();
const handler = (event) => {
if (event.data?.type === 'langbot:api:response' && event.data.requestId === requestId) {
window.removeEventListener('message', handler);
resolve(event.data.data);
}
};
window.addEventListener('message', handler);
window.parent.postMessage({ type: 'langbot:api', requestId, endpoint, body, method }, '*');
});
},
};`,
});
});
let pageApiRequests = 0;
await page.route(
'**/api/v1/plugins/langbot-team/LangRAG/page-api',
async (route) => {
pageApiRequests += 1;
await route.fulfill({
status: 200,
contentType: 'application/json',
body: wrapped({ saved: true }),
});
},
);
@@ -78,11 +126,17 @@ test('loads a Cloud plugin page through the authenticated asset route', async ({
'/home/plugin-pages?id=langbot-team%2FLangRAG%2Fobservability',
);
const pluginFrame = page.frameLocator('iframe');
await expect(
page
.frameLocator('iframe')
.getByRole('heading', { name: 'LangRAG Observability' }),
pluginFrame.getByRole('heading', { name: 'LangRAG Observability' }),
).toBeVisible();
await pluginFrame.getByRole('button', { name: 'Save' }).click();
await expect(pluginFrame.locator('body')).toHaveAttribute(
'data-saved',
'true',
);
expect(authenticatedAssetRequests).toBeGreaterThan(0);
expect(pageSdkRequests).toBe(1);
expect(pageApiRequests).toBe(1);
await expect(page.getByText('Loading...')).toHaveCount(0);
});