Compare commits

..

17 Commits

Author SHA1 Message Date
dadachann 473ba573a3 fix(cloud): retain launch replay records through clock skew 2026-07-30 21:37:06 +00:00
dadachann 93dbd3541e fix(cloud): make direct launch replay-safe 2026-07-30 21:02:38 +00:00
dadachann a5a26f81ee feat(auth): accept direct Space launch assertions 2026-07-30 20:05:17 +00:00
dadachann 92d9db8f95 chore(plugin): pin SDK 0.5.0 2026-07-30 20:05:17 +00:00
dadachann 59db012594 feat(cloud): enforce workspace resource quotas 2026-07-30 18:48:42 +00:00
dadachann 88f328066b fix(cloud): keep runtime sdk ahead of plugin dependencies 2026-07-30 16:42:01 +00:00
dadachann d155d9d5a8 fix(cloud): allow explicitly disabled box runtime 2026-07-30 15:42:49 +00:00
dadachann dd95545309 fix(prod): configure shared Box runtime 2026-07-30 15:17:17 +00:00
dadachann 9066c25729 fix(deploy): let migration own exact runtime ACLs 2026-07-30 14:58:31 +00:00
dadachann 6d2e9d3d72 fix(deploy): do not start disabled Box runtime 2026-07-30 14:48:51 +00:00
dadachann a0b85e11fd fix(deploy): keep runtime role schema read-only 2026-07-30 14:44:57 +00:00
dadachann 122d8fa659 fix(deploy): retry transient image pull failures 2026-07-30 14:34:20 +00:00
dadachann ace8cc67f2 fix(config): preserve typed list environment overrides 2026-07-30 14:26:41 +00:00
dadachann 52c0772806 fix(cloud): pin Space production URL and deployment health 2026-07-30 14:21:04 +00:00
dadachann d5044c2f1e fix(ci): authenticate cloud adapter checkout 2026-07-30 14:08:58 +00:00
dadachann 7baa89254c ops(cloud): deploy exact production stack to jp09 2026-07-30 13:58:37 +00:00
RockChinQ e1ac5e0fc8 feat(tenancy): add Workspace multi-tenant foundation (#2353)
* Document multi-tenant workspace architecture

* Add OSS and commercial workspace boundaries

* docs: redesign multi-tenant workspace architecture

* feat(tenancy): implement workspace isolation

* docs(tenancy): record verification evidence

* docs(tenancy): revise single-instance SaaS topology

* docs(tenancy): refine architecture options

* docs: finalize cloud v2 multi-tenant decisions

* feat(tenancy): establish cloud isolation foundations

* feat(tenancy): harden shared cloud runtime boundaries

* docs(tenancy): record final isolation verification

* fix(tenancy): close isolation and permission gaps

* docs(tenancy): record final isolation verification

* feat(tenancy): connect cloud workspace control plane

* fix(build): install git for pinned SDK

* docs(cloud): update control plane verification

* chore: update multi-tenant SDK pin

* fix(cloud): skip legacy model sync during startup

* test(cloud): preserve minimal model manager fixtures

* fix(cloud): preserve authenticated account context

* fix(cloud): reuse authenticated account for user info

* feat(cloud): complete Workspace settings navigation

* test(web): cover Workspace dropdown menu

* feat(web): place workspace controls in sidebar

* refactor(web): streamline workspace controls

* style(web): format workspace layout test

* fix(cloud): surface runtime and workspace plan status

* fix(plugin): keep runtime identity stable across restarts

* fix(ui): widen and center workspace switcher

* fix(ui): hide roles from workspace switcher

* fix(ui): align workspace switcher with sidebar entries

* feat(workspace): add in-product collaboration and direct Cloud launch

* style: format collaboration changes

* fix(workspace): bind collaboration APIs to tenant UoW

* fix(cloud): preserve Core-owned collaboration state

* test(cloud): require Space identity for invite registration

* feat(cloud): complete secure invitation experience

* style(web): format invitation flows

* fix(cloud): recover box runtime without unscoped skill reload

* feat(oss): enforce invitation account and owner billing flows

* style: format OSS account service

* test(oss): cover invitation logout handoff

* fix(oss): resolve workspace owner in scoped session

* feat(cloud): harden multi-tenant runtime resources

* fix(cloud): bound runtime restart storms

* fix(cloud): eliminate periodic runtime CPU spikes

* fix(cloud): enforce instance capacity ceilings

* fix(cloud): scope public login capability discovery

* fix(cloud): bound tenant maintenance and monitoring work

* fix(runtime): bound tenant resource amplification

* fix(deps): pin green multi-tenant plugin SDK

* fix(cloud): handle unavailable skill capability

* fix(security): require authentication for image file endpoint (H-2)

- Changed /api/v1/files/image from AuthType.NONE to USER_TOKEN_OR_API_KEY
- Added Permission.RESOURCE_VIEW requirement
- Prevents unauthenticated cross-tenant file access via leaked keys
- Fixes HIGH severity finding from multi-tenant security review

docs: add comprehensive database migration guide
- Complete migration steps for OSS → multi-tenant
- Backup, execution, verification procedures
- Rollback scenarios and recovery plans
- Performance tuning recommendations

* test: add comprehensive cross-tenant isolation tests

Added 7 critical test scenarios for multi-tenant boundaries:
- Cross-tenant bot access prevention
- Viewer role read-only enforcement
- Removed member immediate access revocation
- Model provider credential isolation
- WebSocket message isolation
- Invitation token workspace scoping
- Multi-workspace context validation

These tests address P0-2 coverage gaps for:
- workspaces.py (membership & invitation flows)
- user.py (authentication & authorization)
- websocket_chat.py (real-time isolation)
- plugins.py (resource access control)

docs: finalize database migration guide

* fix(security): resolve M-1, M-2, M-3 security findings

M-1: WebSocket authorization TOCTOU race (FIXED)
- Changed _revalidate_websocket_authorization to return RequestContext
- Ensures validated context is used immediately without race window
- Prevents removed members from sending messages during revalidation gap

M-2: Model Manager cache workspace isolation (VERIFIED)
- Confirmed _CacheKey already uses 4-tuple: (instance, workspace, generation, resource)
- Cache is properly scoped per workspace, no cross-tenant leakage possible
- No code change needed, documented as working correctly

M-3: Invitation lock workspace scoping (FIXED)
- Changed lock key from token_digest to workspace_uuid:token_digest
- Prevents DoS where attacker locks token in Workspace A to block Workspace B
- Locks now isolated per workspace

All MEDIUM severity findings from security review now resolved.

* fix(cloud): unblock tenant CI and enforce knowledge quotas

* fix(tenancy): scope rerank model sync

---------

Co-authored-by: dadachann <185672915+dadachann@users.noreply.github.com>
2026-07-30 21:43:35 +08:00
39 changed files with 1545 additions and 56 deletions
+78
View File
@@ -0,0 +1,78 @@
name: Build and deploy production
on:
push:
branches: [deploy/prod]
workflow_dispatch:
permissions:
contents: read
concurrency:
group: langbot-production
cancel-in-progress: false
env:
CORE_IMAGE: ${{ secrets.DOCKER_USERNAME }}/langbot
CLOUD_IMAGE: ${{ secrets.DOCKER_USERNAME }}/langbot-cloud-core
SPACE_REF: e1b261dac45e886efc667b1096a4ec493c6a6111
jobs:
build-and-deploy:
runs-on: ubuntu-latest
environment: production
steps:
- uses: actions/checkout@v4
- uses: docker/setup-buildx-action@v3
- uses: docker/login-action@v3
with:
username: ${{ secrets.DOCKER_USERNAME }}
password: ${{ secrets.DOCKER_PASSWORD }}
- name: Build exact Core image
uses: docker/build-push-action@v6
with:
context: .
push: true
tags: |
${{ env.CORE_IMAGE }}:prod-${{ github.sha }}
${{ env.CORE_IMAGE }}:deploy-prod
cache-from: type=gha,scope=core-prod
cache-to: type=gha,mode=max,scope=core-prod
- name: Checkout production Cloud adapter
uses: actions/checkout@v4
with:
repository: langbot-app/langbot-space
ref: ${{ env.SPACE_REF }}
token: ${{ secrets.CLA_PAT }}
path: .space
- name: Build exact Cloud Core image
uses: docker/build-push-action@v6
with:
context: .space
file: .space/Dockerfile.cloud
push: true
build-args: LANGBOT_CORE_IMAGE=${{ env.CORE_IMAGE }}:prod-${{ github.sha }}
tags: |
${{ env.CLOUD_IMAGE }}:prod-${{ github.sha }}
${{ env.CLOUD_IMAGE }}:deploy-prod
cache-from: type=gha,scope=cloud-core-prod
cache-to: type=gha,mode=max,scope=cloud-core-prod
- name: Configure SSH
env:
SSH_KEY: ${{ secrets.JP09_SSH_KEY }}
KNOWN_HOSTS: ${{ secrets.JP09_KNOWN_HOSTS }}
run: |
install -m 700 -d ~/.ssh
install -m 600 /dev/null ~/.ssh/id_ed25519
printf '%s\n' "$SSH_KEY" > ~/.ssh/id_ed25519
printf '%s\n' "$KNOWN_HOSTS" > ~/.ssh/known_hosts
- name: Upload release manifest and deploy
env:
HOST: ${{ secrets.JP09_HOST }}
USER: ${{ secrets.JP09_USER }}
PORT: ${{ secrets.JP09_PORT }}
run: |
remote="$USER@$HOST"
ssh -p "$PORT" "$remote" 'install -d -m 700 /opt/langbot-cloud-prod'
scp -P "$PORT" deploy/prod/docker-compose.yml deploy/prod/deploy.sh "$remote:/opt/langbot-cloud-prod/"
ssh -p "$PORT" "$remote" "chmod 700 /opt/langbot-cloud-prod/deploy.sh && /opt/langbot-cloud-prod/deploy.sh prod-${GITHUB_SHA}"
+92
View File
@@ -0,0 +1,92 @@
#!/usr/bin/env bash
set -Eeuo pipefail
cd /opt/langbot-cloud-prod
TAG=${1:?usage: deploy.sh prod-<40-char-sha>}
[[ "$TAG" =~ ^prod-[0-9a-f]{40}$ ]] || { echo 'invalid immutable image tag' >&2; exit 2; }
[[ -s .env ]] || { echo '/opt/langbot-cloud-prod/.env is missing' >&2; exit 3; }
rendered_compose=$(docker compose config)
grep -Fq 'LANGBOT_SPACE_CONTROL_PLANE_URL: https://space.langbot.app' <<<"$rendered_compose" || {
echo 'Cloud control-plane URL must be https://space.langbot.app' >&2
exit 4
}
grep -Fq 'SPACE__URL: https://space.langbot.app' <<<"$rendered_compose" || {
echo 'Cloud user-facing Space URL must be https://space.langbot.app' >&2
exit 5
}
update_env() {
local key=$1 value=$2
python3 - "$key" "$value" <<'PY'
from pathlib import Path
import os
import sys
path = Path('.env')
key, value = sys.argv[1:]
lines = path.read_text().splitlines()
updated = False
for index, line in enumerate(lines):
if line.startswith(f'{key}='):
lines[index] = f'{key}={value}'
updated = True
break
if not updated:
lines.append(f'{key}={value}')
temporary = Path('.env.tmp')
temporary.write_text('\n'.join(lines) + '\n')
os.chmod(temporary, 0o600)
temporary.replace(path)
PY
}
update_env LANGBOT_IMAGE_TAG "$TAG"
set -a
. ./.env
set +a
for attempt in 1 2 3 4 5; do
if docker compose pull postgres redis migrate plugin-runtime core; then
break
fi
if [ "$attempt" -eq 5 ]; then
echo "docker compose pull failed after $attempt attempts" >&2
exit 1
fi
delay=$((attempt * 10))
echo "docker compose pull failed (attempt $attempt/5); retrying in ${delay}s" >&2
sleep "$delay"
done
docker compose up -d postgres redis
for _ in $(seq 1 60); do
if docker compose exec -T postgres pg_isready -U langbot_operator -d langbot >/dev/null 2>&1; then break; fi
sleep 2
done
docker compose exec -T postgres pg_isready -U langbot_operator -d langbot >/dev/null
docker compose exec -T postgres psql -v ON_ERROR_STOP=1 -U langbot_operator -d langbot \
-v runtime_password="$POSTGRES_RUNTIME_PASSWORD" <<'SQL'
SELECT format('CREATE ROLE langbot_runtime LOGIN PASSWORD %L', :'runtime_password')
WHERE NOT EXISTS (SELECT 1 FROM pg_roles WHERE rolname = 'langbot_runtime')\gexec
ALTER ROLE langbot_runtime PASSWORD :'runtime_password';
GRANT CONNECT ON DATABASE langbot TO langbot_runtime;
REVOKE CREATE ON SCHEMA public FROM PUBLIC, langbot_runtime;
REVOKE ALL PRIVILEGES ON ALL TABLES IN SCHEMA public FROM langbot_runtime;
REVOKE ALL PRIVILEGES ON ALL SEQUENCES IN SCHEMA public FROM langbot_runtime;
ALTER DEFAULT PRIVILEGES FOR ROLE langbot_operator IN SCHEMA public REVOKE ALL ON TABLES FROM langbot_runtime;
ALTER DEFAULT PRIVILEGES FOR ROLE langbot_operator IN SCHEMA public REVOKE ALL ON SEQUENCES FROM langbot_runtime;
GRANT USAGE ON SCHEMA public TO langbot_runtime;
SQL
docker compose --profile tools run --rm migrate
docker compose up -d --remove-orphans plugin-runtime core
for _ in $(seq 1 90); do
if docker compose exec -T core python -c 'import urllib.request; urllib.request.urlopen("http://127.0.0.1:5300/healthz", timeout=3)' >/dev/null 2>&1; then
docker compose ps
exit 0
fi
sleep 2
done
docker compose logs --tail=200 core plugin-runtime >&2
exit 1
+161
View File
@@ -0,0 +1,161 @@
services:
postgres:
image: pgvector/pgvector:pg17
container_name: langbot-cloud-postgres
restart: unless-stopped
environment:
POSTGRES_DB: langbot
POSTGRES_USER: langbot_operator
POSTGRES_PASSWORD: ${POSTGRES_OPERATOR_PASSWORD}
volumes:
- postgres-data:/var/lib/postgresql/data
healthcheck:
test: [CMD-SHELL, "pg_isready -U langbot_operator -d langbot"]
interval: 5s
timeout: 5s
retries: 30
networks: [internal]
redis:
image: redis:7.4-alpine
container_name: langbot-cloud-redis
restart: unless-stopped
command: [redis-server, --appendonly, "yes", --requirepass, "${REDIS_PASSWORD}"]
volumes:
- redis-data:/data
healthcheck:
test: [CMD-SHELL, "redis-cli -a \"$${REDIS_PASSWORD}\" ping | grep PONG"]
interval: 5s
timeout: 5s
retries: 20
environment:
REDIS_PASSWORD: ${REDIS_PASSWORD}
networks: [internal]
migrate:
image: rockchin/langbot-cloud-core:${LANGBOT_IMAGE_TAG}
profiles: [tools]
command: [uv, run, langbot, migrate, --cloud]
environment: &core-env
TZ: Asia/Shanghai
SYSTEM__INSTANCE_ID: ${CLOUD_V2_INSTANCE_UUID}
SYSTEM__EDITION: cloud
SYSTEM__RECOVERY_KEY: ${SYSTEM_RECOVERY_KEY}
SYSTEM__JWT__SECRET: ${JWT_SECRET}
SYSTEM__LIMITATION__MAX_BOTS: "2"
SYSTEM__LIMITATION__MAX_PIPELINES: "3"
SYSTEM__LIMITATION__MAX_EXTENSIONS: "3"
SYSTEM__LIMITATION__MAX_KNOWLEDGE_BASES: "2"
API__WEBHOOK_PREFIX: https://cloud.langbot.app
API__WEBUI_URL: https://cloud.langbot.app
WORKSPACE__INVITATIONS__PUBLIC_WEB_URL: https://cloud.langbot.app
DATABASE__USE: postgresql
DATABASE__POSTGRESQL__URL: postgresql+asyncpg://langbot_runtime:${POSTGRES_RUNTIME_PASSWORD}@postgres:5432/langbot
DATABASE__CLOUD_MIGRATION__OPERATOR_DSN_ENV: LANGBOT_CLOUD_MIGRATION_DSN
LANGBOT_CLOUD_MIGRATION_DSN: postgresql://langbot_operator:${POSTGRES_OPERATOR_PASSWORD}@postgres:5432/langbot
VDB__USE: pgvector
VDB__PGVECTOR__USE_BUSINESS_DATABASE: "true"
VDB__PGVECTOR__ALLOWED_DIMENSIONS: "384,512,768,1024,1536"
PLUGIN__ENABLE: "true"
PLUGIN__RUNTIME_WS_URL: ws://plugin-runtime:5400/control/ws
PLUGIN__DISPLAY_PLUGIN_DEBUG_URL: wss://cloud.langbot.app/plugin/debug/ws
PLUGIN__WORKER__MAX_CPUS: "0.25"
PLUGIN__WORKER__MAX_MEMORY_MB: "256"
PLUGIN__WORKER__MAX_PIDS: "128"
PLUGIN__WORKER__MAX_WORKERS: "16"
PLUGIN__WORKER__MAX_TOTAL_CPUS: "4.0"
PLUGIN__WORKER__MAX_TOTAL_MEMORY_MB: "4096"
PLUGIN__WORKER__REQUIRE_HARD_LIMITS: "true"
LANGBOT_PLUGIN_RUNTIME_CONTROL_TOKEN: ${PLUGIN_RUNTIME_CONTROL_TOKEN}
# Cloud v2 currently grants no managed Box capability. Keep the shared
# runtime deployed but disable Core integration until a hard-quota-capable
# backend can satisfy the fail-closed Cloud readiness contract.
BOX__ENABLED: "false"
BOX__BACKEND: nsjail
BOX__RUNTIME__ENDPOINT: ws://box:5410
BOX__ADMISSION__REQUIRED: "true"
BOX__ADMISSION__LOGICAL_SESSION_ID: global
BOX__ADMISSION__REQUIRED_BACKEND: nsjail
BOX__ADMISSION__MAX_SESSIONS: "1"
BOX__ADMISSION__MAX_MANAGED_PROCESSES: "0"
BOX__ADMISSION__CPUS: "0.25"
BOX__ADMISSION__MEMORY_MB: "256"
BOX__ADMISSION__WORKSPACE_QUOTA_MB: "256"
BOX__LOCAL__HOST_ROOT: /app/data/box
BOX__LOCAL__DEFAULT_WORKSPACE: /app/data/box
BOX__LOCAL__ALLOWED_MOUNT_ROOTS: /app/data/box
LANGBOT_BOX_CONTROL_TOKEN: ${BOX_CONTROL_TOKEN}
MCP__STDIO__ENABLED: "false"
LANGBOT_SPACE_CONTROL_PLANE_URL: https://space.langbot.app
LANGBOT_SPACE_CONTROL_PLANE_TOKEN: ${CLOUD_V2_CONTROL_PLANE_TOKEN}
LANGBOT_SPACE_CONTROL_PLANE_PUBLIC_KEY: ${CLOUD_V2_MANIFEST_PUBLIC_KEY}
LANGBOT_SPACE_CONTROL_PLANE_KEY_ID: ${CLOUD_V2_MANIFEST_KEY_ID}
SPACE__URL: https://space.langbot.app
depends_on:
postgres: {condition: service_healthy}
networks: [internal]
plugin-runtime:
image: rockchin/langbot:${LANGBOT_IMAGE_TAG}
container_name: langbot-cloud-plugin-runtime
restart: unless-stopped
command: [uv, run, python, -m, langbot_plugin.cli.__init__, rt]
environment:
LANGBOT_PLUGIN_RUNTIME_CONTROL_TOKEN: ${PLUGIN_RUNTIME_CONTROL_TOKEN}
volumes:
- plugin-data:/app/data
- /sys/fs/cgroup:/sys/fs/cgroup:rw
cgroup: host
privileged: true
expose: ["5400"]
networks: [internal]
box:
image: rockchin/langbot:${LANGBOT_IMAGE_TAG}
container_name: langbot-cloud-box
restart: unless-stopped
command: [uv, run, lbp, box, --host, 0.0.0.0, --ws-control-port, "5410"]
environment:
LANGBOT_BOX_CONTROL_TOKEN: ${BOX_CONTROL_TOKEN}
LANGBOT_BOX_ROOT: /app/data/box
volumes:
- box-data:/app/data/box
- /sys/fs/cgroup:/sys/fs/cgroup:rw
cgroup: host
privileged: true
expose: ["5410"]
networks: [internal]
core:
image: rockchin/langbot-cloud-core:${LANGBOT_IMAGE_TAG}
container_name: langbot-cloud-core
restart: unless-stopped
environment: *core-env
volumes:
- core-data:/app/data
- box-data:/app/data/box
depends_on:
postgres: {condition: service_healthy}
redis: {condition: service_healthy}
plugin-runtime: {condition: service_started}
box: {condition: service_started}
expose: ["5300"]
healthcheck:
test: [CMD-SHELL, "python -c 'import urllib.request; urllib.request.urlopen(\"http://127.0.0.1:5300/healthz\", timeout=3)'" ]
interval: 10s
timeout: 5s
retries: 30
start_period: 30s
networks: [internal, shared-network]
networks:
internal:
shared-network:
external: true
volumes:
postgres-data:
redis-data:
plugin-data:
box-data:
core-data:
@@ -0,0 +1,166 @@
# LangBot 多租户数据库迁移指南
## 概述
LangBot 从单租户 OSS 架构迁移到多租户 SaaS 架构,需要执行 7 个数据库迁移(0009-0015)。
## 迁移序列
```
0009_workspace_tenancy_kernel → 创建 Workspace、成员、邀请表
0010_scope_tenant_resources → 所有业务表添加 workspace_uuid
0011_postgres_tenant_rls → PostgreSQL 行级安全策略
0012_plugin_installation_identity → 插件实例租户绑定
0013_tenant_pgvector → RAG 向量存储隔离
0014_cloud_directory_projection → Cloud 控制平面同步
0015_cloud_core_collaboration → 协作和权限功能
```
## 执行步骤
### 1. 备份(必须)
```bash
# SQLite
cp ~/.langbot/data/langbot.db ~/.langbot/data/langbot.db.backup-$(date +%Y%m%d)
# PostgreSQL
pg_dump -U langbot_user -d langbot_db -F c -f langbot_backup_$(date +%Y%m%d).dump
```
### 2. 执行迁移
```bash
# 停止服务
sudo -S -p '' systemctl stop langbot
# 执行迁移
python -m langbot.pkg.persistence.migration upgrade head
# 验证
python -m langbot.pkg.persistence.migration current
# 预期: 0015_cloud_core_collaboration
# 启动服务
sudo -S -p '' systemctl start langbot
```
### 3. OSS 单租户自动迁移
迁移会自动:
- 创建默认 Workspace(名称:"Default Workspace"
- 第一个用户成为 Owner
- 所有现有资源绑定到该 Workspace
### 4. 验证检查
```bash
# 检查 Workspace
python << EOF
from langbot.pkg.persistence import manager
import sqlalchemy as sa
with manager.engine.connect() as conn:
ws = conn.execute(sa.text("SELECT uuid, name FROM workspaces LIMIT 1")).first()
print(f"Workspace: {ws[1]} ({ws[0]})")
# 检查资源绑定
bot_count = conn.execute(sa.text(
f"SELECT COUNT(*) FROM bots WHERE workspace_uuid='{ws[0]}'"
)).scalar()
print(f"Bots: {bot_count}")
EOF
```
## 回滚方案
### 完全回滚(丢失多租户数据)
```bash
# 1. 停止服务
sudo -S -p '' systemctl stop langbot
# 2. 恢复备份
cp ~/.langbot/data/langbot.db.backup-YYYYMMDD ~/.langbot/data/langbot.db
# 3. 回退代码
git checkout v4.10.x
pip install -e .
# 4. 启动
sudo -S -p '' systemctl start langbot
```
### 降级迁移(保留数据但移除多租户)
```bash
# 警告:会移除 Workspace 表但保留资源
python -m langbot.pkg.persistence.migration downgrade 0008_mcp_resource_prefs
```
## 常见问题
### Q: 迁移后无法登录
```bash
# 检查用户 UUID
python << EOF
from langbot.pkg.persistence import manager
import sqlalchemy as sa
with manager.engine.connect() as conn:
users = conn.execute(sa.text("SELECT id, user, uuid, status FROM users")).all()
for u in users:
print(f"{u[1]}: UUID={u[2]}, Status={u[3]}")
EOF
```
### Q: 资源看不见了
检查 Workspace 上下文:
```bash
# 前端请求需要带 X-Workspace-ID header
curl -H "Authorization: Bearer $TOKEN" \
-H "X-Workspace-ID: $WORKSPACE_UUID" \
http://localhost:5200/api/v1/platform/bots
```
### Q: 迁移速度慢
```bash
# SQLite 优化
sqlite3 ~/.langbot/data/langbot.db << EOF
PRAGMA journal_mode=WAL;
PRAGMA synchronous=NORMAL;
VACUUM;
EOF
```
## 性能调优
### PostgreSQL 索引
```sql
-- 迁移后创建
CREATE INDEX CONCURRENTLY idx_model_providers_workspace
ON model_providers(workspace_uuid);
CREATE INDEX CONCURRENTLY idx_bots_workspace
ON bots(workspace_uuid);
CREATE INDEX CONCURRENTLY idx_pipelines_workspace
ON pipelines(workspace_uuid);
```
## 预估时间
- SQLite < 100MB: 2-5 分钟
- SQLite 100MB-1GB: 5-15 分钟
- PostgreSQL: < 5 分钟(取决于数据量)
## 支持
问题反馈:https://github.com/langbot-app/LangBot/issues
**版本**: 1.0
**最后更新**: 2026-07-30
@@ -81,10 +81,10 @@ This log records implementation choices made while delivering the Workspace arch
- Decision: The MCP ASGI mount authenticates the API key once, binds an immutable per-request `RequestContext`, and every tool checks a fixed permission before calling tenant services with that same context.
- Reason: Authenticating the transport without propagating Workspace identity into tool calls would leave the direct service path globally scoped.
### Unreleased SDK protocol is pinned reproducibly without publishing
### Released SDK protocol is pinned from PyPI
- Decision: The SDK tenancy protocol is versioned as 0.4.18. This task does not create a GitHub release or publish PyPI because the user authorized pushing code, not a package release. After the SDK feature branch is final, LangBot's feature branch temporarily pins the exact pushed SDK Git commit. Before merging to master, the release gate is to publish `langbot-plugin==0.4.18` and replace the Git pin with the registry pin.
- Reason: The current registry release does not contain the complete tenant action context and shared Runtime hardening. An exact Git commit is reproducible and keeps the feature branch testable without expanding release authority.
- Decision: The SDK tenancy protocol is released as `langbot-plugin==0.5.0` and LangBot pins that exact registry version.
- Reason: The final PyPI release contains the complete tenant action context and shared Runtime hardening, while the exact version pin keeps production installs reproducible.
### Cloud directory writes stay outside Core
+1 -1
View File
@@ -71,7 +71,7 @@ dependencies = [
"chromadb>=1.0.0,<2.0.0",
"qdrant-client (>=1.15.1,<2.0.0)",
"pyseekdb==1.1.0.post3",
"langbot-plugin @ git+https://github.com/langbot-app/langbot-plugin-sdk.git@1d65ed301a6afc52150a998043f73cd6032c8162",
"langbot-plugin==0.5.0",
"asyncpg>=0.30.0",
"line-bot-sdk>=3.19.0",
"matrix-nio>=0.25.2",
@@ -14,6 +14,7 @@ from ....utils import bounded_executor
from ....workspace.collaboration import MembershipPermissionError, WorkspaceCollaborationError
from ....workspace.errors import WorkspaceNotFoundError
from ....cloud.entitlements import EntitlementUnavailableError
from ....cloud.quotas import WorkspaceQuotaExceededError
from ....core.errors import TaskCapacityError
from ..authz import AuthorizationError, Permission, permissions_for_role, require_permission
from ..context import PrincipalContext, PrincipalType, RequestContext, WorkspaceContext
@@ -219,6 +220,8 @@ class RouterGroup(abc.ABC):
return self.http_status(403, e.code, str(e))
if isinstance(e, WorkspaceCollaborationError):
return self.http_status(400, e.code, str(e))
if isinstance(e, WorkspaceQuotaExceededError):
return self.http_status(409, e.error_code, str(e))
if isinstance(e, TaskCapacityError):
return self.http_status(429, 'task_capacity_exceeded', str(e))
if isinstance(
@@ -23,8 +23,13 @@ def _storage_owner(context: RequestContext) -> str:
@group.group_class('files', '/api/v1/files')
class FilesRouterGroup(group.RouterGroup):
async def initialize(self) -> None:
@self.route('/image/<path:image_key>', methods=['GET'], auth_type=group.AuthType.NONE)
async def _(image_key: str) -> quart.Response:
@self.route(
'/image/<path:image_key>',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def _(image_key: str, request_context: RequestContext) -> quart.Response:
image_bytes = await self.ap.storage_mgr.resolve_public_object(
image_key,
expected_owner_type='upload_image',
@@ -128,7 +128,7 @@ class WebSocketChatRouterGroup(group.RouterGroup):
self,
request_context: RequestContext,
token: str,
) -> None:
) -> RequestContext:
"""Recheck revocable account, membership, permission, and placement state."""
account, _ = await self._authenticate_account(token)
@@ -168,6 +168,7 @@ class WebSocketChatRouterGroup(group.RouterGroup):
entitlement_revision=request_context.entitlement_revision,
)
require_permission(current_context, Permission.RUNTIME_OPERATE)
return current_context
async def _get_scoped_adapter(self, request_context: RequestContext, pipeline_uuid: str):
pipeline = await run_in_workspace_uow(
@@ -410,6 +410,7 @@ class UserRouterGroup(group.RouterGroup):
'token': token,
'user': account.user,
'workspace_uuid': access.workspace.uuid,
'return_path': launch.get('return_path', '/home'),
}
)
except SpaceLaunchError:
+28 -8
View File
@@ -4,6 +4,7 @@ import uuid
import sqlalchemy
from ....core import app
from ....cloud.quotas import require_resource_capacity, resolve_workspace_quota
from ....entity.persistence import bot as persistence_bot
from ....entity.persistence import pipeline as persistence_pipeline
from ....workspace.errors import WorkspaceNotFoundError
@@ -101,20 +102,21 @@ class BotService:
async def create_bot(self, context: TenantContext, bot_data: dict) -> str:
"""Create bot"""
workspace_uuid = require_workspace_uuid(context)
# Check limitation
limitation = self.ap.instance_config.data.get('system', {}).get('limitation', {})
max_bots = limitation.get('max_bots', -1)
if max_bots >= 0:
existing_bots = await self.get_bots(context)
if len(existing_bots) >= max_bots:
raise ValueError(f'Maximum number of bots ({max_bots}) reached')
quota = await resolve_workspace_quota(
self.ap,
workspace_uuid,
'bots.max',
fallback=limitation.get('max_bots', -1),
)
# TODO: 检查配置信息格式
bot_data = bot_data.copy()
bot_data['uuid'] = str(uuid.uuid4())
bot_data['workspace_uuid'] = workspace_uuid
# bind the most recently updated pipeline if any exist
# Preserve the legacy flat-row result shape for this optional lookup;
# quota admission and insertion below still share one transaction.
result = await self.ap.persistence_mgr.execute_async(
scope_statement(
sqlalchemy.select(persistence_pipeline.LegacyPipeline),
@@ -129,7 +131,25 @@ class BotService:
bot_data['use_pipeline_uuid'] = pipeline.uuid
bot_data['use_pipeline_name'] = pipeline.name
await self.ap.persistence_mgr.execute_async(sqlalchemy.insert(persistence_bot.Bot).values(bot_data))
async def persist(execute) -> None:
await require_resource_capacity(
execute,
workspace_uuid=workspace_uuid,
model=persistence_bot.Bot,
quota=quota,
resource_name='bots',
)
await execute(sqlalchemy.insert(persistence_bot.Bot).values(bot_data))
tenant_uow = getattr(self.ap.persistence_mgr, 'tenant_uow', None)
if quota.requires_transaction_lock:
if not callable(tenant_uow):
raise RuntimeError('Cloud bot quota enforcement requires transactional persistence')
async with tenant_uow(workspace_uuid) as uow:
await persist(uow.execute)
else:
await persist(self.ap.persistence_mgr.execute_async)
bot = await self.get_bot(context, bot_data['uuid'], include_secret=True)
@@ -56,6 +56,14 @@ class KnowledgeService:
require_workspace_uuid(context)
# In new architecture, we delegate entirely to RAGManager which uses plugins.
# Legacy internal KB creation is removed.
limitation = (
getattr(getattr(self.ap, 'instance_config', None), 'data', {}).get('system', {}).get('limitation', {})
)
max_knowledge_bases = limitation.get('max_knowledge_bases', -1)
if max_knowledge_bases >= 0:
knowledge_bases = await self.ap.rag_mgr.get_all_knowledge_base_details(context)
if len(knowledge_bases) >= max_knowledge_bases:
raise ValueError(f'Maximum number of knowledge bases ({max_knowledge_bases}) reached')
knowledge_engine_plugin_id = kb_data.get('knowledge_engine_plugin_id')
if not knowledge_engine_plugin_id:
+8 -2
View File
@@ -138,8 +138,14 @@ class VerifiedCloudDeployment:
if plugin_worker.get('require_hard_limits') is not True:
raise CloudBootstrapError('Cloud Runtime requires plugin.worker.require_hard_limits=true')
box_config = config.get('box', {})
if box_config.get('enabled') is not True:
raise CloudBootstrapError('Cloud runtime requires box.enabled=true')
box_enabled = box_config.get('enabled')
if box_enabled is False:
# Explicitly disabling Box removes the sandbox surface entirely and
# therefore does not weaken tenant isolation. Validate the strict
# runtime/admission contract only when the surface is enabled.
return
if box_enabled is not True:
raise CloudBootstrapError('Cloud runtime requires box.enabled to be an explicit boolean')
if box_config.get('backend') != 'nsjail':
raise CloudBootstrapError('Cloud runtime requires box.backend=nsjail')
runtime_endpoint = str(box_config.get('runtime', {}).get('endpoint', '') or '').strip()
+45 -10
View File
@@ -3,9 +3,11 @@ from __future__ import annotations
import asyncio
import base64
import binascii
import datetime
import hashlib
import heapq
import json
import math
import os
import time
import typing
@@ -14,6 +16,10 @@ from collections.abc import Callable, Iterable
from cryptography.exceptions import InvalidSignature
from cryptography.hazmat.primitives import serialization
from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey
import sqlalchemy
from sqlalchemy.dialects.postgresql import insert as pg_insert
from ..entity.persistence.cloud_directory import SpaceLaunchAssertionConsumption
if typing.TYPE_CHECKING:
from ..core.app import Application
@@ -119,21 +125,30 @@ class SpaceLaunchService:
*,
expected_workspace_uuid: str | None = None,
) -> dict[str, str]:
claims = self._verify_assertion(assertion)
claims, clock_skew_seconds = self._verify_assertion(assertion)
payload = claims.get('payload')
if not isinstance(payload, dict):
raise SpaceLaunchError('Launch assertion payload must be a JSON object')
account_uuid = _required_string(payload, 'account_uuid')
workspace_uuid = _required_string(payload, 'workspace_uuid')
return_path = _required_string(payload, 'return_path')
if (
not return_path.startswith('/')
or return_path.startswith('//')
or any(character in return_path for character in ('\\', '\r', '\n', '\t'))
):
raise SpaceLaunchError('Launch assertion return path is invalid')
if expected_workspace_uuid is not None and workspace_uuid != expected_workspace_uuid:
raise SpaceLaunchError('Launch assertion targets another Workspace')
await self._consume_jti(_required_string(claims, 'jti'), _required_int(claims, 'exp', minimum=1))
replay_retention_expires_at = _required_int(claims, 'exp', minimum=1) + math.ceil(clock_skew_seconds)
await self._consume_jti(_required_string(claims, 'jti'), replay_retention_expires_at)
return {
'account_uuid': account_uuid,
'workspace_uuid': workspace_uuid,
'return_path': return_path,
}
def _verify_assertion(self, token: str) -> dict[str, typing.Any]:
def _verify_assertion(self, token: str) -> tuple[dict[str, typing.Any], float]:
if not getattr(getattr(self.ap, 'deployment', None), 'multi_workspace_enabled', False):
raise SpaceLaunchError('Space direct launch requires verified Cloud mode')
public_key, key_id, clock_skew_seconds = self._trust_config()
@@ -184,7 +199,7 @@ class SpaceLaunchService:
raise SpaceLaunchError('Launch assertion is expired')
if expires_at <= max(issued_at, not_before):
raise SpaceLaunchError('Launch assertion expiry must follow issue time')
return claims
return claims, clock_skew_seconds
def _trust_config(self) -> tuple[Ed25519PublicKey, str, float]:
data = getattr(getattr(self.ap, 'instance_config', None), 'data', {}) or {}
@@ -214,19 +229,39 @@ class SpaceLaunchService:
async def _consume_jti(self, jti: str, expires_at: int) -> None:
digest = hashlib.sha256(jti.encode('utf-8')).hexdigest()
now = int(self._wall_time())
persistence_mgr = getattr(self.ap, 'persistence_mgr', None)
instance_uuid = str(self.ap.workspace_service.instance_uuid)
if persistence_mgr is not None:
expires_at_datetime = datetime.datetime.fromtimestamp(expires_at, tz=datetime.timezone.utc)
now_datetime = datetime.datetime.fromtimestamp(now, tz=datetime.timezone.utc)
async with persistence_mgr.directory_projection_uow(instance_uuid) as uow:
await uow.session.execute(
sqlalchemy.delete(SpaceLaunchAssertionConsumption).where(
SpaceLaunchAssertionConsumption.instance_uuid == instance_uuid,
SpaceLaunchAssertionConsumption.expires_at < now_datetime,
)
)
statement = (
pg_insert(SpaceLaunchAssertionConsumption)
.values(instance_uuid=instance_uuid, jti=digest, expires_at=expires_at_datetime)
.on_conflict_do_nothing(index_elements=['instance_uuid', 'jti'])
.returning(SpaceLaunchAssertionConsumption.jti)
)
result = await uow.session.execute(statement)
if result.scalar_one_or_none() is None:
raise SpaceLaunchError('Launch assertion has already been consumed')
return
# Lightweight unit-test and OSS compatibility fallback. Verified Cloud
# runtime always supplies the durable PostgreSQL persistence manager.
async with self._replay_lock:
self._prune_consumed_jtis(now)
if digest in self._consumed_jtis:
raise SpaceLaunchError('Launch assertion has already been consumed')
if len(self._consumed_jtis) >= _CONSUMED_JTI_MAX_ENTRIES:
# Evicting a still-valid digest would make a signed launch
# assertion replayable. Bound memory by failing closed instead.
raise SpaceLaunchError('Launch assertion replay cache capacity reached')
self._consumed_jtis[digest] = expires_at
heapq.heappush(
self._consumed_jti_expiry_heap,
(expires_at, digest),
)
heapq.heappush(self._consumed_jti_expiry_heap, (expires_at, digest))
def _prune_consumed_jtis(self, now: int) -> None:
while self._consumed_jti_expiry_heap:
+82
View File
@@ -0,0 +1,82 @@
from __future__ import annotations
from dataclasses import dataclass
from typing import Any, Awaitable, Callable
import sqlalchemy
from ..entity.persistence import workspace as persistence_workspace
from .entitlements import EntitlementResolver
Execute = Callable[[Any], Awaitable[Any]]
class WorkspaceQuotaExceededError(ValueError):
"""A stable business error raised when a workspace has no free slots."""
error_code = 'workspace_quota_exceeded'
def __init__(self, resource_name: str, limit: int) -> None:
self.resource_name = resource_name
self.limit = limit
super().__init__(f'Maximum number of {resource_name} ({limit}) reached')
@dataclass(frozen=True, slots=True)
class WorkspaceQuota:
limit: int
requires_transaction_lock: bool
async def resolve_workspace_quota(
ap: Any,
workspace_uuid: str,
limit_name: str,
*,
fallback: int = -1,
) -> WorkspaceQuota:
"""Resolve a plan-agnostic Cloud limit while preserving OSS configuration."""
resolver = getattr(ap, 'entitlement_resolver', None)
if isinstance(resolver, EntitlementResolver):
snapshot = await resolver.resolve(workspace_uuid)
return WorkspaceQuota(
limit=snapshot.limit(limit_name),
requires_transaction_lock=True,
)
return WorkspaceQuota(limit=fallback, requires_transaction_lock=False)
async def lock_workspace_for_quota(execute: Execute, workspace_uuid: str) -> None:
"""Serialize quota checks on the durable Workspace row within one transaction."""
result = await execute(
sqlalchemy.select(persistence_workspace.Workspace.uuid)
.where(persistence_workspace.Workspace.uuid == workspace_uuid)
.with_for_update()
)
if result.first() is None:
raise ValueError('Workspace does not exist')
async def require_resource_capacity(
execute: Execute,
*,
workspace_uuid: str,
model: type,
quota: WorkspaceQuota,
resource_name: str,
workspace_locked: bool = False,
) -> None:
if quota.limit < 0:
return
if quota.requires_transaction_lock and not workspace_locked:
await lock_workspace_for_quota(execute, workspace_uuid)
result = await execute(
sqlalchemy.select(sqlalchemy.func.count())
.select_from(model)
.where(model.workspace_uuid == workspace_uuid)
)
if int(result.scalar_one()) >= quota.limit:
raise WorkspaceQuotaExceededError(resource_name, quota.limit)
+8 -3
View File
@@ -186,9 +186,14 @@ def _apply_env_overrides_to_config(cfg: dict) -> dict:
# At the final key
if key in current:
if isinstance(current[key], list):
# Convert comma-separated string to list
# e.g., SYSTEM__DISABLED_ADAPTERS="aiocqhttp,dingtalk"
current[key] = [item.strip() for item in env_value.split(',') if item.strip()]
# Convert comma-separated values while preserving the
# element type declared by a non-empty config default.
items = [item.strip() for item in env_value.split(',') if item.strip()]
if current[key]:
exemplar = current[key][0]
current[key] = [convert_value(item, exemplar) for item in items]
else:
current[key] = items
elif isinstance(current[key], dict):
# Skip dict types
pass
@@ -67,3 +67,26 @@ class DirectoryProjectionInbox(Base):
name='ck_directory_projection_inbox_fingerprint',
),
)
class SpaceLaunchAssertionConsumption(Base):
"""Durable, instance-scoped replay ledger for signed Space launch assertions."""
__tablename__ = 'space_launch_assertion_consumptions'
instance_uuid = sqlalchemy.Column(sqlalchemy.String(255), primary_key=True)
jti = sqlalchemy.Column(sqlalchemy.String(255), primary_key=True)
expires_at = sqlalchemy.Column(sqlalchemy.DateTime(timezone=True), nullable=False)
consumed_at = sqlalchemy.Column(
sqlalchemy.DateTime(timezone=True),
nullable=False,
server_default=sqlalchemy.func.now(),
)
__table_args__ = (
sqlalchemy.Index(
'ix_space_launch_assertion_consumptions_expiry',
'instance_uuid',
'expires_at',
),
)
@@ -0,0 +1,57 @@
"""add durable replay protection for signed Space launch assertions
Revision ID: 0016_space_launch_replay
Revises: 0015_cloud_core_collab
Create Date: 2026-07-31
"""
from __future__ import annotations
import sqlalchemy as sa
from alembic import op
revision = '0016_space_launch_replay'
down_revision = '0015_cloud_core_collab'
branch_labels = None
depends_on = None
_TABLE = 'space_launch_assertion_consumptions'
_POLICY = 'langbot_directory_projection'
_SETTING = "NULLIF(current_setting('langbot.directory_instance_uuid', true), '')"
def upgrade() -> None:
conn = op.get_bind()
if _TABLE not in set(sa.inspect(conn).get_table_names()):
op.create_table(
_TABLE,
sa.Column('instance_uuid', sa.String(255), nullable=False),
sa.Column('jti', sa.String(255), nullable=False),
sa.Column('expires_at', sa.DateTime(timezone=True), nullable=False),
sa.Column('consumed_at', sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False),
sa.PrimaryKeyConstraint('instance_uuid', 'jti'),
)
op.create_index(
'ix_space_launch_assertion_consumptions_expiry',
_TABLE,
['instance_uuid', 'expires_at'],
unique=False,
)
if conn.dialect.name == 'postgresql':
table = conn.dialect.identifier_preparer.quote(_TABLE)
policy = conn.dialect.identifier_preparer.quote(_POLICY)
expression = f'instance_uuid::text = {_SETTING}'
op.execute(sa.text(f'ALTER TABLE {table} ENABLE ROW LEVEL SECURITY'))
op.execute(sa.text(f'ALTER TABLE {table} FORCE ROW LEVEL SECURITY'))
op.execute(sa.text(f'DROP POLICY IF EXISTS {policy} ON {table}'))
op.execute(
sa.text(
f'CREATE POLICY {policy} ON {table} AS PERMISSIVE FOR ALL TO PUBLIC '
f'USING ({expression}) WITH CHECK ({expression})'
)
)
def downgrade() -> None:
if _TABLE in set(sa.inspect(op.get_bind()).get_table_names()):
op.drop_table(_TABLE)
@@ -75,6 +75,7 @@ TENANT_TABLE_COLUMNS: dict[str, str] = {
DIRECTORY_PROJECTION_TABLE_COLUMNS: dict[str, str] = {
'directory_projection_states': 'instance_uuid',
'directory_projection_inbox': 'instance_uuid',
'space_launch_assertion_consumptions': 'instance_uuid',
}
DIRECTORY_PROJECTED_TENANT_TABLES = frozenset(
+1 -1
View File
@@ -892,4 +892,4 @@ class TelegramAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
await self.logger.info('Telegram adapter stopped')
self.msg_stream_id.clear()
self._form_action_titles.clear()
return True
return True
+22
View File
@@ -19,6 +19,11 @@ from urllib.parse import urljoin, urlparse
from langbot_plugin.api.entities.builtin.pipeline.query import provider_session
from ..core import app
from ..cloud.quotas import (
lock_workspace_for_quota,
require_resource_capacity,
resolve_workspace_quota,
)
from . import handler
from .archive import inspect_plugin_archive_metadata
from .github import (
@@ -1295,6 +1300,11 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
install_info: dict[str, Any],
artifact_digest: str,
) -> tuple[InstallationBinding, str | None, bool]:
quota = await resolve_workspace_quota(
self.ap,
execution_context.workspace_uuid,
'plugins.max',
)
safe_install_info = {
key: value
for key, value in install_info.items()
@@ -1316,9 +1326,19 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
)
async def persist(execute):
if quota.requires_transaction_lock:
await lock_workspace_for_quota(execute, execution_context.workspace_uuid)
result = await execute(statement)
setting = result.first()
if setting is None:
await require_resource_capacity(
execute,
workspace_uuid=execution_context.workspace_uuid,
model=persistence_plugin.PluginSetting,
quota=quota,
resource_name='plugins',
workspace_locked=quota.requires_transaction_lock,
)
installation_uuid = str(uuid.uuid4())
runtime_revision = 1
previous_digest = None
@@ -1373,6 +1393,8 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
)
tenant_uow = getattr(self.ap.persistence_mgr, 'tenant_uow', None)
if quota.requires_transaction_lock and not callable(tenant_uow):
raise RuntimeError('Cloud plugin quota enforcement requires transactional persistence')
if callable(tenant_uow):
async with tenant_uow(execution_context.workspace_uuid) as uow:
return await persist(uow.execute)
@@ -529,7 +529,7 @@ class ModelManager:
model['uuid']: model
for model in await self.ap.embedding_models_service.get_embedding_models(context, include_secret=True)
}
existing_rerank_models = {m['uuid']: m for m in await self.ap.rerank_models_service.get_rerank_models()}
existing_rerank_models = {m['uuid']: m for m in await self.ap.rerank_models_service.get_rerank_models(context)}
created = 0
updated = 0
@@ -602,6 +602,7 @@ class ModelManager:
existing = existing_rerank_models.get(space_model.uuid)
if existing is None:
await self.ap.rerank_models_service.create_rerank_model(
context,
{
'uuid': space_model.uuid,
'name': space_model.model_id,
@@ -622,7 +623,9 @@ class ModelManager:
existing.get('name') != desired['name']
or existing.get('prefered_ranking') != desired['prefered_ranking']
):
await self.ap.rerank_models_service.update_rerank_model(space_model.uuid, dict(desired))
await self.ap.rerank_models_service.update_rerank_model(
context, space_model.uuid, dict(desired)
)
updated += 1
if created or updated:
+5 -5
View File
@@ -499,14 +499,14 @@ class WorkspaceCollaborationService:
return await self._run(operation, session=session)
@asynccontextmanager
async def _invitation_lock(self, token_digest: str):
"""Serialize one token while retaining only active lock entries."""
async def _invitation_lock(self, lock_key: str):
"""Serialize one token within workspace scope while retaining only active lock entries."""
async with self._invitation_locks_guard:
entry = self._invitation_locks.get(token_digest)
entry = self._invitation_locks.get(lock_key)
if entry is None:
entry = _InvitationLockEntry(lock=asyncio.Lock())
self._invitation_locks[token_digest] = entry
self._invitation_locks[lock_key] = entry
entry.users += 1
await entry.lock.acquire()
@@ -517,7 +517,7 @@ class WorkspaceCollaborationService:
async with self._invitation_locks_guard:
entry.users -= 1
if entry.users == 0:
self._invitation_locks.pop(token_digest, None)
self._invitation_locks.pop(lock_key, None)
async def revoke_invitation(
self,
+1
View File
@@ -105,6 +105,7 @@ system:
max_bots: -1
max_pipelines: -1
max_extensions: -1
max_knowledge_bases: -1
# When set to a non-empty string, every pipeline is forced to use this
# Box sandbox-scope template regardless of its own configuration, and
# the per-pipeline "Sandbox Scope" selector is locked in the web UI.
@@ -0,0 +1,137 @@
"""PostgreSQL integration coverage for durable workspace quota locking.
Run with TEST_POSTGRES_URL=postgresql+asyncpg://... pytest ...
"""
from __future__ import annotations
import asyncio
import os
import uuid
import pytest
import sqlalchemy as sa
from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
from langbot.pkg.cloud import quotas as quota_module
from langbot.pkg.cloud.quotas import WorkspaceQuota, WorkspaceQuotaExceededError, require_resource_capacity
pytestmark = [pytest.mark.integration, pytest.mark.slow, pytest.mark.asyncio]
class _Base(DeclarativeBase):
pass
class _Workspace(_Base):
__tablename__ = 'quota_integration_workspaces'
uuid: Mapped[str] = mapped_column(sa.String(36), primary_key=True)
class _Resource(_Base):
__tablename__ = 'quota_integration_resources'
uuid: Mapped[str] = mapped_column(sa.String(36), primary_key=True)
workspace_uuid: Mapped[str] = mapped_column(
sa.String(36),
sa.ForeignKey('quota_integration_workspaces.uuid', ondelete='CASCADE'),
nullable=False,
index=True,
)
@pytest.fixture
async def quota_postgres(monkeypatch):
url = os.environ.get('TEST_POSTGRES_URL')
if not url:
pytest.skip('TEST_POSTGRES_URL not set')
if url.startswith('postgresql://'):
url = url.replace('postgresql://', 'postgresql+asyncpg://', 1)
engine = create_async_engine(url, pool_size=5, max_overflow=0)
monkeypatch.setattr(quota_module.persistence_workspace, 'Workspace', _Workspace)
async with engine.begin() as connection:
await connection.run_sync(_Base.metadata.drop_all)
await connection.run_sync(_Base.metadata.create_all)
try:
yield url, engine
finally:
async with engine.begin() as connection:
await connection.run_sync(_Base.metadata.drop_all)
await engine.dispose()
async def test_workspace_row_lock_is_atomic_isolated_and_survives_pool_restart(quota_postgres) -> None:
url, engine = quota_postgres
workspace_a = str(uuid.uuid4())
workspace_b = str(uuid.uuid4())
quota = WorkspaceQuota(limit=1, requires_transaction_lock=True)
sessions = async_sessionmaker(engine, expire_on_commit=False)
async with engine.begin() as connection:
await connection.execute(sa.insert(_Workspace), [{'uuid': workspace_a}, {'uuid': workspace_b}])
lock_acquired = asyncio.Event()
release_first = asyncio.Event()
async def admit(workspace_uuid: str, *, hold: bool = False) -> None:
async with sessions() as session:
async with session.begin():
await require_resource_capacity(
session.execute,
workspace_uuid=workspace_uuid,
model=_Resource,
quota=quota,
resource_name='resources',
)
if hold:
lock_acquired.set()
await release_first.wait()
await session.execute(
sa.insert(_Resource).values(uuid=str(uuid.uuid4()), workspace_uuid=workspace_uuid)
)
first = asyncio.create_task(admit(workspace_a, hold=True))
await asyncio.wait_for(lock_acquired.wait(), timeout=2)
same_workspace = asyncio.create_task(admit(workspace_a))
other_workspace = asyncio.create_task(admit(workspace_b))
await asyncio.wait_for(other_workspace, timeout=2)
assert not same_workspace.done(), 'same-workspace transaction bypassed SELECT FOR UPDATE'
release_first.set()
await first
with pytest.raises(WorkspaceQuotaExceededError, match=r'Maximum number of resources \(1\) reached'):
await same_workspace
async with sessions() as session:
counts = dict(
(
await session.execute(
sa.select(_Resource.workspace_uuid, sa.func.count())
.group_by(_Resource.workspace_uuid)
.order_by(_Resource.workspace_uuid)
)
).all()
)
assert counts == {workspace_a: 1, workspace_b: 1}
await engine.dispose()
restarted_engine = create_async_engine(url, pool_size=2, max_overflow=0)
restarted_sessions = async_sessionmaker(restarted_engine, expire_on_commit=False)
try:
async with restarted_sessions() as session:
async with session.begin():
with pytest.raises(WorkspaceQuotaExceededError):
await require_resource_capacity(
session.execute,
workspace_uuid=workspace_a,
model=_Resource,
quota=quota,
resource_name='resources',
)
finally:
await restarted_engine.dispose()
@@ -9,6 +9,7 @@ import quart
from langbot.pkg.api.http.controller import group
from langbot.pkg.api.http.controller.groups.webhooks import WebhookRouterGroup
from langbot.pkg.cloud.quotas import WorkspaceQuotaExceededError
from langbot.pkg.utils.bounded_executor import (
BlockingWorkCapacityError,
current_blocking_work_scope,
@@ -48,6 +49,16 @@ class _BlockingCapacityRouterGroup(group.RouterGroup):
raise BlockingWorkCapacityError('Workspace blocking executor capacity reached')
class _QuotaRouterGroup(group.RouterGroup):
name = 'quota-test'
path = '/quota-test'
async def initialize(self) -> None:
@self.route('', methods=['POST'], auth_type=group.AuthType.NONE)
async def _():
raise WorkspaceQuotaExceededError('bots', 2)
class _InvalidAccountRouterGroup(group.RouterGroup):
name = 'invalid-account-test'
path = '/invalid-account-test'
@@ -130,6 +141,20 @@ async def test_blocking_work_capacity_maps_to_retryable_http_response():
}
async def test_workspace_quota_maps_to_stable_conflict_response():
application = SimpleNamespace(logger=Mock())
quart_app = quart.Quart(__name__)
await _QuotaRouterGroup(application, quart_app).initialize()
response = await quart_app.test_client().post('/quota-test')
assert response.status_code == 409
assert await response.get_json() == {
'code': 'workspace_quota_exceeded',
'msg': 'Maximum number of bots (2) reached',
}
async def test_public_webhook_carries_scope_without_holding_database_session():
class ScopeOnlyPersistenceManager:
mode = SimpleNamespace(value='cloud_runtime')
@@ -311,10 +311,9 @@ class TestBotServiceCreateBot:
ap.platform_mgr = SimpleNamespace()
ap.platform_mgr.load_bot = AsyncMock()
# Mock get_bots to return 2 bots already
bot1 = _create_mock_bot(bot_uuid='uuid-1')
bot2 = _create_mock_bot(bot_uuid='uuid-2')
mock_result = _create_mock_result([bot1, bot2])
# Mock the atomic count query to report 2 existing bots.
mock_result = _create_mock_result()
mock_result.scalar_one = Mock(return_value=2)
ap.persistence_mgr.execute_async = AsyncMock(return_value=mock_result)
ap.persistence_mgr.serialize_model = Mock(return_value={'uuid': 'uuid-1', 'name': 'Bot 1'})
@@ -0,0 +1,129 @@
from __future__ import annotations
import asyncio
from collections import defaultdict
from types import SimpleNamespace
from unittest.mock import AsyncMock
import pytest
import sqlalchemy
from langbot.pkg.api.http.service.bot import BotService
from langbot.pkg.cloud.entitlements import EntitlementResolver, EntitlementSnapshot
INSTANCE_UUID = 'cloud-instance'
WORKSPACE_A = '11111111-1111-1111-1111-111111111111'
WORKSPACE_B = '22222222-2222-2222-2222-222222222222'
class _Provider:
async def get_workspace_entitlement(self, workspace_uuid: str) -> EntitlementSnapshot:
return EntitlementSnapshot(
instance_uuid=INSTANCE_UUID,
workspace_uuid=workspace_uuid,
entitlement_revision=1,
status='active',
not_before=0,
expires_at=4_102_444_800,
features={},
limits={'bots.max': 2},
)
class _Result:
def __init__(self, *, first=None, scalar=None) -> None:
self._first = first
self._scalar = scalar
def first(self):
return self._first
def scalar_one(self):
return self._scalar
class _TenantUow:
def __init__(self, manager: '_Persistence', workspace_uuid: str) -> None:
self.manager = manager
self.workspace_uuid = workspace_uuid
self.lock = manager.locks[workspace_uuid]
async def __aenter__(self):
await self.lock.acquire()
return self
async def __aexit__(self, exc_type, exc, tb):
self.lock.release()
async def execute(self, statement):
sql = str(statement)
if isinstance(statement, sqlalchemy.sql.dml.Insert):
assert statement.table.name == 'bots'
self.manager.bots[self.workspace_uuid].append(statement.compile().params)
return _Result()
if 'FROM workspaces' in sql:
assert statement._for_update_arg is not None
self.manager.workspace_locks_seen += 1
return _Result(first=(self.workspace_uuid,))
if 'count(' in sql.lower() and 'FROM bots' in sql:
return _Result(scalar=len(self.manager.bots[self.workspace_uuid]))
if 'FROM legacy_pipelines' in sql:
return _Result(first=None)
raise AssertionError(f'unexpected statement: {sql}')
class _Persistence:
def __init__(self) -> None:
self.locks = defaultdict(asyncio.Lock)
self.bots = defaultdict(list)
self.workspace_locks_seen = 0
def tenant_uow(self, workspace_uuid: str) -> _TenantUow:
return _TenantUow(self, workspace_uuid)
async def execute_async(self, statement):
assert 'FROM legacy_pipelines' in str(statement)
return _Result(first=None)
async def _service(manager: _Persistence) -> BotService:
resolver = EntitlementResolver(INSTANCE_UUID, _Provider())
await resolver.reconcile_active_workspaces({WORKSPACE_A, WORKSPACE_B})
ap = SimpleNamespace(
entitlement_resolver=resolver,
persistence_mgr=manager,
instance_config=SimpleNamespace(data={'system': {'limitation': {'max_bots': 99}}}),
platform_mgr=SimpleNamespace(load_bot=AsyncMock()),
)
service = BotService(ap)
service.get_bot = AsyncMock(return_value={'uuid': 'created'})
return service
@pytest.mark.asyncio
async def test_cloud_bot_quota_is_atomic_isolated_and_persists_across_service_restart() -> None:
manager = _Persistence()
service = await _service(manager)
async def create(workspace_uuid: str, index: int):
return await service.create_bot(workspace_uuid, {'name': f'bot-{index}'})
results = await asyncio.gather(
*(create(WORKSPACE_A, index) for index in range(8)),
*(create(WORKSPACE_B, index) for index in range(8)),
return_exceptions=True,
)
successes = [result for result in results if isinstance(result, str)]
failures = [result for result in results if isinstance(result, ValueError)]
assert len(successes) == 4
assert len(failures) == 12
assert len(manager.bots[WORKSPACE_A]) == 2
assert len(manager.bots[WORKSPACE_B]) == 2
assert manager.workspace_locks_seen == 16
restarted_service = await _service(manager)
with pytest.raises(ValueError, match=r'Maximum number of bots \(2\) reached'):
await restarted_service.create_bot(WORKSPACE_A, {'name': 'after-restart'})
assert len(manager.bots[WORKSPACE_A]) == 2
@@ -34,6 +34,7 @@ class _Rows:
def _app():
return SimpleNamespace(
logger=Mock(),
instance_config=SimpleNamespace(data={}),
rag_mgr=SimpleNamespace(
get_all_knowledge_base_details=AsyncMock(return_value=[]),
get_knowledge_base_details=AsyncMock(return_value=None),
@@ -119,6 +120,21 @@ async def test_create_validates_schema_and_binds_context():
)
@pytest.mark.asyncio
async def test_create_enforces_workspace_knowledge_base_limit():
app = _app()
app.instance_config.data = {'system': {'limitation': {'max_knowledge_bases': 2}}}
app.rag_mgr.get_all_knowledge_base_details.return_value = [{'uuid': 'kb-a'}, {'uuid': 'kb-b'}]
service = KnowledgeService(app)
with pytest.raises(ValueError, match=r'Maximum number of knowledge bases \(2\) reached'):
await service.create_knowledge_base(
CONTEXT,
{'knowledge_engine_plugin_id': 'author/engine'},
)
app.rag_mgr.create_knowledge_base.assert_not_awaited()
@pytest.mark.asyncio
async def test_update_rejects_guessed_uuid_and_scopes_reload():
app = _app()
@@ -862,11 +862,11 @@ class TestModelProviderServiceScanProviderModels:
ap.model_mgr.load_provider = AsyncMock(return_value=runtime_provider)
ap.llm_model_service.get_llm_models_by_provider = AsyncMock(return_value=[])
ap.embedding_models_service.get_embedding_models_by_provider = AsyncMock(return_value=[])
ap.rerank_models_service.get_rerank_models_by_provider = AsyncMock(
return_value=[{'name': 'Qwen3-Reranker-8B'}]
)
ap.rerank_models_service.get_rerank_models_by_provider = AsyncMock(return_value=[{'name': 'Qwen3-Reranker-8B'}])
result = await ModelProviderService(ap).scan_provider_models('rerank-scan-uuid', model_type='rerank')
result = await ModelProviderService(ap).scan_provider_models(
WORKSPACE_UUID, 'rerank-scan-uuid', model_type='rerank'
)
assert result['models'][0]['type'] == 'rerank'
assert result['models'][0]['already_added'] is True
+14 -1
View File
@@ -228,10 +228,23 @@ async def test_cloud_pgvector_contract_is_fail_closed(pgvector_config, message):
)
async def test_cloud_runtime_allows_explicitly_disabled_box():
config = _cloud_config()
config['box']['enabled'] = False
deployment = await resolve_deployment(
instance_uuid='instance-a',
instance_config=config,
entry_points=lambda: _EntryPoints([_EntryPoint(_Provider())]),
now=1_000,
)
assert isinstance(deployment, VerifiedCloudDeployment)
@pytest.mark.parametrize(
('mutate', 'message'),
[
(lambda config: config['box'].update(enabled=False), 'box.enabled=true'),
(lambda config: config['box'].update(backend='docker'), 'box.backend=nsjail'),
(lambda config: config['box']['runtime'].update(endpoint=''), 'box.runtime.endpoint'),
(
+34 -1
View File
@@ -48,6 +48,7 @@ def _claims(*, now: int, jti: str | None = None, workspace_uuid: str = WORKSPACE
'payload': {
'account_uuid': ACCOUNT_UUID,
'workspace_uuid': workspace_uuid,
'return_path': '/',
},
}
@@ -81,7 +82,27 @@ async def test_consumes_valid_workspace_launch_assertion_once():
launch = await service.consume_assertion(token, expected_workspace_uuid=WORKSPACE_UUID)
assert launch == {'account_uuid': ACCOUNT_UUID, 'workspace_uuid': WORKSPACE_UUID}
assert launch == {
'account_uuid': ACCOUNT_UUID,
'workspace_uuid': WORKSPACE_UUID,
'return_path': '/',
}
with pytest.raises(SpaceLaunchError, match='already been consumed'):
await service.consume_assertion(token, expected_workspace_uuid=WORKSPACE_UUID)
async def test_consumed_assertion_remains_blocked_through_clock_skew_window():
private_key = Ed25519PrivateKey.generate()
now = int(time.time())
service = _service(private_key, now=now)
claims = _claims(now=now)
claims['iat'] = now - 10
claims['nbf'] = now - 10
claims['exp'] = now - 1
token = _sign(private_key, claims)
await service.consume_assertion(token, expected_workspace_uuid=WORKSPACE_UUID)
with pytest.raises(SpaceLaunchError, match='already been consumed'):
await service.consume_assertion(token, expected_workspace_uuid=WORKSPACE_UUID)
@@ -162,3 +183,15 @@ async def test_rejects_invalid_signature_and_non_cloud_mode():
oss_service.ap.deployment.multi_workspace_enabled = False
with pytest.raises(SpaceLaunchError, match='verified Cloud mode'):
await oss_service.consume_assertion(token, expected_workspace_uuid=WORKSPACE_UUID)
@pytest.mark.asyncio
async def test_rejects_unsafe_signed_return_path() -> None:
private_key = Ed25519PrivateKey.generate()
now = int(time.time())
service = _service(private_key, now=now)
claims = _claims(now=now)
claims['payload']['return_path'] = '//evil.example'
with pytest.raises(SpaceLaunchError, match='return path'):
await service.consume_assertion(_sign(private_key, claims))
@@ -0,0 +1,102 @@
from __future__ import annotations
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock
import pytest
import sqlalchemy
from langbot.pkg.cloud.entitlements import EntitlementResolver, EntitlementSnapshot
from langbot.pkg.cloud.quotas import (
WorkspaceQuota,
require_resource_capacity,
resolve_workspace_quota,
)
from langbot.pkg.entity.persistence.bot import Bot
WORKSPACE_UUID = '11111111-1111-1111-1111-111111111111'
INSTANCE_UUID = 'cloud-instance'
class _Provider:
def __init__(self, limits: dict[str, int]) -> None:
self.limits = limits
async def get_workspace_entitlement(self, workspace_uuid: str) -> EntitlementSnapshot:
return EntitlementSnapshot(
instance_uuid=INSTANCE_UUID,
workspace_uuid=workspace_uuid,
entitlement_revision=1,
status='active',
not_before=0,
expires_at=4_102_444_800,
features={},
limits=self.limits,
plan_name='test',
)
async def _resolver(limits: dict[str, int]) -> EntitlementResolver:
resolver = EntitlementResolver(INSTANCE_UUID, _Provider(limits))
await resolver.reconcile_active_workspaces({WORKSPACE_UUID})
return resolver
@pytest.mark.asyncio
async def test_resolve_workspace_quota_uses_signed_cloud_limit() -> None:
ap = SimpleNamespace(entitlement_resolver=await _resolver({'bots.max': 2}))
quota = await resolve_workspace_quota(ap, WORKSPACE_UUID, 'bots.max', fallback=99)
assert quota == WorkspaceQuota(limit=2, requires_transaction_lock=True)
@pytest.mark.asyncio
async def test_resolve_workspace_quota_preserves_oss_fallback() -> None:
quota = await resolve_workspace_quota(SimpleNamespace(), WORKSPACE_UUID, 'bots.max', fallback=7)
assert quota == WorkspaceQuota(limit=7, requires_transaction_lock=False)
@pytest.mark.asyncio
async def test_require_resource_capacity_locks_workspace_before_counting() -> None:
statements: list[object] = []
lock_result = Mock()
lock_result.first.return_value = (WORKSPACE_UUID,)
count_result = Mock()
count_result.scalar_one.return_value = 1
execute = AsyncMock(side_effect=[lock_result, count_result])
await require_resource_capacity(
execute,
workspace_uuid=WORKSPACE_UUID,
model=Bot,
quota=WorkspaceQuota(limit=2, requires_transaction_lock=True),
resource_name='bots',
)
statements.extend(call.args[0] for call in execute.await_args_list)
assert len(statements) == 2
assert isinstance(statements[0], sqlalchemy.sql.Select)
assert statements[0]._for_update_arg is not None
assert 'workspaces' in str(statements[0])
assert 'count' in str(statements[1]).lower()
@pytest.mark.asyncio
async def test_require_resource_capacity_rejects_at_boundary() -> None:
lock_result = Mock()
lock_result.first.return_value = (WORKSPACE_UUID,)
count_result = Mock()
count_result.scalar_one.return_value = 2
execute = AsyncMock(side_effect=[lock_result, count_result])
with pytest.raises(ValueError, match=r'Maximum number of bots \(2\) reached'):
await require_resource_capacity(
execute,
workspace_uuid=WORKSPACE_UUID,
model=Bot,
quota=WorkspaceQuota(limit=2, requires_transaction_lock=True),
resource_name='bots',
)
+13
View File
@@ -152,6 +152,19 @@ class TestApplyEnvOverridesToConfig:
assert result['system']['disabled_adapters'] == ['aiocqhttp', 'dingtalk', 'telegram']
def test_override_integer_list_preserves_item_type(self):
"""Comma-separated overrides inherit the existing list item type."""
load_config = get_load_config_module()
cfg = {'vdb': {'pgvector': {'allowed_dimensions': [384, 512]}}}
env = {'VDB__PGVECTOR__ALLOWED_DIMENSIONS': '384,512,768'}
with patch.dict(os.environ, env, clear=True):
result = load_config._apply_env_overrides_to_config(cfg)
assert result['vdb']['pgvector']['allowed_dimensions'] == [384, 512, 768]
assert all(isinstance(item, int) for item in result['vdb']['pgvector']['allowed_dimensions'])
def test_override_list_value_empty_items(self):
"""Test that empty items in comma-separated list are filtered."""
load_config = get_load_config_module()
@@ -0,0 +1,191 @@
from __future__ import annotations
import asyncio
from collections import defaultdict
from types import SimpleNamespace
import pytest
import sqlalchemy
from langbot.pkg.api.http.context import ExecutionContext
from langbot.pkg.cloud.entitlements import EntitlementResolver, EntitlementSnapshot
from langbot.pkg.cloud.quotas import WorkspaceQuotaExceededError
from langbot.pkg.plugin.connector import PluginRuntimeConnector
from langbot_plugin.runtime.plugin.mgr import PluginInstallSource
INSTANCE_UUID = 'cloud-instance'
WORKSPACE_A = '11111111-1111-1111-1111-111111111111'
WORKSPACE_B = '22222222-2222-2222-2222-222222222222'
class _Provider:
async def get_workspace_entitlement(self, workspace_uuid: str) -> EntitlementSnapshot:
return EntitlementSnapshot(
instance_uuid=INSTANCE_UUID,
workspace_uuid=workspace_uuid,
entitlement_revision=1,
status='active',
not_before=0,
expires_at=4_102_444_800,
features={},
limits={'plugins.max': 3},
)
class _Result:
def __init__(self, *, first=None, scalar=None) -> None:
self._first = first
self._scalar = scalar
def first(self):
return self._first
def scalar_one(self):
return self._scalar
class _TenantUow:
def __init__(self, manager: '_Persistence', workspace_uuid: str) -> None:
self.manager = manager
self.workspace_uuid = workspace_uuid
self.lock = manager.locks[workspace_uuid]
async def __aenter__(self):
await self.lock.acquire()
return self
async def __aexit__(self, exc_type, exc, tb):
self.lock.release()
async def execute(self, statement):
sql = str(statement)
params = statement.compile().params
if isinstance(statement, sqlalchemy.sql.dml.Insert):
assert statement.table.name == 'plugin_settings'
key = (params['plugin_author'], params['plugin_name'])
self.manager.plugins[self.workspace_uuid][key] = dict(params)
return _Result()
if isinstance(statement, sqlalchemy.sql.dml.Update):
return _Result()
if 'FROM workspaces' in sql:
assert statement._for_update_arg is not None
self.manager.workspace_locks_seen += 1
return _Result(first=(self.workspace_uuid,))
if 'count(' in sql.lower() and 'FROM plugin_settings' in sql:
return _Result(scalar=len(self.manager.plugins[self.workspace_uuid]))
if 'FROM plugin_settings' in sql:
author = next(value for name, value in params.items() if 'plugin_author' in name)
name = next(value for param, value in params.items() if 'plugin_name' in param)
row = self.manager.plugins[self.workspace_uuid].get((author, name))
if row is None:
return _Result(first=None)
return _Result(
first=SimpleNamespace(
installation_uuid=row['installation_uuid'],
runtime_revision=row['runtime_revision'],
artifact_digest=row['artifact_digest'],
install_info=row['install_info'],
)
)
raise AssertionError(f'unexpected statement: {sql}')
class _Persistence:
def __init__(self) -> None:
self.locks = defaultdict(asyncio.Lock)
self.plugins = defaultdict(dict)
self.workspace_locks_seen = 0
def tenant_uow(self, workspace_uuid: str) -> _TenantUow:
return _TenantUow(self, workspace_uuid)
async def _connector(manager: _Persistence) -> PluginRuntimeConnector:
resolver = EntitlementResolver(INSTANCE_UUID, _Provider())
await resolver.reconcile_active_workspaces({WORKSPACE_A, WORKSPACE_B})
connector = object.__new__(PluginRuntimeConnector)
connector.ap = SimpleNamespace(entitlement_resolver=resolver, persistence_mgr=manager)
return connector
def _context(workspace_uuid: str) -> ExecutionContext:
return ExecutionContext(
instance_uuid=INSTANCE_UUID,
workspace_uuid=workspace_uuid,
placement_generation=1,
entitlement_revision=1,
)
@pytest.mark.asyncio
async def test_cloud_plugin_quota_is_atomic_isolated_and_persists_across_connector_restart() -> None:
manager = _Persistence()
connector = await _connector(manager)
async def install(workspace_uuid: str, index: int):
return await connector._persist_installation_package(
_context(workspace_uuid),
plugin_author='test-author',
plugin_name=f'plugin-{index}',
install_source=PluginInstallSource.MARKETPLACE,
install_info={'author': 'test-author', 'name': f'plugin-{index}'},
artifact_digest=f'{index:064x}',
)
results = await asyncio.gather(
*(install(WORKSPACE_A, index) for index in range(10)),
*(install(WORKSPACE_B, index) for index in range(10)),
return_exceptions=True,
)
successes = [result for result in results if isinstance(result, tuple)]
failures = [result for result in results if isinstance(result, WorkspaceQuotaExceededError)]
assert len(successes) == 6
assert len(failures) == 14
assert len(manager.plugins[WORKSPACE_A]) == 3
assert len(manager.plugins[WORKSPACE_B]) == 3
assert manager.workspace_locks_seen == 20
restarted_connector = await _connector(manager)
with pytest.raises(WorkspaceQuotaExceededError, match=r'Maximum number of plugins \(3\) reached'):
await restarted_connector._persist_installation_package(
_context(WORKSPACE_A),
plugin_author='test-author',
plugin_name='after-restart',
install_source=PluginInstallSource.MARKETPLACE,
install_info={},
artifact_digest='f' * 64,
)
assert len(manager.plugins[WORKSPACE_A]) == 3
installed_name = next(iter(manager.plugins[WORKSPACE_A]))[1]
async def reinstall():
return await restarted_connector._persist_installation_package(
_context(WORKSPACE_A),
plugin_author='test-author',
plugin_name=installed_name,
install_source=PluginInstallSource.MARKETPLACE,
install_info={'author': 'test-author', 'name': installed_name, 'revision': 2},
artifact_digest='e' * 64,
)
reinstall_results = await asyncio.gather(reinstall(), reinstall())
assert all(result[2] is True for result in reinstall_results)
mixed_results = await asyncio.gather(
reinstall(),
restarted_connector._persist_installation_package(
_context(WORKSPACE_A),
plugin_author='test-author',
plugin_name='new-at-capacity',
install_source=PluginInstallSource.MARKETPLACE,
install_info={},
artifact_digest='d' * 64,
),
return_exceptions=True,
)
assert isinstance(mixed_results[0], tuple)
assert isinstance(mixed_results[1], WorkspaceQuotaExceededError)
assert len(manager.plugins[WORKSPACE_A]) == 3
@@ -114,9 +114,10 @@ async def test_sync_new_models_from_space_creates_rerank_models(mock_app_for_mod
app.rerank_models_service.get_rerank_models = AsyncMock(return_value=[])
model_mgr = ModelManager(app)
await model_mgr.sync_new_models_from_space()
await model_mgr.sync_new_models_from_space(TEST_EXECUTION_CONTEXT)
app.rerank_models_service.create_rerank_model.assert_awaited_once_with(
TEST_EXECUTION_CONTEXT,
{
'uuid': 'rerank-model-uuid',
'name': 'Qwen3-Reranker-8B',
@@ -1109,4 +1110,4 @@ def test_provider_not_found_error_str():
error = provider_errors.ProviderNotFoundError('test-provider')
assert str(error) == 'Provider test-provider not found'
assert error.provider_name == 'test-provider'
assert error.provider_name == 'test-provider'
Generated
+7 -3
View File
@@ -2116,7 +2116,7 @@ requires-dist = [
{ name = "ebooklib", specifier = ">=0.18" },
{ name = "gewechat-client", specifier = ">=0.1.5" },
{ name = "html2text", specifier = ">=2024.2.26" },
{ name = "langbot-plugin", git = "https://github.com/langbot-app/langbot-plugin-sdk.git?rev=1d65ed301a6afc52150a998043f73cd6032c8162" },
{ name = "langbot-plugin", specifier = "==0.5.0" },
{ name = "langchain", specifier = ">=1.3.9" },
{ name = "langchain-core", specifier = ">=1.3.3" },
{ name = "langchain-text-splitters", specifier = ">=1.1.2" },
@@ -2182,8 +2182,8 @@ dev = [
[[package]]
name = "langbot-plugin"
version = "0.4.18"
source = { git = "https://github.com/langbot-app/langbot-plugin-sdk.git?rev=1d65ed301a6afc52150a998043f73cd6032c8162#1d65ed301a6afc52150a998043f73cd6032c8162" }
version = "0.5.0"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "aiofiles" },
{ name = "aiohttp" },
@@ -2203,6 +2203,10 @@ dependencies = [
{ name = "watchdog" },
{ name = "websockets" },
]
sdist = { url = "https://files.pythonhosted.org/packages/5c/09/697037dea617a235c9b3df85174badfb44886e3c67149a928049839a959a/langbot_plugin-0.5.0.tar.gz", hash = "sha256:9b81fa0f73cde1fe199a746e03ad891f8ced8f4fa0d05aeb6d78bfe7d755a976", size = 464033, upload-time = "2026-07-30T19:22:57.503Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/d2/70/93e544c3a120953f4871266db9ef84ea9a58daf01d493fafa7c45674ad36/langbot_plugin-0.5.0-py3-none-any.whl", hash = "sha256:2b078db96d869de55d08304465974f51765e3cfe582ff79b32f09161292fb877", size = 300758, upload-time = "2026-07-30T19:22:56.248Z" },
]
[[package]]
name = "langchain"
+35 -3
View File
@@ -29,6 +29,7 @@ type SpaceOAuthLoginResult = {
token: string;
user: string;
workspace_uuid?: string;
return_path?: string;
};
const pendingSpaceOAuthLogins = new Map<
@@ -63,6 +64,10 @@ function SpaceOAuthCallbackContent() {
const [searchParams] = useSearchParams();
const { t } = useTranslation();
const isMountedRef = useRef(true);
const directLaunchFragmentRef = useRef<{
workspaceUuid: string | null;
launchAssertion: string | null;
} | null>(null);
const [status, setStatus] = useState<
'loading' | 'confirm' | 'success' | 'error'
@@ -106,7 +111,13 @@ function SpaceOAuthCallbackContent() {
throw new Error('No Workspace is available for this Account');
}
if (response.workspace_uuid) {
navigate('/home', { replace: true });
const returnPath =
typeof response.return_path === 'string' &&
response.return_path.startsWith('/') &&
!response.return_path.startsWith('//')
? response.return_path
: '/home';
navigate(returnPath, { replace: true });
return;
}
setStatus('success');
@@ -207,8 +218,29 @@ function SpaceOAuthCallbackContent() {
const errorDescription = searchParams.get('error_description');
const mode = searchParams.get('mode');
const state = searchParams.get('state');
const workspaceUuid = searchParams.get('workspace_uuid');
const launchAssertion = searchParams.get('launch_assertion');
if (directLaunchFragmentRef.current === null) {
const fragmentParams = new URLSearchParams(
window.location.hash.startsWith('#')
? window.location.hash.slice(1)
: window.location.hash,
);
directLaunchFragmentRef.current = {
workspaceUuid: fragmentParams.get('workspace_uuid'),
launchAssertion: fragmentParams.get('launch_assertion'),
};
if (window.location.hash) {
window.history.replaceState(
null,
'',
`${window.location.pathname}${window.location.search}`,
);
}
}
const workspaceUuid =
directLaunchFragmentRef.current.workspaceUuid ??
searchParams.get('workspace_uuid');
const launchAssertion =
directLaunchFragmentRef.current.launchAssertion;
if (error) {
setStatus('error');
@@ -0,0 +1,24 @@
import assert from 'node:assert/strict';
import fs from 'node:fs';
import test from 'node:test';
const source = fs.readFileSync(
new URL('../../src/app/auth/space/callback/page.tsx', import.meta.url),
'utf8',
);
test('direct launch assertion is fragment-only and removed before exchange', () => {
assert.doesNotMatch(source, /searchParams\.get\(['"]launch_assertion['"]\)/);
const readIndex = source.indexOf("fragmentParams.get('launch_assertion')");
const clearIndex = source.indexOf('window.history.replaceState');
const exchangeIndex = source.indexOf('handleOAuthCallback(', clearIndex);
assert.ok(readIndex >= 0, 'fragment assertion read is missing');
assert.ok(clearIndex > readIndex, 'URL fragment is not cleared after copying the assertion');
assert.ok(exchangeIndex > clearIndex, 'assertion exchange starts before the fragment is cleared');
});
test('direct launch honors only a local signed return path', () => {
assert.match(source, /response\.return_path\.startsWith\(['"]\/['"]\)/);
assert.match(source, /!response\.return_path\.startsWith\(['"]\/\/['"]\)/);
assert.match(source, /navigate\(returnPath, \{ replace: true \}\)/);
});