Compare commits

..

10 Commits

Author SHA1 Message Date
fdc310 cd7acd5769 style(space): apply ruff formatting 2026-08-26 00:13:04 +08:00
fdc310 51b9e1cf54 Merge remote-tracking branch 'origin/master' into feat/rework-agent-onboarding 2026-08-25 22:36:44 +08:00
fdc310 1f07471f81 feat(wizard): label page bot test preview 2026-08-25 22:34:45 +08:00
fdc310 7f29d51c69 feat(wizard): streamline custom model onboarding 2026-08-25 19:33:36 +08:00
fdc310 6b60ad1678 fix(wizard): repair HTTP bot inbound test setup 2026-08-23 20:41:46 +08:00
langbot-dev 7a4fadc375 feat(wizard): add floating page bot verification 2026-08-14 23:58:39 +08:00
langbot-dev f2ba540ffb feat(wizard): add inbound bot verification 2026-08-14 23:14:22 +08:00
langbot-dev 97b176aef2 fix(wizard): parse ranked model selection entries 2026-08-14 17:43:36 +08:00
langbot-dev 7c387f75e1 fix(web): support LAN development access 2026-08-14 17:29:21 +08:00
langbot-dev b0566f4c9d feat(wizard): rework agent onboarding flow 2026-08-14 01:10:00 +08:00
142 changed files with 450 additions and 3994 deletions
+2 -2
View File
@@ -1,5 +1,5 @@
name: 漏洞反馈
description: 【供中文用户】报错或漏洞请使用这个模板创建,不使用此模板创建的异常、漏洞相关issue将被直接关闭。由于自己操作不当/不甚了解所用技术栈引起的网络连接问题恕无法解决,请勿提 issue。容器间网络连接问题,参考文档 https://langbot.app/docs/zh/workshop/network-details
description: 【供中文用户】报错或漏洞请使用这个模板创建,不使用此模板创建的异常、漏洞相关issue将被直接关闭。由于自己操作不当/不甚了解所用技术栈引起的网络连接问题恕无法解决,请勿提 issue。容器间网络连接问题,参考文档 https://link.langbot.app/zh/docs/network
title: "[Bug]: "
labels: ["bug?"]
body:
@@ -22,7 +22,7 @@ body:
- type: textarea
attributes:
label: 异常情况
description: 完整描述异常情况,什么时候发生的、发生了什么。**请附带日志信息。**
description: 完整描述异常情况,什么时候发生的、发生了什么。**请附带日志信息。**
validations:
required: true
- type: textarea
+1 -1
View File
@@ -1,5 +1,5 @@
name: Bug report
description: Report bugs or vulnerabilities using this template. For container network connection issues, refer to the documentation https://langbot.app/docs/en/workshop/network-details
description: Report bugs or vulnerabilities using this template. For container network connection issues, refer to the documentation https://link.langbot.app/en/docs/network
title: "[Bug]: "
labels: ["bug?"]
body:
+2 -2
View File
@@ -43,8 +43,8 @@ Run the narrowest useful test first, then broader checks when confidence is need
## Where to Look
- Architecture map: `ARCHITECTURE.md`.
- Dev environment guide: https://langbot.app/docs/zh/develop/dev-config.
- Plugin runtime / CLI / SDK debugging: https://langbot.app/docs/zh/develop/plugin-runtime.
- Dev environment guide: https://docs.langbot.app/zh/develop/dev-config.
- Plugin runtime / CLI / SDK debugging: https://docs.langbot.app/zh/develop/plugin-runtime.
- API-key auth: `docs/API_KEY_AUTH.md`.
- Box deep-dive notes: `docs/review/box-architecture.md` and related files.
- In-repo skills: `skills/` is the single source of truth for LangBot agent skills.
+6 -6
View File
@@ -19,9 +19,9 @@ English / [简体中文](README_CN.md) / [繁體中文](README_TW.md) / [日本
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Website</a>
<a href="https://langbot.app/docs/en/insight/features">Features</a>
<a href="https://langbot.app/docs/en/insight/guide">Docs</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Features</a>
<a href="https://link.langbot.app/en/docs/guide">Docs</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app/cloud">Cloud</a>
<a href="https://space.langbot.app">Plugin Market</a>
<a href="https://langbot.featurebase.app/roadmap">Roadmap</a>
@@ -49,7 +49,7 @@ LangBot is an **open-source, production-grade platform** for building AI-powered
- **Web Management Panel** — Configure, manage, and monitor your bots through an intuitive browser interface. No YAML editing required.
- **Multi-Pipeline Architecture** — Different bots for different scenarios, with comprehensive monitoring and exception handling.
[→ Learn more about all features](https://langbot.app/docs/en/insight/features)
[→ Learn more about all features](https://link.langbot.app/en/docs/features)
📍 Practical guides: [deploy a multi-platform AI bot in 5 minutes](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [connect DeepSeek to WeChat, Discord, and Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [run a Dify Agent in Discord, Telegram, and Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/), and [build an n8n-powered chatbot](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -89,7 +89,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**More options:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Manual](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**More options:** [Docker](https://link.langbot.app/en/docs/docker) · [Manual](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
@@ -151,7 +151,7 @@ _Note: Public demo environment. Do not enter sensitive information._
| [302.AI](https://share.302ai.cn/SuTG99) | Gateway | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Gateway | ✅ |
[→ View all integrations](https://langbot.app/docs/en/insight/features)
[→ View all integrations](https://link.langbot.app/en/docs/features)
---
+6 -6
View File
@@ -21,9 +21,9 @@
[![star](https://gitcode.com/RockChinQ/LangBot/star/badge.svg)](https://gitcode.com/RockChinQ/LangBot)
<a href="https://langbot.app">官网</a>
<a href="https://langbot.app/docs/zh/insight/features">特性</a>
<a href="https://langbot.app/docs/zh/insight/guide">文档</a>
<a href="https://langbot.app/docs/zh/tags/readme">API</a>
<a href="https://link.langbot.app/zh/docs/features">特性</a>
<a href="https://link.langbot.app/zh/docs/guide">文档</a>
<a href="https://link.langbot.app/zh/docs/api">API</a>
<a href="https://space.langbot.app/cloud">Cloud</a>
<a href="https://space.langbot.app">扩展市场</a>
<a href="https://langbot.featurebase.app/roadmap">路线图</a>
@@ -49,7 +49,7 @@ LangBot 是一个**开源的生产级平台**,用于构建 AI 驱动的即时
- **Web 管理面板** — 通过浏览器直观地配置、管理和监控机器人,无需手动编辑配置文件。
- **多流水线架构** — 不同机器人用于不同场景,具备全面的监控和异常处理能力。
[→ 了解更多功能特性](https://langbot.app/docs/zh/insight/features)
[→ 了解更多功能特性](https://link.langbot.app/zh/docs/features)
📍 实践指南:[5 分钟部署多平台 AI 机器人](https://langbot.app/zh/blog/deploy-ai-bot-in-5-minutes/)、[将 DeepSeek 接入微信、企业微信与 Discord](https://langbot.app/zh/blog/connect-deepseek-to-wechat/)、[让 Dify Agent 跑在 Discord、Telegram 和 Slack 上](https://langbot.app/zh/blog/dify-agent-discord-telegram-slack/),以及[用 n8n 构建多平台 AI 聊天机器人](https://langbot.app/zh/blog/n8n-multi-platform-ai-chatbot/)。
@@ -89,7 +89,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/zh-CN/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**更多方式:** [Docker](https://langbot.app/docs/zh/deploy/langbot/docker) · [手动部署](https://langbot.app/docs/zh/deploy/langbot/manual) · [宝塔面板](https://langbot.app/docs/zh/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/zh/deploy/langbot/kubernetes)
**更多方式:** [Docker](https://link.langbot.app/zh/docs/docker) · [手动部署](https://link.langbot.app/zh/docs/manual-deploy) · [宝塔面板](https://link.langbot.app/zh/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/zh/deploy/langbot/kubernetes)
---
@@ -152,7 +152,7 @@ docker compose --profile all up -d
| [百宝箱Tbox](https://www.tbox.cn/open) | 智能体平台 | ✅ |
| [七牛云Qiniu](https://www.qiniu.com/ai/agent) | 聚合平台 | ✅ |
[→ 查看完整集成列表](https://langbot.app/docs/zh/insight/features)
[→ 查看完整集成列表](https://link.langbot.app/zh/docs/features)
### TTS(语音合成)
+6 -6
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Inicio</a>
<a href="https://langbot.app/docs/en/insight/features">Características</a>
<a href="https://langbot.app/docs/en/insight/guide">Documentación</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Características</a>
<a href="https://link.langbot.app/en/docs/guide">Documentación</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Mercado de Plugins</a>
<a href="https://langbot.featurebase.app/roadmap">Hoja de Ruta</a>
@@ -48,7 +48,7 @@ LangBot es una **plataforma de código abierto y grado de producción** para con
- **Panel de Gestión Web** — Configure, gestione y monitoree sus bots a través de una interfaz de navegador intuitiva. Sin necesidad de editar YAML.
- **Arquitectura Multi-Pipeline** — Diferentes bots para diferentes escenarios, con monitoreo completo y manejo de excepciones.
[→ Conocer más sobre todas las funcionalidades](https://langbot.app/docs/en/insight/features)
[→ Conocer más sobre todas las funcionalidades](https://link.langbot.app/en/docs/features)
📍 Guías prácticas: [desplegar un bot de IA multiplataforma en 5 minutos](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [conectar DeepSeek a WeChat, Discord y Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [ejecutar un Dify Agent en Discord, Telegram y Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) y [crear un chatbot con n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -88,7 +88,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Más opciones:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Manual](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Más opciones:** [Docker](https://link.langbot.app/en/docs/docker) · [Manual](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
@@ -149,7 +149,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | Pasarela | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Pasarela | ✅ |
[→ Ver todas las integraciones](https://langbot.app/docs/en/insight/features)
[→ Ver todas las integraciones](https://link.langbot.app/en/docs/features)
---
+6 -6
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Accueil</a>
<a href="https://langbot.app/docs/en/insight/features">Fonctionnalités</a>
<a href="https://langbot.app/docs/en/insight/guide">Documentation</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Fonctionnalités</a>
<a href="https://link.langbot.app/en/docs/guide">Documentation</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Marché des Plugins</a>
<a href="https://langbot.featurebase.app/roadmap">Feuille de Route</a>
@@ -48,7 +48,7 @@ LangBot est une **plateforme open-source de niveau production** pour créer des
- **Panneau de Gestion Web** — Configurez, gérez et surveillez vos bots via une interface navigateur intuitive. Aucune édition de YAML requise.
- **Architecture Multi-Pipeline** — Différents bots pour différents scénarios, avec surveillance complète et gestion des exceptions.
[→ En savoir plus sur toutes les fonctionnalités](https://langbot.app/docs/en/insight/features)
[→ En savoir plus sur toutes les fonctionnalités](https://link.langbot.app/en/docs/features)
📍 Guides pratiques : [déployer un bot IA multiplateforme en 5 minutes](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [connecter DeepSeek à WeChat, Discord et Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [exécuter un Dify Agent dans Discord, Telegram et Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) et [créer un chatbot avec n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -88,7 +88,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Plus d'options :** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Manuel](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Plus d'options :** [Docker](https://link.langbot.app/en/docs/docker) · [Manuel](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
@@ -149,7 +149,7 @@ docker compose --profile all up -d
| [ShengSuanYun](https://www.shengsuanyun.com/?from=CH_KYIPP758) | Plateforme GPU | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Passerelle | ✅ |
[→ Voir toutes les intégrations](https://langbot.app/docs/en/insight/features)
[→ Voir toutes les intégrations](https://link.langbot.app/en/docs/features)
---
+6 -6
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">ホーム</a>
<a href="https://langbot.app/docs/ja/insight/features">機能</a>
<a href="https://langbot.app/docs/ja/insight/guide">ドキュメント</a>
<a href="https://langbot.app/docs/ja/tags/readme">API</a>
<a href="https://link.langbot.app/ja/docs/features">機能</a>
<a href="https://link.langbot.app/ja/docs/guide">ドキュメント</a>
<a href="https://link.langbot.app/ja/docs/api">API</a>
<a href="https://space.langbot.app">プラグインマーケット</a> |
<a href="https://langbot.featurebase.app/roadmap">ロードマップ</a>
@@ -48,7 +48,7 @@ LangBot は、AI搭載のインスタントメッセージングボットを構
- **Web管理パネル** — 直感的なブラウザインターフェースからボットの設定、管理、監視が可能。YAML編集は不要。
- **マルチパイプラインアーキテクチャ** — 異なるシナリオに異なるボットを配置し、包括的な監視と例外処理を実現。
[→ すべての機能について詳しく見る](https://langbot.app/docs/ja/insight/features)
[→ すべての機能について詳しく見る](https://link.langbot.app/ja/docs/features)
📍 実践ガイド: [5分でマルチプラットフォームAIボットをデプロイ](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/)、[DeepSeekをWeChat・Discord・Telegramに接続](https://langbot.app/en/blog/connect-deepseek-to-wechat/)、[Dify AgentをDiscord・Telegram・Slackで動かす](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/)、[n8n連携チャットボットを構築](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/)。
@@ -88,7 +88,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**その他:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [手動デプロイ](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**その他:** [Docker](https://link.langbot.app/en/docs/docker) · [手動デプロイ](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
@@ -149,7 +149,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | ゲートウェイ | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | ゲートウェイ | ✅ |
[→ すべての統合を表示](https://langbot.app/docs/en/insight/features)
[→ すべての統合を表示](https://link.langbot.app/en/docs/features)
---
+6 -6
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">홈</a>
<a href="https://langbot.app/docs/en/insight/features">기능</a>
<a href="https://langbot.app/docs/en/insight/guide">문서</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">기능</a>
<a href="https://link.langbot.app/en/docs/guide">문서</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">플러그인 마켓</a>
<a href="https://langbot.featurebase.app/roadmap">로드맵</a>
@@ -48,7 +48,7 @@ LangBot은 AI 기반 인스턴트 메시징 봇을 구축하기 위한 **오픈
- **웹 관리 패널** — 직관적인 브라우저 인터페이스로 봇을 구성, 관리 및 모니터링. YAML 편집 불필요.
- **멀티 파이프라인 아키텍처** — 다양한 시나리오에 맞는 다양한 봇 구성, 종합 모니터링 및 예외 처리.
[→ 모든 기능 자세히 보기](https://langbot.app/docs/en/insight/features)
[→ 모든 기능 자세히 보기](https://link.langbot.app/en/docs/features)
📍 실전 가이드: [5분 만에 멀티 플랫폼 AI 봇 배포하기](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [DeepSeek를 WeChat, Discord, Telegram에 연결하기](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [Dify Agent를 Discord, Telegram, Slack에서 실행하기](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/), [n8n 기반 챗봇 만들기](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -88,7 +88,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**더 많은 옵션:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [수동 배포](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**더 많은 옵션:** [Docker](https://link.langbot.app/en/docs/docker) · [수동 배포](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
@@ -149,7 +149,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | 게이트웨이 | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | 게이트웨이 | ✅ |
[→ 모든 통합 보기](https://langbot.app/docs/en/insight/features)
[→ 모든 통합 보기](https://link.langbot.app/en/docs/features)
---
+6 -6
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Главная</a>
<a href="https://langbot.app/docs/en/insight/features">Возможности</a> |
<a href="https://langbot.app/docs/en/insight/guide">Документация</a> |
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Возможности</a> |
<a href="https://link.langbot.app/en/docs/guide">Документация</a> |
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Магазин плагинов</a>
<a href="https://langbot.featurebase.app/roadmap">Дорожная карта</a>
@@ -48,7 +48,7 @@ LangBot — это **платформа с открытым исходным к
- **Веб-панель управления** — Настраивайте, управляйте и мониторьте ваших ботов через интуитивный браузерный интерфейс. Ручное редактирование YAML не требуется.
- **Мультиконвейерная архитектура** — Разные боты для разных сценариев с комплексным мониторингом и обработкой исключений.
[→ Подробнее обо всех возможностях](https://langbot.app/docs/en/insight/features)
[→ Подробнее обо всех возможностях](https://link.langbot.app/en/docs/features)
📍 Практические руководства: [развернуть мультиплатформенного ИИ-бота за 5 минут](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [подключить DeepSeek к WeChat, Discord и Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [запустить Dify Agent в Discord, Telegram и Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) и [создать чат-бота на n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -88,7 +88,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Другие варианты:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Ручная установка](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Другие варианты:** [Docker](https://link.langbot.app/en/docs/docker) · [Ручная установка](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
@@ -149,7 +149,7 @@ docker compose --profile all up -d
| [ShengSuanYun](https://www.shengsuanyun.com/?from=CH_KYIPP758) | Платформа GPU | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Шлюз | ✅ |
[→ Смотреть все интеграции](https://langbot.app/docs/en/insight/features)
[→ Смотреть все интеграции](https://link.langbot.app/en/docs/features)
---
+6 -6
View File
@@ -21,9 +21,9 @@
[![star](https://gitcode.com/RockChinQ/LangBot/star/badge.svg)](https://gitcode.com/RockChinQ/LangBot)
<a href="https://langbot.app">官網</a>
<a href="https://langbot.app/docs/zh/insight/features">特性</a>
<a href="https://langbot.app/docs/zh/insight/guide">文件</a>
<a href="https://langbot.app/docs/zh/tags/readme">API</a>
<a href="https://link.langbot.app/zh/docs/features">特性</a>
<a href="https://link.langbot.app/zh/docs/guide">文件</a>
<a href="https://link.langbot.app/zh/docs/api">API</a>
<a href="https://space.langbot.app">外掛市場</a>
<a href="https://langbot.featurebase.app/roadmap">路線圖</a>
@@ -50,7 +50,7 @@ LangBot 是一個**開源的生產級平台**,用於建構 AI 驅動的即時
- **Web 管理面板** — 透過瀏覽器直觀地配置、管理和監控機器人,無需手動編輯設定檔。
- **多流水線架構** — 不同機器人用於不同場景,具備全面的監控和異常處理能力。
[→ 了解更多功能特性](https://langbot.app/docs/zh/insight/features)
[→ 了解更多功能特性](https://link.langbot.app/zh/docs/features)
📍 實踐指南:[5 分鐘部署多平台 AI 機器人](https://langbot.app/zh/blog/deploy-ai-bot-in-5-minutes/)、[將 DeepSeek 接入微信、企業微信與 Discord](https://langbot.app/zh/blog/connect-deepseek-to-wechat/)、[讓 Dify Agent 跑在 Discord、Telegram 和 Slack 上](https://langbot.app/zh/blog/dify-agent-discord-telegram-slack/),以及[用 n8n 建構多平台 AI 聊天機器人](https://langbot.app/zh/blog/n8n-multi-platform-ai-chatbot/)。
@@ -90,7 +90,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/zh-CN/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**更多方式:** [Docker](https://langbot.app/docs/zh/deploy/langbot/docker) · [手動部署](https://langbot.app/docs/zh/deploy/langbot/manual) · [寶塔面板](https://langbot.app/docs/zh/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/zh/deploy/langbot/kubernetes)
**更多方式:** [Docker](https://link.langbot.app/zh/docs/docker) · [手動部署](https://link.langbot.app/zh/docs/manual-deploy) · [寶塔面板](https://link.langbot.app/zh/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/zh/deploy/langbot/kubernetes)
---
@@ -165,7 +165,7 @@ docker compose --profile all up -d
|-----------|------|
| 阿里雲百煉 | [外掛](https://github.com/Thetail001/LangBot_BailianTextToImagePlugin) |
[→ 查看完整整合列表](https://langbot.app/docs/zh/insight/features)
[→ 查看完整整合列表](https://link.langbot.app/zh/docs/features)
---
+6 -6
View File
@@ -19,9 +19,9 @@
[![GitHub stars](https://img.shields.io/github/stars/langbot-app/LangBot?style=social)](https://github.com/langbot-app/LangBot/stargazers)
<a href="https://langbot.app">Trang chủ</a>
<a href="https://langbot.app/docs/en/insight/features">Tính năng</a>
<a href="https://langbot.app/docs/en/insight/guide">Tài liệu</a>
<a href="https://langbot.app/docs/en/tags/readme">API</a>
<a href="https://link.langbot.app/en/docs/features">Tính năng</a>
<a href="https://link.langbot.app/en/docs/guide">Tài liệu</a>
<a href="https://link.langbot.app/en/docs/api">API</a>
<a href="https://space.langbot.app">Chợ Plugin</a>
<a href="https://langbot.featurebase.app/roadmap">Lộ trình</a>
@@ -48,7 +48,7 @@ LangBot là một **nền tảng mã nguồn mở, cấp sản xuất** để x
- **Bảng quản lý Web** — Cấu hình, quản lý và giám sát bot thông qua giao diện trình duyệt trực quan. Không cần chỉnh sửa YAML.
- **Kiến trúc đa Pipeline** — Các bot khác nhau cho các kịch bản khác nhau, với giám sát toàn diện và xử lý ngoại lệ.
[→ Tìm hiểu thêm về tất cả tính năng](https://langbot.app/docs/en/insight/features)
[→ Tìm hiểu thêm về tất cả tính năng](https://link.langbot.app/en/docs/features)
📍 Hướng dẫn thực hành: [triển khai bot AI đa nền tảng trong 5 phút](https://langbot.app/en/blog/deploy-ai-bot-in-5-minutes/), [kết nối DeepSeek với WeChat, Discord và Telegram](https://langbot.app/en/blog/connect-deepseek-to-wechat/), [chạy Dify Agent trên Discord, Telegram và Slack](https://langbot.app/en/blog/dify-agent-discord-telegram-slack/) và [xây dựng chatbot với n8n](https://langbot.app/en/blog/n8n-multi-platform-ai-chatbot/).
@@ -88,7 +88,7 @@ docker compose --profile all up -d
[![Deploy on Zeabur](https://zeabur.com/button.svg)](https://zeabur.com/en-US/templates/ZKTBDH)
[![Deploy on Railway](https://railway.com/button.svg)](https://railway.app/template/yRrAyL?referralCode=vogKPF)
**Thêm tùy chọn:** [Docker](https://langbot.app/docs/en/deploy/langbot/docker) · [Thủ công](https://langbot.app/docs/en/deploy/langbot/manual) · [BTPanel](https://langbot.app/docs/en/deploy/langbot/one-click/bt) · [Kubernetes](https://langbot.app/docs/en/deploy/langbot/kubernetes)
**Thêm tùy chọn:** [Docker](https://link.langbot.app/en/docs/docker) · [Thủ công](https://link.langbot.app/en/docs/manual-deploy) · [BTPanel](https://link.langbot.app/en/docs/bt-panel) · [Kubernetes](https://docs.langbot.app/en/deploy/langbot/kubernetes)
---
@@ -149,7 +149,7 @@ docker compose --profile all up -d
| [302.AI](https://share.302ai.cn/SuTG99) | Cổng | ✅ |
| [Qiniu](https://www.qiniu.com/ai/agent) | Cổng | ✅ |
[→ Xem tất cả tích hợp](https://langbot.app/docs/en/insight/features)
[→ Xem tất cả tích hợp](https://link.langbot.app/en/docs/features)
---
+1 -1
View File
@@ -1,5 +1,5 @@
# Docker Compose configuration for LangBot
# For Kubernetes deployment, see kubernetes.yaml and the deployment guide at https://langbot.app/docs
# For Kubernetes deployment, see kubernetes.yaml and the deployment guide at https://docs.langbot.app
version: "3"
services:
+1 -1
View File
@@ -1,7 +1,7 @@
# Kubernetes Deployment for LangBot
# This file provides Kubernetes deployment manifests for LangBot based on docker-compose.yaml
#
# Full deployment guide (zh/en/ja): https://langbot.app/docs -> Installation -> Kubernetes
# Full deployment guide (zh/en/ja): https://docs.langbot.app -> Installation -> Kubernetes
#
# Usage:
# kubectl -n langbot create secret generic langbot-plugin-runtime-control \
-17
View File
@@ -88,23 +88,6 @@ Each endpoint accepts **either**:
1. **User Token** (via `Authorization: Bearer <user_jwt_token>`) - for web UI and authenticated users
2. **API Key** (via `X-API-Key` or `Authorization: Bearer <api_key>`) - for external services
### Inspecting API Key Identity
`GET /api/v1/system/context` validates an API key (user JWT not accepted) and returns its bound identity without requiring resource permissions:
```json
{
"code": 0,
"msg": "ok",
"data": {
"instance_uuid": "...",
"workspace_uuid": "...",
"api_key_id": "...",
"permissions": ["..."]
}
}
```
## Example: Model Management
### List All LLM Models
+2 -2
View File
@@ -218,8 +218,8 @@ metadata:
spec:
categories: [popular, global]
help_links:
zh: https://langbot.app/docs/zh/platforms/http-bot
en: https://langbot.app/docs/en/platforms/http-bot
zh: https://docs.langbot.app/zh/platforms/http-bot
en: https://docs.langbot.app/en/platforms/http-bot
config:
- { name: inbound_secret, type: string, required: true, default: "" }
- { name: callback_url, type: string, required: false, default: "" }
+1 -1
View File
@@ -243,7 +243,7 @@ For large datasets:
- SeekDB GitHub: https://github.com/oceanbase/seekdb
- pyseekdb SDK: https://github.com/oceanbase/pyseekdb
- OceanBase Documentation: https://oceanbase.ai
- LangBot Documentation: https://langbot.app/docs
- LangBot Documentation: https://docs.langbot.app
## License
+1 -1
View File
@@ -6,7 +6,7 @@ Minimal, dependency-light clients for the LangBot **HTTP Bot** platform adapter.
They show the whole loop: signing a request, pushing a message, and receiving
multi-part replies on a callback endpoint.
Full guide: [docs.langbot.app — HTTP Bot](https://langbot.app/docs/en/usage/platforms/http-bot).
Full guide: [docs.langbot.app — HTTP Bot](https://docs.langbot.app/en/usage/platforms/http-bot).
Machine-readable contract: [`docs/http-bot-openapi.json`](../../docs/http-bot-openapi.json).
## Files
+1 -1
View File
@@ -6,7 +6,7 @@
它们完整展示了整条链路:对请求签名、推送一条消息、在回调端点接收
1→M 的多段回复。
完整指南:[docs.langbot.app —— HTTP Bot](https://langbot.app/docs/zh/usage/platforms/http-bot)。
完整指南:[docs.langbot.app —— HTTP Bot](https://docs.langbot.app/zh/usage/platforms/http-bot)。
机器可读的接口契约:[`docs/http-bot-openapi.json`](../../docs/http-bot-openapi.json)。
## 文件清单
+1 -1
View File
@@ -6,7 +6,7 @@ A single self-contained HTML page that demos the LangBot **Page Bot**
(`web_page_bot`) embeddable chat widget — the one you drop onto any website with
a single `<script>` tag.
Full guide: [docs.langbot.app — Page Bot](https://langbot.app/docs/en/usage/platforms/webpage).
Full guide: [docs.langbot.app — Page Bot](https://docs.langbot.app/en/usage/platforms/webpage).
## Files
+1 -1
View File
@@ -6,7 +6,7 @@
(`web_page_bot`) 的可嵌入聊天组件 —— 也就是你用一行 `<script>` 标签就能放到任意
网站上的那个组件。
完整指南:[docs.langbot.app —— 页面机器人](https://langbot.app/docs/zh/usage/platforms/webpage)。
完整指南:[docs.langbot.app —— 页面机器人](https://docs.langbot.app/zh/usage/platforms/webpage)。
## 文件清单
+3 -3
View File
@@ -1,6 +1,6 @@
[project]
name = "langbot"
version = "4.10.10"
version = "4.10.8"
description = "Production-grade platform for building agentic IM bots"
readme = "README.md"
license-files = ["LICENSE"]
@@ -70,7 +70,7 @@ dependencies = [
"langchain-text-splitters>=1.1.2",
"chromadb>=1.0.0,<2.0.0",
"qdrant-client (>=1.15.1,<2.0.0)",
"langbot-plugin==0.5.7",
"langbot-plugin==0.5.5",
"asyncpg>=0.30.0",
"line-bot-sdk>=3.19.0",
"matrix-nio>=0.25.2",
@@ -114,7 +114,7 @@ seekdb = [
[project.urls]
Homepage = "https://langbot.app"
Documentation = "https://langbot.app/docs"
Documentation = "https://docs.langbot.app"
Repository = "https://github.com/langbot-app/LangBot"
[project.scripts]
+1 -2
View File
@@ -1349,8 +1349,7 @@
"local-agent",
"tools",
"e2b",
"nsjail",
"host"
"nsjail"
],
"automation": "",
"setup_automation": [],
+2 -8
View File
@@ -48,7 +48,7 @@ tools, skill add/edit, and stdio MCP are disabled. Set `box.enabled: false`
## Kubernetes
See `docker/kubernetes.yaml` and the deployment guide at
https://langbot.app/docs. `docker/deploy-k8s-test.sh` is a test helper.
https://docs.langbot.app. `docker/deploy-k8s-test.sh` is a test helper.
## config.yaml (generated at `data/config.yaml` on first run)
@@ -63,7 +63,7 @@ Key settings:
| `api.global_api_key` | **Global API key** for the HTTP API + MCP server. Non-empty = accepted with no login/DB record; no `lbk_` prefix required. Empty = disabled. Plaintext — trusted/internal only, serve over HTTPS. |
| `plugin.runtime_ws_url` | Standalone plugin runtime WS URL (e.g. `ws://langbot_plugin_runtime:5400/control/ws`) |
| `box.enabled` | Master switch for the Box sandbox runtime |
| `box.backend` | `local` (Docker/nsjail autopick) / `docker` / `nsjail` / `e2b` / explicit unsafe `host`; env override `BOX__BACKEND` |
| `box.backend` | `local` (Docker/nsjail autopick) / `docker` / `nsjail` / `e2b`; env override `BOX__BACKEND` |
| `box.runtime.endpoint` | External Box runtime URL (e.g. `ws://127.0.0.1:5410`); empty = local auto-managed |
Many keys have `ENV__SUBKEY` overrides (e.g. `BOX__BACKEND`, `BOX__ENABLED`).
@@ -75,10 +75,6 @@ Many keys have `ENV__SUBKEY` overrides (e.g. `BOX__BACKEND`, `BOX__ENABLED`).
with `--standalone-runtime`.
- Box has a parallel `--standalone-box` flag; the Docker box host is
`langbot_box:5410`.
- `box.backend: host` runs commands directly as the Box Runtime system user.
It is never auto-selected, provides no sandbox isolation, and is only for
trusted local development. A WebSocket-controlled host backend requires
`LANGBOT_BOX_CONTROL_TOKEN`; local stdio control is allowed.
## Global API key — enabling for agents/automation
@@ -97,7 +93,5 @@ login session. See `langbot-mcp-ops` for using it, and `docs/API_KEY_AUTH.md`.
- "No supported sandbox backend (Docker / nsjail / E2B)" with Docker running
usually means the user isn't in the `docker` group →
`sudo usermod -aG docker <user>` and restart in a new shell.
- Do not use `box.backend: host` as a production fallback. It cannot enforce
image, filesystem, network, PID, CPU, memory, or storage isolation.
- Box root host/container path mismatch breaks sandbox container creation.
- Don't commit a non-empty `api.global_api_key` to version control.
-2
View File
@@ -43,8 +43,6 @@ Two kinds of key are accepted:
Invalid, revoked, or expired keys get `401 Unauthorized`. A valid key whose
scopes do not authorize a tool gets `403 Forbidden`.
To inspect key identity and permissions, call `GET /api/v1/system/context` with the API key.
## Client configuration
```json
@@ -13,7 +13,6 @@ tags:
- tools
- e2b
- nsjail
- host
skills:
- langbot-env-setup
- langbot-testing
@@ -24,7 +23,7 @@ env:
- LANGBOT_LOCAL_AGENT_PIPELINE_NAME
preconditions:
- "LANGBOT_LOCAL_AGENT_PIPELINE_URL or LANGBOT_LOCAL_AGENT_PIPELINE_NAME points to the local-agent pipeline under test."
- "LangBot is started with the Box backend intended for this run, such as e2b, nsjail, or explicit host development mode."
- "LangBot is started with the sandbox backend intended for this run, such as e2b or nsjail."
- "The selected model route supports tool/function calling strongly enough to invoke sandbox tools."
steps:
- "Start LangBot with the target sandbox backend and confirm the Box status UI or LANGBOT_BACKEND_URL /api/v1/box/status reports the expected backend."
@@ -34,7 +33,7 @@ steps:
checks:
- "UI: Debug Chat final assistant response contains E2E_OK:<skill-name>."
- "Logs: The model called exec, register_skill, activate, then exec again from the activated skill path."
- "Logs: The selected backend name is the expected one, such as e2b, nsjail, or host."
- "Logs: The selected backend name is the expected one, such as e2b or nsjail."
- "Skill store: The registered package and activated writeback match references/sandbox-skill-authoring.md."
- "Box status: recent_error_count is 0 after the run."
evidence_required:
@@ -4,7 +4,7 @@
Verify that Local Agent can use sandbox tools to create, register, activate, and use a LangBot skill package through the same path a user would exercise in Debug Chat.
This flow applies to Docker, nsjail, E2B, and the explicit host development backend. Host runs commands directly as the Box Runtime user and must never be treated as sandbox-isolation coverage. API calls are useful diagnostics, but the primary pass/fail signal is the model-driven Debug Chat tool sequence.
This flow applies to Docker, nsjail, and E2B backends. API calls are useful diagnostics, but the primary pass/fail signal is the model-driven Debug Chat tool sequence.
## Preconditions
@@ -13,7 +13,6 @@ This flow applies to Docker, nsjail, E2B, and the explicit host development back
- `BOX_BACKEND=e2b` when validating E2B.
- `BOX_BACKEND=nsjail` when validating nsjail.
- `BOX_BACKEND=local` or `docker` when validating local container fallback.
- `BOX_BACKEND=host` only when validating explicit, trusted local direct execution.
3. Confirm `/api/v1/box/status` reports `available: true` and the expected backend name.
4. Confirm Debug Chat uses a model with function-calling ability.
5. Confirm backend logs say native sandbox tools are available.
@@ -72,7 +71,7 @@ Backend logs should show:
- `register_skill`
- `activate`
- a second `exec` whose workdir is `/workspace/.skills/<skill-name>`
- `backend=e2b`, `backend=nsjail`, `backend=host`, or the expected local backend
- `backend=e2b`, `backend=nsjail`, or the expected local backend
After the run, verify the skill store through the UI or API:
@@ -126,8 +125,6 @@ For E2B raw HTTP diagnostics, include a valid template id such as `base`; a miss
- Session metadata should keep LangBot logical paths such as `/workspace`; storing provider-internal paths can make later requests look incompatible.
- nsjail versions differ. Some expose only `--disable_clone_new*` flags and use `--bindmount` instead of `--rw_bind`.
- On WSL, cgroup v2 may exist but not be writable. The backend should warn and fall back to rlimits rather than fail the sandbox.
- The host backend does not honor sandbox image, network, rootfs, process, or
resource isolation. Use a disposable workspace and low-privilege account.
- If `ALL_PROXY` uses a SOCKS URL and `socksio` is not installed, some Python HTTP clients can fail during startup. Prefer consistent HTTP proxy variables unless SOCKS support is installed.
## Related Troubleshooting
@@ -3,7 +3,7 @@ title: "Native sandbox tools are unavailable even though a backend is configured
date: 2026-05-18
symptoms:
- "Backend logs show Native sandbox tools (exec/read/write/edit/glob/grep) are NOT available."
- "The Box runtime later reports that E2B, nsjail, Docker, or explicit host mode is configured."
- "The Box runtime later reports that E2B, nsjail, or Docker is configured."
- "Debug Chat does not expose exec, register_skill, or activate as usable tools."
patterns:
- "Native sandbox tools ... are NOT available"
@@ -19,7 +19,6 @@ fix_steps:
- "Ensure the Box runtime reselects a backend when get_backend_info is called and the cached backend is empty."
- "For E2B, verify the key without printing it and confirm any required template setting."
- "For nsjail, run nsjail --help and confirm the binary is on PATH for the LangBot process."
- "For trusted local development only, explicitly set box.backend=host; never use host as a production sandbox fallback."
verification: "Run sandbox-skill-authoring-e2e. Logs should show Native sandbox tools are available and /api/v1/box/status should report available=true with the expected backend."
related_cases:
- sandbox-skill-authoring-e2e
+1 -1
View File
@@ -16,7 +16,7 @@ asciiart = r"""
|___/
Open Source 开源地址: https://github.com/langbot-app/LangBot
📖 Documentation 文档地址: https://langbot.app/docs
📖 Documentation 文档地址: https://docs.langbot.app
"""
+11 -38
View File
@@ -1,14 +1,13 @@
from __future__ import annotations
import asyncio
import json
import os
import typing
from pathlib import Path
import httpx
import typing
import json
from .errors import DifyAPIError
from pathlib import Path
import os
_MAX_DIFY_RESPONSE_BYTES = 1024 * 1024
_MAX_DIFY_SSE_LINE_BYTES = 1024 * 1024
@@ -16,32 +15,6 @@ _MAX_DIFY_STREAM_BYTES = 16 * 1024 * 1024
_MAX_DIFY_UPLOAD_BYTES = 10 * 1024 * 1024
def _decode_sse_data(line: bytes) -> dict[str, typing.Any] | None:
data = line[5:].strip()
if not data or data == b'[DONE]':
return None
try:
payload = json.loads(data.decode('utf-8'))
except (json.JSONDecodeError, UnicodeDecodeError) as exc:
raise DifyAPIError('Dify SSE data line is not valid JSON') from exc
if not isinstance(payload, dict):
raise DifyAPIError('Dify SSE event is not a JSON object')
return payload
def _decode_upload_response(body: bytes) -> dict[str, typing.Any]:
try:
response = json.loads(body)
except (json.JSONDecodeError, UnicodeDecodeError) as exc:
raise DifyAPIError('Dify upload response is not valid JSON') from exc
if not isinstance(response, dict):
raise DifyAPIError('Dify upload response is not a JSON object')
payload = response.get('data', response)
if not isinstance(payload, dict) or not isinstance(payload.get('id'), str) or not payload['id']:
raise DifyAPIError('Dify upload response does not contain a valid file id')
return payload
async def _read_limited_response(
response: httpx.Response,
*,
@@ -83,16 +56,16 @@ async def _iter_sse_json(
line = raw_line.rstrip(b'\r').strip()
if not line or not line.startswith(b'data:'):
continue
payload = _decode_sse_data(line)
if payload is not None:
payload = json.loads(line[5:].decode('utf-8', errors='replace'))
if isinstance(payload, dict):
yield payload
if len(buffer) > _MAX_DIFY_SSE_LINE_BYTES:
raise DifyAPIError('Dify SSE event exceeds the runtime limit')
line = bytes(buffer).rstrip(b'\r').strip()
if line.startswith(b'data:'):
payload = _decode_sse_data(line)
if payload is not None:
payload = json.loads(line[5:].decode('utf-8', errors='replace'))
if isinstance(payload, dict):
yield payload
@@ -269,7 +242,7 @@ class AsyncDifyServiceClient:
file: httpx._types.FileTypes,
user: str,
timeout: float = 30.0,
) -> dict[str, typing.Any]:
) -> str:
# 处理 Path 对象
if isinstance(file, Path):
if not file.exists():
@@ -298,6 +271,6 @@ class AsyncDifyServiceClient:
timeout=timeout,
) as response:
body = await _read_limited_response(response)
if response.status_code not in (200, 201):
if response.status_code != 201:
raise DifyAPIError(f'{response.status_code} {body.decode(errors="replace")}')
return _decode_upload_response(body)
return json.loads(body)
+2 -3
View File
@@ -697,10 +697,9 @@ class DingTalkClient:
if not await self.check_access_token():
await self.get_access_token()
template_params = dict(card_param_map or {})
cardData: dict = {'cardParamMap': _stringify_card_param_map(card_param_map)}
if card_data_config is not None:
template_params['config'] = card_data_config
cardData: dict = {'cardParamMap': _stringify_card_param_map(template_params)}
cardData['config'] = json.dumps(card_data_config)
body: dict = {
'cardTemplateId': card_template_id,
@@ -936,13 +936,6 @@ class WecomBotWsClient:
'chat_type': message_data.get('type', 'single'),
}
self._prune_stream_state()
# Send an initial empty stream frame so the WeCom client
# shows its built-in loading spinner while the pipeline
# processes the message (e.g. RAG retrieval).
try:
await self.reply_stream(req_id, stream_id, '', finish=False)
except Exception:
await self.logger.warning(f'Failed to send initial stream frame: {traceback.format_exc()}')
message_data['stream_id'] = stream_id
message_data['req_id'] = req_id
@@ -295,34 +295,6 @@ class WecomCSClient:
raise Exception('Failed to send message')
return data
@_bounded_token_retry
async def send_image_msg(self, open_kfid: str, external_userid: str, msgid: str, media_id: str):
if not await self.check_access_token():
self.access_token = await self.get_access_token(self.secret)
url = f'{self.base_url}/kf/send_msg?access_token={self.access_token}'
payload = {
'touser': external_userid,
'open_kfid': open_kfid,
'msgid': msgid,
'msgtype': 'image',
'image': {
'media_id': media_id,
},
}
async with self._http_client_context() as client:
response = await client.post(url, json=payload)
data = await httpclient.parse_json_response(response)
if data['errcode'] == 40014 or data['errcode'] == 42001:
self.access_token = await self.get_access_token(self.secret)
return await self.send_image_msg(open_kfid, external_userid, msgid, media_id)
if data['errcode'] != 0:
await self.logger.error(f'发送图片失败:{data}')
raise Exception('Failed to send image message')
return data
async def handle_callback_request(self):
"""处理回调请求(独立端口模式,使用全局 request)。"""
return await self._handle_callback_internal(request)
@@ -218,7 +218,6 @@ class MonitoringRouterGroup(group.RouterGroup):
pipeline_ids = quart.request.args.getlist('pipelineId')
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
user_query = quart.request.args.get('userQuery')
is_active_str = quart.request.args.get('isActive')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
@@ -238,7 +237,6 @@ class MonitoringRouterGroup(group.RouterGroup):
pipeline_ids=pipeline_ids if pipeline_ids else None,
start_time=start_time,
end_time=end_time,
user_query=user_query,
is_active=is_active,
limit=limit,
offset=offset,
@@ -398,14 +396,7 @@ class MonitoringRouterGroup(group.RouterGroup):
@self.route('/sessions/<session_id>/analysis', methods=['GET'], permission=Permission.RESOURCE_VIEW)
async def get_session_analysis(session_id: str, request_context: RequestContext) -> str:
"""Get detailed analysis for a specific session"""
start_time = parse_iso_datetime(quart.request.args.get('startTime'))
end_time = parse_iso_datetime(quart.request.args.get('endTime'))
analysis = await self.ap.monitoring_service.get_session_analysis(
request_context,
session_id,
start_time=start_time,
end_time=end_time,
)
analysis = await self.ap.monitoring_service.get_session_analysis(request_context, session_id)
# Always return success with the analysis data
# The frontend will handle the 'found: false' case
@@ -15,17 +15,6 @@ from .....workspace.invitation_delivery import InvitationDeliveryService
@group.group_class('system', '/api/v1/system')
class SystemRouterGroup(group.RouterGroup):
async def initialize(self) -> None:
@self.route('/context', methods=['GET'], auth_type=group.AuthType.API_KEY)
async def _(request_context: RequestContext) -> str:
return self.success(
data={
'instance_uuid': request_context.instance_uuid,
'workspace_uuid': request_context.workspace_uuid,
'api_key_id': request_context.principal.api_key_uuid,
'permissions': sorted(request_context.workspace.permissions),
}
)
@self.route('/info', methods=['GET'], auth_type=group.AuthType.NONE)
async def _() -> str:
# Read wizard_status and wizard_progress from metadata table
@@ -186,9 +186,6 @@ class UserRouterGroup(group.RouterGroup):
json_data = await quart.request.json
code = json_data.get('code')
state = json_data.get('state')
redirect_uri = json_data.get('redirect_uri') or (
quart.request.url_root.rstrip('/') + '/auth/space/callback'
)
launch_assertion = json_data.get('launch_assertion')
workspace_uuid = json_data.get('workspace_uuid')
@@ -202,11 +199,8 @@ class UserRouterGroup(group.RouterGroup):
return self.fail(1, 'Missing authorization code')
if not state:
return self.fail(1, 'Missing state parameter')
if not str(code).startswith('v4_'):
return self.fail(1, 'Unsupported Space OAuth code contract')
try:
redirect_uri = self._validate_space_redirect_uri(str(redirect_uri), bind=False)
consumed_state = await self.ap.user_service.consume_space_oauth_state_details(state, 'login')
# Exchange code for tokens
launch_workspace_uuid = consumed_state.launch_workspace_uuid
@@ -224,36 +218,24 @@ class UserRouterGroup(group.RouterGroup):
code,
workspace_uuids,
workspace_created_ats,
redirect_uri=redirect_uri,
)
access_token = token_data.get('access_token')
refresh_token = token_data.get('refresh_token')
expires_in = token_data.get('expires_in', 0)
cloud_workspace_uuid = token_data.get('cloud_workspace_uuid')
if not access_token:
return self.fail(1, 'Failed to get access token from Space')
cloud_mode = getattr(getattr(self.ap, 'deployment', None), 'mode', 'oss') == 'cloud'
if cloud_mode and launch_workspace_uuid and launch_workspace_uuid != cloud_workspace_uuid:
return self.fail(1, 'Space OAuth Workspace binding mismatch')
target_workspace_uuid = launch_workspace_uuid or cloud_workspace_uuid
if cloud_mode:
if not target_workspace_uuid:
return self.fail(1, 'Space OAuth response is missing the Cloud Workspace binding')
await self.ap.directory_projection_service.reconcile_workspaces((target_workspace_uuid,))
# Authenticate only after the signed, exact Workspace delta has
# established the Account and membership runtime shadow rows.
# Authenticate and create/update local user
jwt_token, user_obj = await self.ap.user_service.authenticate_space_user(
access_token, refresh_token, expires_in
)
if target_workspace_uuid:
if launch_workspace_uuid:
try:
access = await self.ap.workspace_collaboration_service.resolve_account_workspace(
user_obj.uuid,
target_workspace_uuid,
launch_workspace_uuid,
)
except Exception:
self.ap.logger.warning('Rejected Space OAuth launch for unauthorized Workspace')
@@ -385,17 +367,12 @@ class UserRouterGroup(group.RouterGroup):
json_data = await quart.request.json
code = json_data.get('code')
state = json_data.get('state')
redirect_uri = json_data.get('redirect_uri') or (
quart.request.url_root.rstrip('/') + '/auth/space/callback?mode=bind'
)
if not code:
return self.http_status(400, -1, 'Missing authorization code')
if not state:
return self.http_status(400, -1, 'Missing state parameter')
if not str(code).startswith('v4_'):
return self.http_status(400, -1, 'Unsupported Space OAuth code contract')
try:
user_obj = await self.ap.user_service.consume_space_oauth_state(state, 'bind')
@@ -408,10 +385,7 @@ class UserRouterGroup(group.RouterGroup):
return self.http_status(400, -1, 'Only local accounts can bind to Space')
try:
redirect_uri = self._validate_space_redirect_uri(str(redirect_uri), bind=True)
updated_user = await self.ap.user_service.bind_space_account(
user_obj.user, code, redirect_uri=redirect_uri
)
updated_user = await self.ap.user_service.bind_space_account(user_obj.user, code)
jwt_token = await self.ap.user_service.generate_jwt_token(updated_user)
return self.success(
data={
@@ -454,10 +428,6 @@ class UserRouterGroup(group.RouterGroup):
}
)
projection_service = self.ap.directory_projection_service
if projection_service is None:
raise SpaceLaunchError('Cloud directory projection is unavailable')
await projection_service.reconcile_workspaces((launch['workspace_uuid'],))
account = await self.ap.user_service.get_user_by_uuid(launch['account_uuid'])
if account is None:
raise SpaceLaunchError('Launch Account is not projected into Core')
+1 -10
View File
@@ -137,16 +137,7 @@ class BotService:
bot = await self.get_bot(context, bot_data['uuid'], include_secret=True)
try:
await self.ap.platform_mgr.load_bot(context, bot)
except Exception:
# The bot row was already inserted above; without this rollback a
# failing adapter constructor (e.g. a missing optional credential
# key) would leave a permanently disabled orphan bot in the DB.
await self.ap.persistence_mgr.execute_async(
sqlalchemy.delete(persistence_bot.Bot).where(persistence_bot.Bot.uuid == bot_data['uuid'])
)
raise
await self.ap.platform_mgr.load_bot(context, bot)
return bot_data['uuid']
+4 -20
View File
@@ -1257,7 +1257,6 @@ class MonitoringService:
pipeline_ids: list[str] | None = None,
start_time: datetime.datetime | None = None,
end_time: datetime.datetime | None = None,
user_query: str | None = None,
is_active: bool | None = None,
limit: int = 100,
offset: int = 0,
@@ -1275,14 +1274,6 @@ class MonitoringService:
conditions.append(persistence_monitoring.MonitoringSession.start_time >= start_time)
if end_time:
conditions.append(persistence_monitoring.MonitoringSession.start_time <= end_time)
if user_query and user_query.strip():
user_pattern = f'%{user_query.strip()}%'
conditions.append(
sqlalchemy.or_(
persistence_monitoring.MonitoringSession.user_id.ilike(user_pattern),
persistence_monitoring.MonitoringSession.user_name.ilike(user_pattern),
)
)
if is_active is not None:
conditions.append(persistence_monitoring.MonitoringSession.is_active == is_active)
@@ -1374,8 +1365,6 @@ class MonitoringService:
self,
context: TenantContext,
session_id: str,
start_time: datetime.datetime | None = None,
end_time: datetime.datetime | None = None,
) -> dict:
"""Get bounded session details with full statistics computed in SQL."""
workspace_uuid = require_workspace_uuid(context)
@@ -1489,17 +1478,12 @@ class MonitoringService:
)
)
tool_stats = tool_stats_result.one()
tool_conditions = [
persistence_monitoring.MonitoringToolCall.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringToolCall.session_id == session_id,
]
if start_time is not None:
tool_conditions.append(persistence_monitoring.MonitoringToolCall.timestamp >= start_time)
if end_time is not None:
tool_conditions.append(persistence_monitoring.MonitoringToolCall.timestamp <= end_time)
tool_query = (
sqlalchemy.select(persistence_monitoring.MonitoringToolCall)
.where(*tool_conditions)
.where(
persistence_monitoring.MonitoringToolCall.workspace_uuid == workspace_uuid,
persistence_monitoring.MonitoringToolCall.session_id == session_id,
)
.order_by(persistence_monitoring.MonitoringToolCall.timestamp.asc())
.limit(detail_limit + 1)
)
+1 -4
View File
@@ -119,7 +119,7 @@ class SpaceService:
space_config = self._get_space_config()
authorize_url = space_config['oauth_authorize_url']
params = {'redirect_uri': redirect_uri, 'code_contract': 'redirect-v1'}
params = {'redirect_uri': redirect_uri}
if state:
params['state'] = state
return f'{authorize_url}?{urlencode(params)}'
@@ -129,8 +129,6 @@ class SpaceService:
code: str,
workspace_uuids: list[str] | None = None,
workspace_created_ats: dict[str, int] | None = None,
*,
redirect_uri: str = '',
) -> typing.Dict:
"""Exchange OAuth authorization code for tokens"""
from langbot.pkg.utils import constants
@@ -143,7 +141,6 @@ class SpaceService:
f'{space_url}/api/v1/accounts/oauth/token',
json={
'code': code,
'redirect_uri': redirect_uri,
'instance_id': constants.instance_id,
# Sending an explicit empty list tells new Space servers not to
# synthesize a legacy instance-derived Workspace binding.
+2 -3
View File
@@ -774,7 +774,7 @@ class UserService:
f'email:{normalized_email}',
)
async def bind_space_account(self, user_email: str, code: str, *, redirect_uri: str = '') -> user.User:
async def bind_space_account(self, user_email: str, code: str) -> user.User:
"""Bind Space account to existing local account"""
local_account = await self.get_user_by_email(user_email)
if local_account is None:
@@ -794,13 +794,12 @@ class UserService:
code,
[binding.workspace_uuid],
{binding.workspace_uuid: created_ts},
redirect_uri=redirect_uri,
)
else:
# Compatibility for early/bootstrap call sites that have not wired
# WorkspaceService yet; old Space servers still derive the legacy
# Workspace identity from instance_id when the field is omitted.
token_data = await self.ap.space_service.exchange_oauth_code(code, redirect_uri=redirect_uri)
token_data = await self.ap.space_service.exchange_oauth_code(code)
access_token = token_data.get('access_token')
refresh_token = token_data.get('refresh_token')
expires_in = token_data.get('expires_in', 0)
+3 -10
View File
@@ -455,9 +455,7 @@ class BoxService:
async def _require_validated_workspace_sandbox(self, execution_context: ExecutionContext) -> None:
if not self._available:
raise BoxError(
'Box runtime is not available. Configure an available Box backend before using Box features.'
)
raise BoxError('Box runtime is not available. Install and start Docker to use sandbox features.')
if self._cloud_managed:
if self._admission is None:
raise BoxAdmissionError('Cloud Box sandbox admission is unavailable')
@@ -567,9 +565,7 @@ class BoxService:
skip_host_mount_validation: bool = False,
) -> dict:
if not self._available:
raise BoxError(
'Box runtime is not available. Configure an available Box backend before using Box features.'
)
raise BoxError('Box runtime is not available. Install and start Docker to use sandbox features.')
execution_context = await self._validated_execution_context(self._query_execution_context(query))
spec_payload = self._managed_policy_payload(execution_context, spec_payload)
await self._require_validated_workspace_sandbox(execution_context)
@@ -2146,8 +2142,5 @@ class BoxService:
if backend_name:
payload['connector_error'] = f'Configured sandbox backend "{backend_name}" is unavailable'
else:
payload['connector_error'] = (
'No supported sandbox backend (Docker / nsjail / E2B) is available. '
'Trusted local development may explicitly select the unsafe host backend.'
)
payload['connector_error'] = 'No supported sandbox backend (Docker / nsjail / E2B) is available'
return payload
+2 -100
View File
@@ -125,21 +125,10 @@ class DirectoryProjectionService:
# The database cursor remains the shared projection high-water mark,
# while this cursor tracks what this process has actually observed.
self._consumer_cursor: int | None = None
self._sync_lock = asyncio.Lock()
async def initialize(self) -> None:
"""Block Cloud startup until one full signed snapshot is committed."""
async with self._sync_lock:
await self._refresh_snapshot()
async def refresh_snapshot(self) -> None:
"""Refresh from one full signed snapshot within the sync single-flight."""
async with self._sync_lock:
await self._refresh_snapshot()
async def _refresh_snapshot(self) -> None:
last_superseded: _DirectorySnapshotSuperseded | None = None
for _attempt in range(5):
snapshot = await self.provider.fetch_snapshot(self.instance_uuid)
@@ -170,84 +159,9 @@ class DirectoryProjectionService:
delay = min(max(delay * 2, self.sync_interval_seconds), self.max_staleness_seconds / 2)
async def sync_once(self) -> None:
async with self._sync_lock:
await self._sync_once()
async def reconcile_workspaces(self, workspace_uuids: Iterable[str]) -> None:
"""Synchronously project an exact Workspace set without moving the event cursor."""
requested = tuple(sorted({str(value).strip() for value in workspace_uuids if str(value).strip()}))
if not requested:
raise DirectoryProjectionUnavailableError('Targeted directory reconciliation requires a Workspace')
if len(requested) > self.event_limit:
raise DirectoryProjectionUnavailableError('Targeted directory reconciliation exceeds the batch limit')
async with self._sync_lock:
delta = await self.provider.fetch_workspaces(self.instance_uuid, requested)
await self._apply_targeted_delta(delta, requested)
async def _apply_targeted_delta(
self,
delta: DirectoryDelta,
requested_workspace_uuids: tuple[str, ...],
) -> None:
if not isinstance(delta, DirectoryDelta):
raise DirectoryProjectionUnavailableError('Directory provider returned an invalid delta')
workspace_count, membership_count = self._validate_batch_capacity(
delta.workspaces,
full_snapshot=False,
)
delta = DirectoryDelta.model_validate(delta.model_dump())
if delta.instance_uuid != self.instance_uuid:
raise DirectoryProjectionUnavailableError('Directory delta targets another LangBot instance')
requested = set(requested_workspace_uuids)
if set(delta.requested_workspace_uuids) != requested:
raise DirectoryProjectionUnavailableError('Directory delta does not match the requested Workspaces')
if {workspace.uuid for workspace in delta.workspaces} != requested:
raise DirectoryProjectionUnavailableError('Directory delta omitted a requested Workspace')
directory_uow = getattr(self.ap.persistence_mgr, 'directory_projection_uow', None)
if not callable(directory_uow):
raise DirectoryProjectionUnavailableError('Directory projection persistence scope is unavailable')
async with directory_uow(self.instance_uuid) as uow:
session = uow.session
state = await session.scalar(
sqlalchemy.select(DirectoryProjectionState)
.where(DirectoryProjectionState.instance_uuid == self.instance_uuid)
.with_for_update()
)
if state is None:
raise DirectoryProjectionUnavailableError('Directory projection is not initialized')
snapshot = DirectorySnapshot(
instance_uuid=self.instance_uuid,
cursor=state.cursor,
generated_at=delta.generated_at,
workspaces=delta.workspaces,
)
accounts_by_uuid = await self._apply_accounts(session, snapshot, preserve_existing=True)
await self._apply_workspaces(session, snapshot, accounts_by_uuid=accounts_by_uuid)
active_workspace_count = await self._enforce_active_workspace_capacity(session)
await session.flush()
await self._update_entitlement_workspace_activity(
snapshot.workspaces,
requested_workspace_uuids=requested,
)
self._publish_runtime_execution_projection(
snapshot.workspaces,
affected_workspace_uuids=requested,
)
self._request_model_catalog_sync()
self._record_batch_cardinality(
active_workspaces=active_workspace_count,
workspaces=workspace_count,
memberships=membership_count,
)
async def _sync_once(self) -> None:
cursor = self._consumer_cursor
if cursor is None:
await self._refresh_snapshot()
await self.initialize()
return
batch = await self.provider.fetch_events(
self.instance_uuid,
@@ -794,13 +708,7 @@ class DirectoryProjectionService:
for row in inbox_rows:
row.applied_at = now
async def _apply_accounts(
self,
session: Any,
snapshot: DirectorySnapshot,
*,
preserve_existing: bool = False,
) -> dict[str, User]:
async def _apply_accounts(self, session: Any, snapshot: DirectorySnapshot) -> dict[str, User]:
selected: dict[str, DirectoryMember] = {}
emails: dict[str, str] = {}
for workspace in snapshot.workspaces:
@@ -865,12 +773,6 @@ class DirectoryProjectionService:
continue
if account.source != AccountSource.CLOUD_PROJECTION.value:
raise DirectoryProjectionUnavailableError('Directory account UUID collides with a local Core account')
if preserve_existing:
# A targeted Workspace fetch has no independently monotonic
# Account revision. It may create a missing runtime shadow, but
# ordered event/snapshot projection remains the only updater of
# existing Account identity and status fields.
continue
if account.projection_revision > snapshot.cursor:
raise DirectoryProjectionUnavailableError('Directory account revision rolled back')
projected_account = self._account_projection(member)
+2 -2
View File
@@ -635,9 +635,9 @@ class Application:
frontend_path = paths.get_frontend_path()
if not os.path.exists(frontend_path):
self.logger.warning('WebUI 文件缺失,请根据文档部署:https://langbot.app/docs/zh')
self.logger.warning('WebUI 文件缺失,请根据文档部署:https://docs.langbot.app/zh')
self.logger.warning(
'WebUI files are missing, please deploy according to the documentation: https://langbot.app/docs/en'
'WebUI files are missing, please deploy according to the documentation: https://docs.langbot.app/en'
)
return
@@ -3,7 +3,6 @@
from __future__ import annotations
import asyncio
import contextlib
import dataclasses
import datetime
import json
@@ -83,7 +82,7 @@ def _verify_connection(connection: sqlite3.Connection, expected_revision: str) -
def _verify_file(path: pathlib.Path, expected_revision: str) -> None:
with contextlib.closing(_open_read_only(path)) as connection:
with _open_read_only(path) as connection:
_verify_connection(connection, expected_revision)
@@ -120,16 +119,12 @@ def _write_manifest(backup: SQLiteMigrationBackup, status: str, **extra: typing.
def _fsync_file(path: pathlib.Path, *, reopen_attempts: int = 20) -> None:
"""Sync a file, tolerating delayed visibility after replace on bind mounts.
Uses O_RDWR so os.fsync works on Windows (where _commit requires write
access to the file descriptor).
"""
"""Sync a file, tolerating delayed visibility after replace on bind mounts."""
descriptor: int | None = None
for attempt in range(reopen_attempts):
try:
descriptor = os.open(path, os.O_RDWR)
descriptor = os.open(path, os.O_RDONLY)
break
except FileNotFoundError:
if attempt + 1 >= reopen_attempts:
@@ -143,37 +138,13 @@ def _fsync_file(path: pathlib.Path, *, reopen_attempts: int = 20) -> None:
def _fsync_directory(path: pathlib.Path) -> None:
if os.name == 'nt':
# Windows cannot fsync directory handles opened through os.open.
return
descriptor = os.open(path, os.O_RDONLY | getattr(os, 'O_DIRECTORY', 0))
descriptor = os.open(path, os.O_RDONLY)
try:
os.fsync(descriptor)
finally:
os.close(descriptor)
def _remove_stale_temporary_files(
directory: pathlib.Path,
*,
prefix: str,
suffix: str,
) -> None:
"""Remove temporary files left by an interrupted backup or restore."""
for candidate in directory.iterdir():
if candidate.is_dir() or not candidate.name.startswith(prefix) or not candidate.name.endswith(suffix):
continue
try:
candidate.unlink()
except FileNotFoundError:
continue
except PermissionError:
# Another process may still own this file. Do not turn harmless
# cleanup into a migration failure; its unique name cannot collide.
continue
def _create_backup(
database_path: pathlib.Path,
source_revision: str,
@@ -182,11 +153,6 @@ def _create_backup(
backup_directory = database_path.parent / 'migration-backups'
backup_directory.mkdir(mode=0o700, parents=True, exist_ok=True)
os.chmod(backup_directory, 0o700)
_remove_stale_temporary_files(
backup_directory,
prefix=f'.{database_path.stem}-pre-',
suffix='.creating',
)
created_at = datetime.datetime.now(datetime.UTC).strftime('%Y-%m-%dT%H-%M-%S.%fZ')
stem = (
f'{database_path.stem}-pre-{_safe_label(target_revision)}-'
@@ -203,8 +169,11 @@ def _create_backup(
temporary_path = pathlib.Path(temporary_name)
try:
with (
contextlib.closing(_open_read_only(database_path)) as source,
contextlib.closing(sqlite3.connect(temporary_path, timeout=30)) as destination,
_open_read_only(database_path) as source,
sqlite3.connect(
temporary_path,
timeout=30,
) as destination,
):
source.execute('PRAGMA busy_timeout = 30000')
source.backup(destination)
@@ -252,11 +221,6 @@ async def create_verified_backup(
def _restore_backup(backup: SQLiteMigrationBackup) -> None:
_verify_file(backup.backup_path, backup.source_revision)
_remove_stale_temporary_files(
backup.database_path.parent,
prefix=f'.{backup.database_path.name}.',
suffix='.restoring',
)
descriptor, temporary_name = tempfile.mkstemp(
prefix=f'.{backup.database_path.name}.',
suffix='.restoring',
@@ -266,8 +230,11 @@ def _restore_backup(backup: SQLiteMigrationBackup) -> None:
temporary_path = pathlib.Path(temporary_name)
try:
with (
contextlib.closing(_open_read_only(backup.backup_path)) as source,
contextlib.closing(sqlite3.connect(temporary_path, timeout=30)) as destination,
_open_read_only(backup.backup_path) as source,
sqlite3.connect(
temporary_path,
timeout=30,
) as destination,
):
source.backup(destination)
destination.commit()
+1 -1
View File
@@ -209,7 +209,7 @@ _ALLOWED_SCOPED_BUILTIN_FUNCTION_TYPES = {
'now': sqlalchemy.sql.functions.now,
'sum': sqlalchemy.sql.functions.sum,
}
_ALLOWED_SCOPED_GENERIC_FUNCTIONS = frozenset({'date_trunc', 'length', 'nullif', 'strftime'})
_ALLOWED_SCOPED_GENERIC_FUNCTIONS = frozenset({'date_trunc', 'length', 'nullif'})
_ALLOWED_SCOPED_CUSTOM_OPERATORS = frozenset({'<=>'})
_ALLOWED_SCOPED_STATEMENT_TYPES = (
sqlalchemy.sql.dml.UpdateBase,
@@ -5,11 +5,6 @@ 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):
@@ -19,17 +14,12 @@ 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, current = await mask_patterns(
words,
found, message = await mask_patterns(
self.ap.sensitive_meta.data['words'],
message,
mask=mask,
mask_word=mask_word,
max_pattern_count=_MAX_SENSITIVE_WORD_PATTERNS,
mask=self.ap.sensitive_meta.data['mask'],
mask_word=self.ap.sensitive_meta.data['mask_word'],
)
except SafeRegexError as exc:
return entities.FilterResult(
@@ -41,7 +31,7 @@ class BanWordFilter(filter_model.ContentFilter):
return entities.FilterResult(
level=entities.ResultLevel.MASKED if found else entities.ResultLevel.PASS,
replacement=current,
replacement=message,
user_notice='消息中存在不合适的内容, 请修改' if found else '',
console_notice='',
)
@@ -15,9 +15,9 @@ spec:
categories:
- protocol
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/qq/aiocqhttp/napcat
en: https://langbot.app/docs/en/usage/platforms/qq/aiocqhttp/napcat
ja: https://langbot.app/docs/ja/usage/platforms/qq/aiocqhttp/napcat
zh: https://link.langbot.app/zh/platforms/aiocqhttp
en: https://link.langbot.app/en/platforms/aiocqhttp
ja: https://link.langbot.app/ja/platforms/aiocqhttp
config:
- name: host
label:
@@ -15,9 +15,9 @@ spec:
categories:
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/dingtalk
en: https://langbot.app/docs/en/usage/platforms/dingtalk
ja: https://langbot.app/docs/ja/usage/platforms/dingtalk
zh: https://link.langbot.app/zh/platforms/dingtalk
en: https://link.langbot.app/en/platforms/dingtalk
ja: https://link.langbot.app/ja/platforms/dingtalk
config:
- name: one-click-create
label:
@@ -24,9 +24,9 @@ spec:
- popular
- global
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/discord
en: https://langbot.app/docs/en/usage/platforms/discord
ja: https://langbot.app/docs/ja/usage/platforms/discord
zh: https://link.langbot.app/zh/platforms/discord
en: https://link.langbot.app/en/platforms/discord
ja: https://link.langbot.app/ja/platforms/discord
config:
- name: client_id
label:
@@ -18,9 +18,9 @@ spec:
- popular
- global
help_links:
zh: https://langbot.app/docs/zh/platforms/http-bot
en: https://langbot.app/docs/en/platforms/http-bot
ja: https://langbot.app/docs/ja/platforms/http-bot
zh: https://docs.langbot.app/zh/platforms/http-bot
en: https://docs.langbot.app/en/platforms/http-bot
ja: https://docs.langbot.app/ja/platforms/http-bot
config:
- name: webhook_url
label:
+3 -3
View File
@@ -15,9 +15,9 @@ spec:
categories:
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/kook
en: https://langbot.app/docs/en/usage/platforms/kook
ja: https://langbot.app/docs/ja/usage/platforms/kook
zh: https://link.langbot.app/zh/platforms/kook
en: https://link.langbot.app/en/platforms/kook
ja: https://link.langbot.app/ja/platforms/kook
config:
- name: token
label:
+4 -32
View File
@@ -160,29 +160,6 @@ def _lark_should_update_stream_element(
return not resume_from and not form_data and (msg_seq % 8 == 0 or is_final)
def _lark_final_layout_texts(
*,
resume_from: bool,
text_message: str,
pre_pause_cached: str | None,
resume_cached: str,
) -> tuple[str, str]:
"""Return (main_text, resume_placeholder_text) for the final card update.
Non-resume round: the full reply belongs in the main streaming element
only also rendering the resume placeholder duplicates the reply, since
both hold the same accumulated text. Resume round (Dify HITL): keep the
pre-pause text in the main element and the resumed text in the
placeholder, as they are distinct segments.
"""
if resume_from:
# An empty pre-pause cache is valid (Dify paused before emitting any
# text); only a missing entry (None) falls back to the full text.
main_text = text_message if pre_pause_cached is None else pre_pause_cached
return main_text, resume_cached
return text_message, ''
def _lark_display_input_value(field: dict, value: typing.Any) -> str:
field_type = _dify_field_type(field)
if field_type == 'file':
@@ -2381,21 +2358,16 @@ class LarkAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
self.card_form_input_defs[card_id] = _lark_form_input_defs(form_data)
self.card_form_inputs[card_id] = dict(form_data.get('inputs') or {})
else:
# Normal finish: remove buttons/notice and finalize the card.
main_text, resume_text = _lark_final_layout_texts(
resume_from=resume_from,
text_message=text_message,
pre_pause_cached=self.card_pre_pause_text.get(card_id),
resume_cached=resume_cached,
)
# Normal finish: keep pre-pause + resume content visible,
# remove buttons/notice, drop the resume placeholder.
await self._update_card_layout(
card_id=card_id,
message_source=message_source,
text_message=main_text,
text_message=pre_pause,
sequence=final_seq,
form_data=None,
notice_text=selected_notice if resume_from else '',
resume_placeholder_text=resume_text,
resume_placeholder_text=resume_cached,
)
self._drop_card_state(card_id)
self.card_id_dict.pop(message_id, None)
+3 -3
View File
@@ -19,9 +19,9 @@ spec:
- china
- global
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/lark
en: https://langbot.app/docs/en/usage/platforms/lark
ja: https://langbot.app/docs/ja/usage/platforms/lark
zh: https://link.langbot.app/zh/platforms/lark
en: https://link.langbot.app/en/platforms/lark
ja: https://link.langbot.app/ja/platforms/lark
config:
- name: domain
label:
+8 -51
View File
@@ -25,7 +25,6 @@ from linebot.v3.webhooks import (
ImageMessageContent,
VideoMessageContent,
AudioMessageContent,
UserMentionee,
)
# from linebot import WebhookParser
@@ -59,19 +58,15 @@ class LINEMessageConverter(abstract_platform_adapter.AbstractMessageConverter):
return content_list
def __init__(self, bot_account_id: str = ''):
self.bot_account_id = bot_account_id
async def target2yiri(self, message, bot_client) -> platform_message.MessageChain:
@staticmethod
async def target2yiri(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.extend(
self._build_text_components(message.message.text, getattr(message.message, 'mention', None))
)
lb_msg_list.append(platform_message.Plain(text=message.message.text))
elif isinstance(message.message, AudioMessageContent):
pass
elif isinstance(message.message, VideoMessageContent):
@@ -91,55 +86,17 @@ 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
async def target2yiri(self, event, bot_client) -> platform_events.Event:
message_chain = await self.message_converter.target2yiri(event, bot_client)
@staticmethod
async def target2yiri(event, bot_client) -> platform_events.Event:
message_chain = await LINEMessageConverter.target2yiri(event, bot_client)
if event.source.type == 'user':
return platform_events.FriendMessage(
@@ -212,8 +169,8 @@ class LINEAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
listeners={},
card_id_dict={},
seq=1,
event_converter=LINEEventConverter(bot_account_id),
message_converter=LINEMessageConverter(bot_account_id),
event_converter=LINEEventConverter(),
message_converter=LINEMessageConverter(),
line_webhook=line_webhook,
parser=parser,
configuration=configuration,
+3 -3
View File
@@ -22,9 +22,9 @@ spec:
categories:
- global
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/line
en: https://langbot.app/docs/en/usage/platforms/line
ja: https://langbot.app/docs/ja/usage/platforms/line
zh: https://link.langbot.app/zh/platforms/line
en: https://link.langbot.app/en/platforms/line
ja: https://link.langbot.app/ja/platforms/line
config:
- name: webhook_url
label:
@@ -15,9 +15,9 @@ spec:
categories:
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/wxoa
en: https://langbot.app/docs/en/usage/platforms/wxoa
ja: https://langbot.app/docs/ja/usage/platforms/wxoa
zh: https://link.langbot.app/zh/platforms/officialaccount
en: https://link.langbot.app/en/platforms/officialaccount
ja: https://link.langbot.app/ja/platforms/officialaccount
config:
- name: webhook_url
label:
@@ -16,9 +16,9 @@ spec:
- popular
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/wechat/weixin
en: https://langbot.app/docs/en/usage/platforms/readme
ja: https://langbot.app/docs/ja/usage/platforms/readme
zh: https://link.langbot.app/zh/platforms/openclaw_weixin
en: https://link.langbot.app/en/platforms/openclaw_weixin
ja: https://link.langbot.app/ja/platforms/openclaw_weixin
config:
- name: base_url
label:
@@ -205,7 +205,7 @@ class QQOfficialAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter
bot = QQOfficialClient(
app_id=config['appid'],
secret=config['secret'],
token=config.get('token', ''),
token=config['token'],
logger=logger,
unified_mode=enable_webhook,
)
@@ -15,9 +15,9 @@ spec:
categories:
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/qq/official_webhook
en: https://langbot.app/docs/en/usage/platforms/qq/official_webhook
ja: https://langbot.app/docs/ja/usage/platforms/qq/official_webhook
zh: https://link.langbot.app/zh/platforms/qqofficial
en: https://link.langbot.app/en/platforms/qqofficial
ja: https://link.langbot.app/ja/platforms/qqofficial
config:
- name: __system.outbound_ips
label:
+3 -3
View File
@@ -21,9 +21,9 @@ spec:
categories:
- protocol
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/readme
en: https://langbot.app/docs/en/usage/platforms/readme
ja: https://langbot.app/docs/ja/usage/platforms/readme
zh: https://link.langbot.app/zh/platforms/satori
en: https://link.langbot.app/en/platforms/satori
ja: https://link.langbot.app/ja/platforms/satori
config:
- name: platform
label:
+3 -3
View File
@@ -24,9 +24,9 @@ spec:
- popular
- global
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/slack
en: https://langbot.app/docs/en/usage/platforms/slack
ja: https://langbot.app/docs/ja/usage/platforms/slack
zh: https://link.langbot.app/zh/platforms/slack
en: https://link.langbot.app/en/platforms/slack
ja: https://link.langbot.app/ja/platforms/slack
config:
- name: webhook_url
label:
@@ -24,9 +24,9 @@ spec:
- popular
- global
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/telegram
en: https://langbot.app/docs/en/usage/platforms/telegram
ja: https://langbot.app/docs/ja/usage/platforms/telegram
zh: https://link.langbot.app/zh/platforms/telegram
en: https://link.langbot.app/en/platforms/telegram
ja: https://link.langbot.app/ja/platforms/telegram
config:
- name: token
label:
@@ -15,9 +15,9 @@ spec:
categories:
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/wechat/wechatpad
en: https://langbot.app/docs/en/usage/platforms/readme
ja: https://langbot.app/docs/ja/usage/platforms/readme
zh: https://link.langbot.app/zh/platforms/wechatpad
en: https://link.langbot.app/en/platforms/wechatpad
ja: https://link.langbot.app/ja/platforms/wechatpad
config:
- name: wechatpad_url
label:
+3 -3
View File
@@ -274,11 +274,11 @@ class WecomAdapter(abstract_platform_adapter.AbstractMessagePlatformAdapter):
if content['type'] == 'text':
await self.bot.send_private_msg(user_id, agent_id, content['content'])
if content['type'] == 'image':
await self.bot.send_image(user_id, agent_id, content['media_id'])
await self.bot.send_image(user_id, agent_id, content['media'])
if content['type'] == 'voice':
await self.bot.send_voice(user_id, agent_id, content['media_id'])
await self.bot.send_voice(user_id, agent_id, content['media'])
if content['type'] == 'file':
await self.bot.send_file(user_id, agent_id, content['media_id'])
await self.bot.send_file(user_id, agent_id, content['media'])
def register_listener(
self,
+3 -3
View File
@@ -16,9 +16,9 @@ spec:
- popular
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/wecom/wecom
en: https://langbot.app/docs/en/usage/platforms/wecom/wecom
ja: https://langbot.app/docs/ja/usage/platforms/wecom/wecom
zh: https://link.langbot.app/zh/platforms/wecom
en: https://link.langbot.app/en/platforms/wecom
ja: https://link.langbot.app/ja/platforms/wecom
config:
- name: webhook_url
label:
@@ -15,9 +15,9 @@ spec:
categories:
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/wecom/wecombot
en: https://langbot.app/docs/en/usage/platforms/wecom/wecombot
ja: https://langbot.app/docs/ja/usage/platforms/wecom/wecombot
zh: https://link.langbot.app/zh/platforms/wecombot
en: https://link.langbot.app/en/platforms/wecombot
ja: https://link.langbot.app/ja/platforms/wecombot
config:
- name: one-click-create
label:
+3 -10
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'{event.receiver_id}|u{event.user_id}',
id=f'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'{event.receiver_id}|u{event.user_id}',
id=f'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'{uuid.uuid4().hex}'
msgid = f'langbot_{uuid.uuid4().hex}'
if content['type'] == 'text':
await self.bot.send_text_msg(
open_kfid=open_kfid,
@@ -205,13 +205,6 @@ 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"""
@@ -15,9 +15,9 @@ spec:
categories:
- china
help_links:
zh: https://langbot.app/docs/zh/usage/platforms/wecom/wecomcs
en: https://langbot.app/docs/en/usage/platforms/wecom/wecomcs
ja: https://langbot.app/docs/ja/usage/platforms/wecom/wecomcs
zh: https://link.langbot.app/zh/platforms/wecomcs
en: https://link.langbot.app/en/platforms/wecomcs
ja: https://link.langbot.app/ja/platforms/wecomcs
config:
- name: webhook_url
label:
+2 -7
View File
@@ -1913,14 +1913,9 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
return plugins
async def get_plugin_info(self, author: str, plugin_name: str) -> dict[str, Any] | None:
async def get_plugin_info(self, author: str, plugin_name: str) -> dict[str, Any]:
runtime_handler = self._runtime_handler()
try:
binding = await self._target_binding(author, plugin_name)
except ValueError as exc:
if str(exc) == f'Plugin {author}/{plugin_name} is not installed in this Workspace':
return None
raise
binding = await self._target_binding(author, plugin_name)
with runtime_handler.installation_scope(binding):
return await runtime_handler.get_plugin_info(author, plugin_name)
@@ -573,7 +573,7 @@ class LiteLLMRequester(requester.ProviderAPIRequester):
levels = ['provider_default', 'disabled', 'enabled']
elif family == 'doubao':
levels = ['provider_default', 'disabled', 'low', 'medium', 'high']
elif family in ('ollama', 'ollama_chat'):
elif family == 'ollama':
levels = ['provider_default']
levels.append('disabled')
if normalized_name.startswith('gpt-oss') or '/gpt-oss' in normalized_name:
@@ -747,24 +747,9 @@ class LiteLLMRequester(requester.ProviderAPIRequester):
converted_parts = []
for part in content:
if isinstance(part, dict) and part.get('type') == '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
part['image_url'] = {'url': part['image_base64']}
part['type'] = 'image_url'
del part['image_base64']
# 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
@@ -1345,14 +1330,7 @@ class LiteLLMRequester(requester.ProviderAPIRequester):
extra_args: dict[str, typing.Any] = {},
) -> tuple[list[list[float]], dict]:
"""Invoke embedding and return vectors with usage info."""
# litellm's embedding routing has no "ollama_chat" branch (that provider
# exists only for /api/chat completions) — embeddings still go through
# the plain "ollama" provider. Requesters configured for ollama_chat
# (to get native tool-calling on the chat path) must fall back to
# "ollama" here specifically, or embedding calls raise "Unmapped LLM
# provider for this endpoint".
embedding_provider = 'ollama' if self._get_custom_llm_provider() == 'ollama_chat' else None
model_name = self._build_litellm_model_name(model.model_entity.name, embedding_provider)
model_name = self._build_litellm_model_name(model.model_entity.name)
api_key = model.provider.token_mgr.get_token()
args = {
@@ -1548,12 +1526,6 @@ class LiteLLMRequester(requester.ProviderAPIRequester):
event_hooks=httpclient.httpx_response_limit_hooks(),
) as client:
response = await client.get(models_url, headers=headers)
if response.status_code == 404 and not base_url.rstrip('/').endswith('/v1'):
# Some OpenAI-compatible servers (notably a bare Ollama host,
# e.g. http://host:11434) expose the model list under /v1/models
# rather than /models. Providers whose configured base_url
# already ends in /v1 keep their original (working) URL.
response = await client.get(f'{base_url}/v1/models', headers=headers)
response.raise_for_status()
payload = await httpclient.parse_json_response(response)
@@ -7,7 +7,7 @@ metadata:
zh_Hans: Ollama
icon: ollama.svg
spec:
litellm_provider: ollama_chat
litellm_provider: ollama
config:
- name: base_url
label:
@@ -619,9 +619,7 @@ class LocalAgentRunner(runner.RequestRunner):
and len(func_ret) > 0
and isinstance(func_ret[0], provider_message.ContentElement)
):
# 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)
tool_content = func_ret
else:
tool_content = json.dumps(func_ret, ensure_ascii=False)
+1 -13
View File
@@ -39,9 +39,6 @@ class N8nServiceAPIRunner(runner.RequestRunner):
# 获取输出键名,默认为response
self.output_key = self.pipeline_config['ai']['n8n-service-api'].get('output-key', 'response')
self.response_handling = self.pipeline_config['ai']['n8n-service-api'].get('response-handling', 'reply')
if self.response_handling not in {'reply', 'ignore'}:
raise ValueError(f'Invalid n8n response-handling: {self.response_handling}')
# 获取认证类型,默认为none
self.auth_type = self.pipeline_config['ai']['n8n-service-api'].get('auth-type', 'none')
@@ -265,11 +262,7 @@ class N8nServiceAPIRunner(runner.RequestRunner):
async with session.post(
self.webhook_url, json=payload, headers=headers, auth=auth, timeout=self.timeout
) as response:
if self.response_handling == 'ignore':
status_ok = 200 <= response.status < 300
else:
status_ok = response.status == 200
if not status_ok:
if response.status != 200:
error_text = (
await httpclient.read_limited(
response,
@@ -279,11 +272,6 @@ class N8nServiceAPIRunner(runner.RequestRunner):
self.ap.logger.error(f'n8n webhook call failed: {response.status}, {error_text}')
raise Exception(f'n8n webhook call failed: {response.status}, {error_text}')
if self.response_handling == 'ignore':
response.release()
self.ap.logger.debug('n8n async webhook accepted; response body ignored')
return
async for chunk in self._process_response(response):
if is_stream:
yield chunk
@@ -222,7 +222,6 @@ class NativeToolLoader(loader.ToolLoader):
self.ap.logger.warning(
'Native sandbox tools (exec/read/write/edit/glob/grep) are NOT available. '
'No sandbox backend (Docker/nsjail/E2B) is ready. '
'Trusted local development may explicitly select box.backend=host. '
'The LLM will not have access to code execution or file operation tools.'
)
@@ -42,8 +42,7 @@ class SkillToolLoader(loader.ToolLoader):
else:
self.ap.logger.info(
'Skill tools (activate/register_skill) are NOT available. '
'No sandbox backend (Docker/nsjail/E2B) is ready. '
'Trusted local development may explicitly select box.backend=host.'
'No sandbox backend (Docker/nsjail/E2B) is ready.'
)
async def _check_sandbox_available(self) -> bool:
+4 -13
View File
@@ -27,16 +27,10 @@ class SafeRegexTimeoutError(SafeRegexError):
"""Raised when the regex engine exhausts the operation CPU budget."""
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')
def _validate_patterns(patterns: Sequence[str]) -> tuple[str, ...]:
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')
@@ -121,9 +115,8 @@ def _mask_patterns_sync(
mask: str,
mask_word: str,
timeout_seconds: float,
max_pattern_count: int,
) -> tuple[bool, str]:
normalized_patterns = _validate_patterns(patterns, max_pattern_count=max_pattern_count)
normalized_patterns = _validate_patterns(patterns)
_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')
@@ -169,7 +162,6 @@ 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."""
@@ -182,5 +174,4 @@ async def mask_patterns(
mask=mask,
mask_word=mask_word,
timeout_seconds=timeout_seconds,
max_pattern_count=max_pattern_count,
)
+1 -1
View File
@@ -83,7 +83,7 @@ class VersionManager:
try:
if await self.is_new_version_available():
return (
'New version available. Update guide: https://langbot.app/docs/en/deploy/update',
'New version available. Update guide: https://link.langbot.app/en/docs/update',
logging.INFO,
)
except Exception as e:
@@ -269,7 +269,7 @@ class InvitationDeliveryService:
<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://langbot.app/docs/langbot-logo.png" alt="LangBot" width="34" height="34" style="display:inline-block;width:34px;height:34px;border:0;vertical-align:middle;">
<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>
+1 -4
View File
@@ -331,10 +331,7 @@ box:
# skill tool, skill add/edit, and stdio-mode MCP servers. Skills can still
# be listed read-only and http/sse MCP servers continue to work.
enabled: true
# 'host' runs commands directly as the Box Runtime user without sandbox
# isolation. It is never auto-selected and is only for trusted local
# development. Can be written via BOX__BACKEND.
backend: 'local' # 'local' (Docker/nsjail), 'docker', 'nsjail', 'e2b', or explicit unsafe 'host'.
backend: 'local' # 'local' (Docker/nsjail), 'docker', 'nsjail', or 'e2b'. Can be written via BOX__BACKEND.
runtime:
# LANGBOT_BOX_CONTROL_TOKEN is optional for OSS external WebSocket
# runtimes. To protect an exposed endpoint, set the same strong secret
@@ -80,8 +80,7 @@
"header-name": "",
"header-value": "",
"timeout": 120,
"output-key": "response",
"response-handling": "reply"
"output-key": "response"
},
"langflow-api": {
"base-url": "http://localhost:7860",
+3 -4
View File
@@ -642,10 +642,9 @@
.replace(/\s+/g, " ")
.trim();
if (
prevContent &&
(prevContent === content ||
prevContent.indexOf(content) >= 0 ||
content.indexOf(prevContent) >= 0)
prevContent === content ||
prevContent.indexOf(content) >= 0 ||
content.indexOf(prevContent) >= 0
)
return;
}
@@ -325,7 +325,7 @@ stages:
zh_Hans: API 密钥
type: string
required: true
default: ''
default: 'your-api-key'
- name: n8n-service-api
label:
en_US: n8n Workflow API
@@ -475,25 +475,6 @@ stages:
type: string
required: false
default: 'response'
- name: response-handling
label:
en_US: Webhook Response Handling
zh_Hans: Webhook 响应处理方式
description:
en_US: Choose whether LangBot forwards the n8n webhook response to the chat user. Ignore mode requires the n8n Webhook node to use Respond Immediately.
zh_Hans: 选择是否将 n8n Webhook 响应转发给聊天用户。忽略模式要求 n8n Webhook 节点使用“立即响应”。
type: select
required: false
default: 'reply'
options:
- name: reply
label:
en_US: Forward as chat reply
zh_Hans: 转发为聊天回复
- name: ignore
label:
en_US: Ignore response body (asynchronous workflow)
zh_Hans: 忽略响应正文(异步工作流)
- name: coze-api
label:
en_US: coze API
+2 -24
View File
@@ -242,22 +242,6 @@ class TestMonitoringSessionsEndpoint:
assert response.status_code == 200
@pytest.mark.asyncio
async def test_get_sessions_forwards_user_search_and_page_window(self, quart_test_client, fake_monitoring_app):
fake_monitoring_app.monitoring_service.get_sessions.reset_mock()
response = await quart_test_client.get(
'/api/v1/monitoring/sessions?botId=bot-1&userQuery=alice&limit=20&offset=40',
headers={'Authorization': 'Bearer test_token'},
)
assert response.status_code == 200
kwargs = fake_monitoring_app.monitoring_service.get_sessions.await_args.kwargs
assert kwargs['bot_ids'] == ['bot-1']
assert kwargs['user_query'] == 'alice'
assert kwargs['limit'] == 20
assert kwargs['offset'] == 40
@pytest.mark.usefixtures('mock_circular_import_chain')
class TestMonitoringErrorsEndpoint:
@@ -294,19 +278,13 @@ class TestMonitoringDetailsEndpoints:
"""Tests for detail endpoints."""
@pytest.mark.asyncio
async def test_get_session_analysis(self, quart_test_client, fake_monitoring_app):
async def test_get_session_analysis(self, quart_test_client):
"""GET /api/v1/monitoring/sessions/{id}/analysis."""
response = await quart_test_client.get(
'/api/v1/monitoring/sessions/sess-1/analysis'
'?startTime=2026-08-31T16%3A00%3A00.000Z'
'&endTime=2026-09-01T15%3A59%3A59.999Z',
headers={'Authorization': 'Bearer test_token'},
'/api/v1/monitoring/sessions/sess-1/analysis', headers={'Authorization': 'Bearer test_token'}
)
assert response.status_code == 200
kwargs = fake_monitoring_app.monitoring_service.get_session_analysis.await_args.kwargs
assert kwargs['start_time'].isoformat() == '2026-08-31T16:00:00'
assert kwargs['end_time'].isoformat() == '2026-09-01T15:59:59.999000'
@pytest.mark.asyncio
async def test_get_message_details(self, quart_test_client):
+8 -198
View File
@@ -27,8 +27,7 @@ async def space_oauth_api():
execution=SimpleNamespace(instance_uuid='instance-a', placement_generation=1),
)
application = Mock()
application.deployment = SimpleNamespace(multi_workspace_enabled=False, mode='oss')
application.directory_projection_service = None
application.deployment = SimpleNamespace(multi_workspace_enabled=False)
application.persistence_mgr = None
application.user_service.get_authenticated_account = AsyncMock(return_value=account)
application.user_service.issue_space_oauth_state = AsyncMock(
@@ -126,26 +125,6 @@ async def test_cloud_launch_state_is_server_issued_and_workspace_bound(space_oau
)
@pytest.mark.asyncio
async def test_cloud_login_entry_uses_normal_stateful_oauth(space_oauth_api):
application, client = space_oauth_api
application.deployment.mode = 'cloud'
response = await client.get(
'/api/v1/user/space/authorize-url',
query_string={
'redirect_uri': 'http://localhost/auth/space/callback',
'cloud_entry': '1',
},
headers={'Origin': 'http://localhost'},
)
assert response.status_code == 200
authorize_url = (await response.get_json())['data']['authorize_url']
assert authorize_url.startswith('https://space.example/authorize?state=')
application.user_service.issue_space_oauth_state.assert_awaited_once_with('login')
@pytest.mark.asyncio
async def test_public_login_rejects_caller_supplied_state(space_oauth_api):
application, client = space_oauth_api
@@ -270,14 +249,10 @@ async def test_server_side_webhook_origin_supports_bundled_ui(space_oauth_api):
async def test_login_callback_requires_and_consumes_server_state(space_oauth_api):
application, client = space_oauth_api
missing = await client.post('/api/v1/user/space/callback', json={'code': 'v4_oauth-code'})
missing = await client.post('/api/v1/user/space/callback', json={'code': 'oauth-code'})
response = await client.post(
'/api/v1/user/space/callback',
json={
'code': 'v4_oauth-code',
'state': 'opaque-login-state',
'redirect_uri': 'https://oss.example/auth/space/callback',
},
json={'code': 'oauth-code', 'state': 'opaque-login-state'},
)
assert (await missing.get_json())['code'] == 1
@@ -285,146 +260,12 @@ async def test_login_callback_requires_and_consumes_server_state(space_oauth_api
assert (await response.get_json())['data']['token'] == 'space-login-token'
application.user_service.consume_space_oauth_state_details.assert_awaited_once_with('opaque-login-state', 'login')
application.space_service.exchange_oauth_code.assert_awaited_once_with(
'v4_oauth-code',
'oauth-code',
[WORKSPACE_UUID],
{WORKSPACE_UUID: int(WORKSPACE_CREATED_AT.timestamp())},
redirect_uri='https://oss.example/auth/space/callback',
)
@pytest.mark.asyncio
async def test_login_callback_rejects_downgraded_legacy_code(space_oauth_api):
application, client = space_oauth_api
response = await client.post(
'/api/v1/user/space/callback',
json={'code': 'v2_legacy-code', 'state': 'opaque-login-state'},
)
payload = await response.get_json()
assert response.status_code == 200
assert payload['code'] == 1
assert 'code contract' in payload['msg']
application.space_service.exchange_oauth_code.assert_not_awaited()
@pytest.mark.asyncio
async def test_cloud_login_callback_reconciles_authorized_workspace_before_local_authentication(space_oauth_api):
application, client = space_oauth_api
application.deployment.mode = 'cloud'
calls: list[str] = []
application.directory_projection_service = SimpleNamespace(
reconcile_workspaces=AsyncMock(side_effect=lambda _workspace_uuids: calls.append('reconcile'))
)
application.space_service.exchange_oauth_code.return_value = {
'access_token': 'space-access-token',
'refresh_token': 'space-refresh-token',
'expires_in': 3600,
'cloud_workspace_uuid': WORKSPACE_UUID,
}
authenticated_account = application.user_service.authenticate_space_user.return_value[1]
async def authenticate(*_args):
calls.append('authenticate')
return 'space-login-token', authenticated_account
application.user_service.authenticate_space_user.side_effect = authenticate
response = await client.post(
'/api/v1/user/space/callback',
json={'code': 'v4_oauth-code', 'state': 'opaque-login-state'},
)
assert response.status_code == 200
assert (await response.get_json())['data']['workspace_uuid'] == WORKSPACE_UUID
assert calls == ['reconcile', 'authenticate']
application.directory_projection_service.reconcile_workspaces.assert_awaited_once_with((WORKSPACE_UUID,))
@pytest.mark.asyncio
async def test_cloud_login_callback_fails_closed_without_workspace_binding(space_oauth_api):
application, client = space_oauth_api
application.deployment.mode = 'cloud'
application.directory_projection_service = SimpleNamespace(reconcile_workspaces=AsyncMock())
response = await client.post(
'/api/v1/user/space/callback',
json={'code': 'v4_oauth-code', 'state': 'opaque-login-state'},
)
payload = await response.get_json()
assert response.status_code == 200
assert payload['code'] == 1
assert 'Cloud Workspace binding' in payload['msg']
application.directory_projection_service.reconcile_workspaces.assert_not_awaited()
application.user_service.authenticate_space_user.assert_not_awaited()
@pytest.mark.asyncio
async def test_cloud_login_callback_requires_code_binding_for_launch_state(space_oauth_api):
application, client = space_oauth_api
application.deployment.mode = 'cloud'
application.directory_projection_service = SimpleNamespace(reconcile_workspaces=AsyncMock())
application.user_service.consume_space_oauth_state_details.return_value = SimpleNamespace(
launch_workspace_uuid=WORKSPACE_UUID
)
response = await client.post(
'/api/v1/user/space/callback',
json={'code': 'v4_oauth-code', 'state': 'opaque-login-state'},
)
payload = await response.get_json()
assert response.status_code == 200
assert payload['code'] == 1
assert 'Workspace binding' in payload['msg']
application.directory_projection_service.reconcile_workspaces.assert_not_awaited()
application.user_service.authenticate_space_user.assert_not_awaited()
@pytest.mark.asyncio
async def test_cloud_login_callback_rejects_conflicting_state_and_code_workspace_bindings(space_oauth_api):
application, client = space_oauth_api
application.deployment.mode = 'cloud'
application.directory_projection_service = SimpleNamespace(reconcile_workspaces=AsyncMock())
application.user_service.consume_space_oauth_state_details.return_value = SimpleNamespace(
launch_workspace_uuid=WORKSPACE_UUID
)
application.space_service.exchange_oauth_code.return_value = {
'access_token': 'space-access-token',
'refresh_token': 'space-refresh-token',
'expires_in': 3600,
'cloud_workspace_uuid': 'workspace-from-another-flow',
}
response = await client.post(
'/api/v1/user/space/callback',
json={'code': 'v4_oauth-code', 'state': 'opaque-login-state'},
)
payload = await response.get_json()
assert response.status_code == 200
assert payload['code'] == 1
assert 'Workspace binding' in payload['msg']
application.directory_projection_service.reconcile_workspaces.assert_not_awaited()
application.user_service.authenticate_space_user.assert_not_awaited()
@pytest.mark.asyncio
async def test_oss_login_callback_does_not_request_cloud_reconciliation(space_oauth_api):
application, client = space_oauth_api
application.directory_projection_service = SimpleNamespace(reconcile_workspaces=AsyncMock())
response = await client.post(
'/api/v1/user/space/callback',
json={'code': 'v4_oauth-code', 'state': 'opaque-login-state'},
)
assert response.status_code == 200
application.directory_projection_service.reconcile_workspaces.assert_not_awaited()
@pytest.mark.asyncio
async def test_login_callback_launch_state_selects_asserted_workspace(space_oauth_api):
application, client = space_oauth_api
@@ -435,7 +276,7 @@ async def test_login_callback_launch_state_selects_asserted_workspace(space_oaut
response = await client.post(
'/api/v1/user/space/callback',
json={'code': 'v4_oauth-code', 'state': 'opaque-login-state'},
json={'code': 'oauth-code', 'state': 'opaque-login-state'},
)
assert response.status_code == 200
@@ -534,22 +375,18 @@ async def test_bind_callback_uses_opaque_state_and_never_treats_it_as_jwt(space_
rejected = await client.post(
'/api/v1/user/bind-space',
json={'code': 'v4_attacker-code', 'state': 'jwt.must-not-be-used'},
json={'code': 'attacker-code', 'state': 'jwt.must-not-be-used'},
)
response = await client.post(
'/api/v1/user/bind-space',
json={'code': 'v4_oauth-code', 'state': 'opaque-bind-state'},
json={'code': 'oauth-code', 'state': 'opaque-bind-state'},
)
assert rejected.status_code == 401
assert response.status_code == 200
assert (await response.get_json())['data']['token'] == 'rotated-account-token'
application.user_service.verify_jwt_token.assert_not_awaited()
application.user_service.bind_space_account.assert_awaited_once_with(
'owner@example.com',
'v4_oauth-code',
redirect_uri='http://localhost/auth/space/callback?mode=bind',
)
application.user_service.bind_space_account.assert_awaited_once_with('owner@example.com', 'oauth-code')
@pytest.mark.asyncio
@@ -557,7 +394,6 @@ async def test_direct_launch_assertion_does_not_consume_normal_oauth_state(space
application, client = space_oauth_api
application.user_service.consume_space_oauth_state.reset_mock()
application.space_service.exchange_oauth_code.reset_mock()
application.directory_projection_service = SimpleNamespace(reconcile_workspaces=AsyncMock())
response = await client.post(
'/api/v1/user/space/callback',
@@ -578,29 +414,3 @@ async def test_direct_launch_assertion_does_not_consume_normal_oauth_state(space
)
application.user_service.consume_space_oauth_state.assert_not_awaited()
application.space_service.exchange_oauth_code.assert_not_awaited()
@pytest.mark.asyncio
async def test_direct_launch_reconciles_exact_workspace_before_resolving_access(space_oauth_api):
application, client = space_oauth_api
projected_account = SimpleNamespace(
uuid='account-a',
user='owner@example.com',
account_type='space',
status='active',
)
application.user_service.get_user_by_uuid = AsyncMock(return_value=projected_account)
application.directory_projection_service = SimpleNamespace(reconcile_workspaces=AsyncMock())
response = await client.post(
'/api/v1/user/space/callback',
json={
'workspace_uuid': WORKSPACE_UUID,
'launch_assertion': 'signed-launch-token',
},
)
assert response.status_code == 200
assert (await response.get_json())['data']['workspace_uuid'] == WORKSPACE_UUID
application.directory_projection_service.reconcile_workspaces.assert_awaited_once_with((WORKSPACE_UUID,))
application.user_service.get_user_by_uuid.assert_awaited_once_with('account-a')
-64
View File
@@ -440,70 +440,6 @@ async def test_api_key_secret_is_one_time_and_viewer_cannot_manage_keys(workspac
assert (await forbidden.get_json())['code'] == 'permission_denied'
async def test_api_key_context_returns_bound_identity_without_workspace_permission(workspace_api):
application, client, _, owner_token = workspace_api
current_response = await client.get('/api/v1/workspaces/current', headers=_auth(owner_token))
workspace_uuid = (await current_response.get_json())['data']['workspace']['uuid']
create_response = await client.post(
'/api/v1/apikeys',
headers=_auth(owner_token, workspace_uuid),
json={'name': 'Context probe', 'scopes': []},
)
assert create_response.status_code == 200
created = (await create_response.get_json())['data']['key']
missing_auth = await client.get('/api/v1/system/context')
assert missing_auth.status_code == 401
invalid_auth = await client.get(
'/api/v1/system/context',
headers={'X-API-Key': 'lbk_invalid'},
)
assert invalid_auth.status_code == 401
response = await client.get(
'/api/v1/system/context',
headers={
'X-API-Key': created['key'],
'X-Workspace-Id': 'caller-selected-workspace-must-be-ignored',
},
)
assert response.status_code == 200
assert (await response.get_json())['data'] == {
'instance_uuid': application.workspace_service.instance_uuid,
'workspace_uuid': workspace_uuid,
'api_key_id': created['uuid'],
'permissions': [],
}
bearer_response = await client.get(
'/api/v1/system/context',
headers={'Authorization': f'Bearer {created["key"]}'},
)
assert bearer_response.status_code == 200
assert (await bearer_response.get_json())['data']['api_key_id'] == created['uuid']
jwt_response = await client.get(
'/api/v1/system/context',
headers={'Authorization': f'Bearer {owner_token}'},
)
assert jwt_response.status_code == 401
revoke_response = await client.delete(
f'/api/v1/apikeys/{created["id"]}',
headers=_auth(owner_token, workspace_uuid),
)
assert revoke_response.status_code == 200
revoked_response = await client.get(
'/api/v1/system/context',
headers={'X-API-Key': created['key']},
)
assert revoked_response.status_code == 401
async def test_cloud_projection_is_selected_explicitly_and_collaboration_runs_in_core(
workspace_api,
):
@@ -115,7 +115,6 @@ class _CapacityPluginRuntimeHandler:
def __init__(self) -> None:
self.bindings: dict[str, typing.Any] = {}
self.reconciled: tuple[typing.Any, ...] = ()
self.reconcile_timeout: float | None = None
def register_installation_binding(
self,
@@ -133,14 +132,8 @@ class _CapacityPluginRuntimeHandler:
def unregister_installation_binding(self, binding) -> None:
self.bindings.pop(binding.installation_uuid, None)
async def reconcile_plugin_installations(
self,
desired_states,
*,
timeout: float | None = None,
) -> dict:
async def reconcile_plugin_installations(self, desired_states) -> dict:
self.reconciled = tuple(desired_states)
self.reconcile_timeout = timeout
return {
'applied': [],
'removed': [],
@@ -1041,7 +1034,6 @@ class TestPostgreSQLTenantRuntime:
assert not mcp_loader._hosted_mcp_tasks
assert len(plugin_handler.reconciled) == workspace_count
assert len(plugin_handler.bindings) == workspace_count
assert plugin_handler.reconcile_timeout == 300.0
assert all(count == workspace_count for count in statement_counts.values()), statement_counts
if max_elapsed is not None:
assert elapsed <= max_elapsed
@@ -39,35 +39,6 @@ def _assert_verified_backup(payload: dict) -> None:
assert connection.execute('SELECT version_num FROM alembic_version').fetchone()[0] == payload['source_revision']
def _temporary_sqlite_files(root: pathlib.Path) -> list[pathlib.Path]:
return [*root.rglob('*.creating'), *root.rglob('*.restoring')]
async def test_backup_removes_stale_temporary_file_from_interrupted_run(tmp_path):
database_path = tmp_path / 'legacy-stale-backup.db'
engine = create_async_engine(f'sqlite+aiosqlite:///{database_path}')
try:
await create_legacy_resource_schema(engine, instance_uuid='stale-backup')
await alembic_runner.run_alembic_stamp(engine, '0008_mcp_resource_prefs')
backup_directory = tmp_path / 'migration-backups'
backup_directory.mkdir()
stale_path = backup_directory / '.legacy-stale-backup-pre-0009-old.creating'
unrelated_path = backup_directory / '.another-database-pre-0009-old.creating'
stale_path.write_bytes(b'interrupted backup')
unrelated_path.write_bytes(b'unrelated backup')
await sqlite_migration_backup.create_verified_backup(
engine,
source_revision='0008_mcp_resource_prefs',
target_revision='0009_workspace_tenancy',
)
assert not stale_path.exists()
assert unrelated_path.read_bytes() == b'unrelated backup'
finally:
await engine.dispose()
async def test_tenancy_migrations_retain_verified_boundary_backups(tmp_path):
database_path = tmp_path / 'legacy-with-backups.db'
engine = create_async_engine(f'sqlite+aiosqlite:///{database_path}')
@@ -88,7 +59,6 @@ async def test_tenancy_migrations_retain_verified_boundary_backups(tmp_path):
}
for payload in payloads:
_assert_verified_backup(payload)
assert _temporary_sqlite_files(tmp_path) == []
finally:
await engine.dispose()
@@ -130,7 +100,6 @@ async def test_failed_tenancy_migration_restores_backup_and_revision(
assert restored[0]['status'] == 'restored_after_failure'
assert restored[0]['source_revision'] == '0009_workspace_tenancy'
_assert_verified_backup(restored[0])
assert _temporary_sqlite_files(tmp_path) == []
monkeypatch.setattr(alembic_runner, 'run_alembic_upgrade', real_upgrade)
await _manager(engine)._run_alembic_migrations()
@@ -139,41 +108,6 @@ async def test_failed_tenancy_migration_restores_backup_and_revision(
await engine.dispose()
async def test_restore_publish_failure_preserves_current_database(tmp_path, monkeypatch):
database_path = tmp_path / 'restore-publish-failure.db'
engine = create_async_engine(f'sqlite+aiosqlite:///{database_path}')
try:
await create_legacy_resource_schema(engine, instance_uuid='restore-publish-failure')
await alembic_runner.run_alembic_stamp(engine, '0008_mcp_resource_prefs')
backup = await sqlite_migration_backup.create_verified_backup(
engine,
source_revision='0008_mcp_resource_prefs',
target_revision='0009_workspace_tenancy',
)
stale_restore_path = tmp_path / f'.{database_path.name}.interrupted.restoring'
stale_restore_path.write_bytes(b'interrupted restore')
async with engine.begin() as connection:
await connection.execute(sa.text("UPDATE alembic_version SET version_num = 'failed-revision'"))
await engine.dispose()
database_before_restore = database_path.read_bytes()
real_replace = os.replace
def fail_restore_publish(source, destination):
if pathlib.Path(destination) == database_path:
raise OSError('simulated atomic publish failure')
return real_replace(source, destination)
monkeypatch.setattr(sqlite_migration_backup.os, 'replace', fail_restore_publish)
with pytest.raises(OSError, match='atomic publish failure'):
await sqlite_migration_backup.restore_verified_backup(engine, backup)
assert database_path.read_bytes() == database_before_restore
assert _temporary_sqlite_files(tmp_path) == []
finally:
await engine.dispose()
async def test_backup_retries_transient_reopen_failure_after_replace(tmp_path, monkeypatch):
database_path = tmp_path / 'legacy-bind-mount.db'
engine = create_async_engine(f'sqlite+aiosqlite:///{database_path}')
@@ -12,7 +12,6 @@ import pytest
from unittest.mock import AsyncMock, MagicMock, Mock, patch
from types import SimpleNamespace
import json
import sqlalchemy
import uuid
from langbot.pkg.api.http.service.bot import BotService
@@ -450,58 +449,10 @@ class TestBotServiceCreateBot:
insert_statement = ap.persistence_mgr.execute_async.await_args_list[1].args[0]
insert_values = insert_statement.compile().params
assert insert_values['workspace_uuid'] == WORKSPACE_UUID
assert insert_values['use_pipeline_uuid'] == 'default-pipeline-uuid'
assert insert_values['use_pipeline_name'] == 'Default Pipeline'
assert bot_uuid is not None # Verify UUID was returned
async def test_create_bot_rolls_back_insert_when_load_bot_fails(self):
"""Deletes the inserted row when the adapter fails to load.
Regression: a failing adapter constructor (e.g. KeyError on a missing
optional credential key) used to leave a permanently disabled orphan
bot in the DB the insert was already committed and the HTTP layer
surfaced a 500 without any cleanup.
"""
# Setup
ap = SimpleNamespace()
ap.persistence_mgr = SimpleNamespace()
ap.instance_config = SimpleNamespace()
ap.instance_config.data = {'system': {'limitation': {'max_bots': -1}}}
ap.platform_mgr = SimpleNamespace()
ap.platform_mgr.load_bot = AsyncMock(side_effect=KeyError('token'))
pipeline_result = Mock()
pipeline_result.first = Mock(return_value=None)
bot_result = Mock()
bot_result.first = Mock(return_value=_create_mock_bot())
executed_statements = []
async def mock_execute(query):
executed_statements.append(query)
if len(executed_statements) <= 2:
return pipeline_result # 1: limitation bots query, 2: pipeline query
if len(executed_statements) == 3:
return Mock() # insert
return bot_result # get_bot after insert
ap.persistence_mgr.execute_async = AsyncMock(side_effect=mock_execute)
ap.persistence_mgr.serialize_model = Mock(return_value={'uuid': 'new-uuid', 'name': 'New Bot'})
service = BotService(ap)
# Execute & Verify: the adapter error propagates
with pytest.raises(KeyError, match='token'):
await service.create_bot(
WORKSPACE_UUID, {'name': 'New Bot', 'adapter': 'telegram', 'adapter_config': {}}
)
# And the inserted row is rolled back via a DELETE on the new uuid
# (no limitation query runs because max_bots=-1)
assert len(executed_statements) == 4 # pipeline select, insert, bot select, delete
delete_statement = executed_statements[-1]
assert isinstance(delete_statement, sqlalchemy.sql.dml.Delete)
compiled = delete_statement.compile()
assert compiled.params['uuid_1'] is not None
class TestBotServiceUpdateBot:
"""Tests for update_bot method."""
@@ -138,39 +138,6 @@ async def test_same_session_and_resource_ids_do_not_collide(service):
assert (await service.get_message_details(context_a, message_b))['found'] is False
async def test_session_search_matches_user_id_or_name_within_workspace(service):
context_a = _context(WORKSPACE_A)
context_b = _context(WORKSPACE_B)
fixtures = [
(context_a, 'session-id-match', 'customer-42', 'Alice'),
(context_a, 'session-name-match', 'customer-99', 'Bob Alice Cooper'),
(context_a, 'session-no-match', 'customer-7', 'Bob'),
(context_b, 'session-other-workspace', 'customer-42', 'Alice'),
]
for context, session_id, user_id, user_name in fixtures:
await service.record_session_start(
context,
session_id=session_id,
bot_id='same-bot',
bot_name='Same Bot',
pipeline_id='same-pipeline',
pipeline_name='Same Pipeline',
user_id=user_id,
user_name=user_name,
)
by_id, id_total = await service.get_sessions(context_a, user_query='customer-42')
by_name, name_total = await service.get_sessions(context_a, user_query='alice')
assert id_total == 1
assert [session['session_id'] for session in by_id] == ['session-id-match']
assert name_total == 2
assert {session['session_id'] for session in by_name} == {
'session-id-match',
'session-name-match',
}
async def test_tool_call_inherits_context_from_connection_message_row(service):
context = _context(WORKSPACE_A)
message_id = await _record_message(service, context, 'tool context')
@@ -95,9 +95,7 @@ class TestSpaceServiceGetOAuthAuthorizeUrl:
result = service.get_oauth_authorize_url('http://localhost/callback')
# Verify
query = parse_qs(urlsplit(result).query)
assert query['redirect_uri'] == ['http://localhost/callback']
assert query['code_contract'] == ['redirect-v1']
assert parse_qs(urlsplit(result).query)['redirect_uri'] == ['http://localhost/callback']
assert 'https://space.langbot.app/auth/authorize' in result
def test_get_oauth_authorize_url_with_state(self):
@@ -580,14 +578,12 @@ class TestSpaceServiceExchangeOAuthCode:
'auth_code',
['workspace-1'],
{'workspace-1': 1_700_000_000},
redirect_uri='https://oss.example/auth/space/callback',
)
# Verify
assert result['access_token'] == 'new_access_token'
assert mock_session_obj.post.call_args.kwargs['json'] == {
'code': 'auth_code',
'redirect_uri': 'https://oss.example/auth/space/callback',
'instance_id': constants.instance_id,
'workspace_uuids': ['workspace-1'],
'workspace_created_ats': {'workspace-1': 1_700_000_000},
@@ -850,7 +846,10 @@ class TestSpaceServiceGetModelSelection:
if response_shape == 'models-envelope':
data = {'models': models}
elif response_shape == 'availability-wrapper':
data = [{'model': model, 'latency_ms': index + 10, 'http_code': 200} for index, model in enumerate(models)]
data = [
{'model': model, 'latency_ms': index + 10, 'http_code': 200}
for index, model in enumerate(models)
]
else:
data = models
payload = {'code': 0, 'data': data}
@@ -1,10 +1,9 @@
from __future__ import annotations
import asyncio
import datetime
import logging
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock
from unittest.mock import Mock
import pytest
import sqlalchemy
@@ -215,88 +214,6 @@ async def test_directory_delta_requests_model_catalog_sync_after_commit(projecti
request_sync.assert_called_once_with()
async def test_targeted_reconciliation_projects_new_workspace_without_advancing_event_cursor(projection_context):
application, session_factory = projection_context
provider = _Provider(
[_snapshot(7, workspaces=[])],
deltas=[_delta(workspaces=[_workspace(revision=8, name='JIT Workspace')])],
)
service = DirectoryProjectionService(application, provider, INSTANCE_UUID)
await service.initialize()
await service.reconcile_workspaces((WORKSPACE_UUID,))
async with session_factory() as session:
account = await session.scalar(sqlalchemy.select(User).where(User.uuid == ACCOUNT_UUID))
workspace = await session.get(Workspace, WORKSPACE_UUID)
membership = await session.scalar(
sqlalchemy.select(WorkspaceMembership).where(
WorkspaceMembership.workspace_uuid == WORKSPACE_UUID,
WorkspaceMembership.account_uuid == ACCOUNT_UUID,
)
)
state = await session.get(DirectoryProjectionState, INSTANCE_UUID)
assert account is not None
assert workspace is not None and workspace.name == 'JIT Workspace'
assert membership is not None and membership.status == 'active'
assert state is not None and state.cursor == 7
assert provider.delta_calls == 1
assert provider.after_cursors == []
async def test_targeted_reconciliation_preserves_existing_account_until_ordered_event_projection(projection_context):
application, session_factory = projection_context
targeted_workspace = _workspace(revision=8, name='Renamed Workspace').model_copy(
update={
'members': [
_member(revision=8).model_copy(update={'display_name': 'Changed Account Name'})
]
}
)
provider = _Provider(
[_snapshot(7)],
deltas=[_delta(workspaces=[targeted_workspace])],
)
service = DirectoryProjectionService(application, provider, INSTANCE_UUID)
await service.initialize()
await service.reconcile_workspaces((WORKSPACE_UUID,))
async with session_factory() as session:
account = await session.scalar(sqlalchemy.select(User).where(User.uuid == ACCOUNT_UUID))
workspace = await session.get(Workspace, WORKSPACE_UUID)
state = await session.get(DirectoryProjectionState, INSTANCE_UUID)
assert account is not None and account.user == 'Workspace Owner'
assert account.projection_revision == 7
assert workspace is not None and workspace.name == 'Renamed Workspace'
assert state is not None and state.cursor == 7
async def test_targeted_reconciliation_only_updates_requested_workspace_side_effects(projection_context):
application, _session_factory = projection_context
provider = _Provider(
[_snapshot(7, workspaces=[])],
deltas=[_delta(workspaces=[_workspace(revision=8, name='JIT Workspace')])],
)
service = DirectoryProjectionService(application, provider, INSTANCE_UUID)
await service.initialize()
service._reconcile_entitlement_snapshot_set = AsyncMock()
service._update_entitlement_workspace_activity = AsyncMock()
service._publish_runtime_execution_projection = Mock()
await service.reconcile_workspaces((WORKSPACE_UUID,))
service._reconcile_entitlement_snapshot_set.assert_not_awaited()
service._update_entitlement_workspace_activity.assert_awaited_once()
assert service._update_entitlement_workspace_activity.await_args.kwargs == {
'requested_workspace_uuids': {WORKSPACE_UUID},
}
service._publish_runtime_execution_projection.assert_called_once()
assert service._publish_runtime_execution_projection.call_args.kwargs == {
'affected_workspace_uuids': {WORKSPACE_UUID},
}
async def test_initial_snapshot_projects_core_owned_rows(projection_context):
application, session_factory = projection_context
reconcile_execution_projection = Mock()
@@ -818,72 +735,6 @@ async def test_each_replica_consumes_events_with_its_own_cursor(projection_conte
assert second_provider.after_cursors == [1, 2]
async def test_concurrent_sync_once_calls_are_serialized_per_service(projection_context):
application, _session_factory = projection_context
class _ConcurrentProvider(_Provider):
def __init__(self) -> None:
super().__init__([_snapshot(1)])
self.first_fetch_started = asyncio.Event()
self.release_first_fetch = asyncio.Event()
self.active_fetches = 0
self.max_active_fetches = 0
async def fetch_events(
self,
instance_uuid: str,
after_cursor: int,
limit: int,
) -> DirectoryEventBatch:
assert instance_uuid == INSTANCE_UUID
assert limit == 100
self.after_cursors.append(after_cursor)
self.active_fetches += 1
self.max_active_fetches = max(self.max_active_fetches, self.active_fetches)
try:
if len(self.after_cursors) == 1:
self.first_fetch_started.set()
await self.release_first_fetch.wait()
cursor = after_cursor + 1
return DirectoryEventBatch(
instance_uuid=instance_uuid,
after_cursor=after_cursor,
cursor=cursor,
high_water_cursor=cursor,
events=(
DirectoryEvent(
cursor=cursor,
uuid=f'40000000-0000-4000-8000-{cursor:012d}',
aggregate_uuid=WORKSPACE_UUID,
event_type='entitlement.changed',
revision=cursor,
payload={
'workspace_uuid': WORKSPACE_UUID,
'entitlement_revision': cursor,
},
created_at=datetime.datetime(2026, 7, 24, 12, cursor, tzinfo=datetime.UTC),
),
),
)
finally:
self.active_fetches -= 1
provider = _ConcurrentProvider()
service = DirectoryProjectionService(application, provider, INSTANCE_UUID)
await service.initialize()
first = asyncio.create_task(service.sync_once())
await provider.first_fetch_started.wait()
second = asyncio.create_task(service.sync_once())
await asyncio.sleep(0)
provider.release_first_fetch.set()
await asyncio.gather(first, second)
assert provider.max_active_fetches == 1
assert provider.after_cursors == [1, 2]
assert service._consumer_cursor == 3
async def test_snapshot_coverage_allows_lagging_replica_to_replay_receipts(projection_context):
application, session_factory = projection_context
event_two = DirectoryEvent(
@@ -964,7 +964,6 @@ async def test_scoped_session_rejects_raw_or_unapproved_sql(
sa.func.date_trunc('hour', sa.column('timestamp')),
sa.func.length(sa.literal('value')),
sa.func.nullif(sa.literal('value'), sa.literal('')),
sa.func.strftime('%Y-%m-%d %H:00', sa.column('timestamp')),
),
sa.select(sa.column('embedding').op('<=>')(sa.literal([0.1]))),
sa.select(sa.cast(sa.column('embedding'), Vector(384))),
@@ -978,19 +977,6 @@ async def test_scoped_sql_structure_allows_only_the_production_vocabulary(statem
_validate_scoped_statement_call((statement,), {})
async def test_scoped_session_executes_sqlite_strftime() -> None:
engine = create_async_engine('sqlite+aiosqlite:///:memory:')
try:
async with TenantUnitOfWork(engine, 'workspace-a') as uow:
result = await uow.session.execute(
sa.select(sa.func.strftime('%Y-%m-%d %H:00', sa.literal('2026-08-28 03:45:00')))
)
assert result.scalar_one() == '2026-08-28 03:00'
finally:
await engine.dispose()
async def test_scoped_sql_rejects_public_execution_options() -> None:
statement = sa.select(sa.literal(1))
with pytest.raises(ScopedSessionTransactionError, match='execution options'):
-113
View File
@@ -1,113 +0,0 @@
"""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 == ''
+1 -88
View File
@@ -55,7 +55,7 @@ finally:
# ---------------------------------------------------------------------------
def make_runner(output_key: str = 'response', response_handling: str = 'reply') -> N8nServiceAPIRunner:
def make_runner(output_key: str = 'response') -> N8nServiceAPIRunner:
ap = Mock()
ap.logger = Mock()
pipeline_config = {
@@ -63,7 +63,6 @@ def make_runner(output_key: str = 'response', response_handling: str = 'reply')
'n8n-service-api': {
'webhook-url': 'http://test-n8n/webhook',
'output-key': output_key,
'response-handling': response_handling,
'auth-type': 'none',
}
}
@@ -288,7 +287,6 @@ def make_http_session_mock(response_bytes: bytes, status: int = 200):
"""Mock httpclient.get_session() returning a session whose post() yields response_bytes."""
mock_response = make_mock_response([response_bytes], status=status)
mock_response.status = status
mock_response.headers = {}
mock_cm = AsyncMock()
mock_cm.__aenter__ = AsyncMock(return_value=mock_response)
@@ -316,91 +314,6 @@ async def test_call_webhook_nonstream_adapter_plain_json():
assert results[0].content == 'result text'
@pytest.mark.asyncio
@pytest.mark.parametrize('status', [200, 201, 202, 204])
@pytest.mark.parametrize(
'response_body',
[
b'{"message":"Workflow was started"}',
b'{"response":"must not be forwarded"}',
b'plain acknowledgement',
],
)
async def test_call_webhook_ignore_response_body(response_body: bytes, status: int):
"""Ignore mode accepts any HTTP 2xx response without emitting chat output."""
runner = make_runner(response_handling='ignore')
query = make_query(is_stream=False)
http_session = make_http_session_mock(response_body, status=status)
with patch('langbot.pkg.provider.runners.n8nsvapi.httpclient.get_session', return_value=http_session):
results = []
async for message in runner._call_webhook(query):
results.append(message)
assert results == []
@pytest.mark.asyncio
async def test_call_webhook_ignore_releases_without_reading_response_body():
"""Ignore mode returns after the success status without waiting for the body."""
runner = make_runner(response_handling='ignore')
query = make_query(is_stream=False)
mock_response = make_mock_response([], status=202)
mock_response.headers = {}
mock_response.release = Mock()
async def fail_if_read(_size):
raise AssertionError('ignore mode must not read the response body')
yield b''
mock_response.content.iter_chunked = fail_if_read
mock_cm = AsyncMock()
mock_cm.__aenter__ = AsyncMock(return_value=mock_response)
mock_cm.__aexit__ = AsyncMock(return_value=False)
mock_session = Mock()
mock_session.post = Mock(return_value=mock_cm)
with patch('langbot.pkg.provider.runners.n8nsvapi.httpclient.get_session', return_value=mock_session):
results = [message async for message in runner._call_webhook(query)]
assert results == []
mock_response.release.assert_called_once_with()
@pytest.mark.asyncio
@pytest.mark.parametrize('status', [201, 202, 204])
async def test_call_webhook_reply_mode_preserves_http_200_contract(status: int):
"""Reply mode remains backward compatible and rejects non-200 statuses."""
runner = make_runner(response_handling='reply')
query = make_query(is_stream=False)
http_session = make_http_session_mock(b'', status=status)
with patch('langbot.pkg.provider.runners.n8nsvapi.httpclient.get_session', return_value=http_session):
with pytest.raises(N8nAPIError, match=f'n8n webhook call failed: {status}'):
async for _ in runner._call_webhook(query):
pass
@pytest.mark.asyncio
async def test_call_webhook_ignore_mode_preserves_http_error():
"""Ignore mode must not swallow a failed n8n webhook response."""
runner = make_runner(response_handling='ignore')
query = make_query(is_stream=False)
http_session = make_http_session_mock(b'{"error":"unavailable"}', status=500)
with patch('langbot.pkg.provider.runners.n8nsvapi.httpclient.get_session', return_value=http_session):
with pytest.raises(N8nAPIError, match='n8n webhook call exception'):
async for _ in runner._call_webhook(query):
pass
@pytest.mark.asyncio
async def test_invalid_response_handling_is_rejected():
"""Configuration errors should fail fast instead of silently changing reply behavior."""
with pytest.raises(ValueError, match='Invalid n8n response-handling'):
make_runner(response_handling='unexpected')
@pytest.mark.asyncio
async def test_call_webhook_stream_adapter_stream_format():
"""Stream adapter + stream format → MessageChunks, last is_final."""
+1 -42
View File
@@ -1,11 +1,8 @@
"""Tests for DingTalk API payload helpers."""
import json
from contextlib import asynccontextmanager
from unittest.mock import AsyncMock
from langbot.libs.dingtalk_api.api import DingTalkClient, _stringify_card_param_map
from langbot.pkg.utils import httpclient
from langbot.libs.dingtalk_api.api import _stringify_card_param_map
def test_dingtalk_card_param_map_stringifies_select_component_arrays():
@@ -43,41 +40,3 @@ def test_dingtalk_card_param_map_stringifies_unregistered_structures():
assert params['other'] == '["A"]'
assert params['empty'] == ''
async def test_create_card_embeds_layout_config_as_template_parameter(monkeypatch):
response = type('Response', (), {'status_code': 200})()
post = AsyncMock(return_value=response)
@asynccontextmanager
async def client_context():
yield type('HttpClient', (), {'post': post})()
client = object.__new__(DingTalkClient)
client.access_token = 'access-token'
client.robot_code = 'robot-code'
client.key = 'client-id'
client.logger = None
client.check_access_token = AsyncMock(return_value=True)
client._http_client_context = client_context
monkeypatch.setattr(httpclient, 'response_text', AsyncMock(return_value='{}'))
original_params = {'content': 'hello'}
delivered = await client.create_and_deliver_card(
card_template_id='template-id',
out_track_id='track-id',
open_space_id='dtv1.card//IM_ROBOT.user-id',
is_group=False,
card_param_map=original_params,
card_data_config={'autoLayout': True},
)
request_body = post.await_args.kwargs['json']
assert delivered is True
assert request_body['cardData'] == {
'cardParamMap': {
'content': 'hello',
'config': '{"autoLayout": true}',
}
}
assert original_params == {'content': 'hello'}
+1 -120
View File
@@ -1,7 +1,7 @@
"""Tests for Lark adapter helper behavior."""
import threading
from unittest.mock import AsyncMock, MagicMock
from unittest.mock import MagicMock
import pytest
@@ -12,7 +12,6 @@ from langbot.pkg.platform.sources.lark import (
_lark_completed_input_lines,
_lark_current_input_defs,
_lark_extract_action_form_inputs,
_lark_final_layout_texts,
_lark_should_update_stream_element,
_lark_visible_form_content,
)
@@ -222,121 +221,3 @@ def test_lark_completed_input_lines_display_select_value_from_object():
)
assert lines == ['✅ xialaB']
def test_lark_final_layout_texts_normal_round_drops_resume_placeholder():
"""Non-resume final chunk: the reply must land in the main element only.
Regression: rendering the resume placeholder too duplicated the reply,
because the accumulated streaming text equals the final text on a normal
round (e.g. 'It is Sep 1, 2026.\nIt is Sep 1, 2026.' in the card).
"""
main_text, resume_text = _lark_final_layout_texts(
resume_from=False,
text_message='It is Sep 1, 2026, 15:09:15.',
pre_pause_cached=None,
resume_cached='It is Sep 1, 2026, 15:09:15.',
)
assert main_text == 'It is Sep 1, 2026, 15:09:15.'
assert resume_text == ''
def test_lark_final_layout_texts_resume_round_keeps_both_segments():
"""Dify HITL resume final chunk: pre-pause text and resumed text differ,
both segments stay visible."""
main_text, resume_text = _lark_final_layout_texts(
resume_from=True,
text_message='resumed answer',
pre_pause_cached='partial answer before pause',
resume_cached='resumed answer',
)
assert main_text == 'partial answer before pause'
assert resume_text == 'resumed answer'
def test_lark_final_layout_texts_resume_round_without_pre_pause_falls_back():
main_text, resume_text = _lark_final_layout_texts(
resume_from=True,
text_message='answer',
pre_pause_cached=None,
resume_cached='answer',
)
assert main_text == 'answer'
assert resume_text == 'answer'
def test_lark_final_layout_texts_resume_round_empty_pre_pause_kept_empty():
"""Dify paused before emitting any text: the pre-pause cache is a valid
empty string and must NOT be treated as a cache miss.
Regression: `pre_pause_cached or text_message` fell back to the full
text, so the final card rendered ('resumed answer', 'resumed answer')
and duplicated the reply.
"""
main_text, resume_text = _lark_final_layout_texts(
resume_from=True,
text_message='resumed answer',
pre_pause_cached='',
resume_cached='resumed answer',
)
assert main_text == ''
assert resume_text == 'resumed answer'
def _build_resume_final_chunk_adapter(message_text: str):
"""Build a LarkAdapter whose card state mimics a Dify HITL round that
paused before emitting any text, then resumed and completed."""
adapter = LarkAdapter.model_construct(
api_client=MagicMock(),
message_converter=MagicMock(yiri2target=AsyncMock(return_value=([[{'tag': 'text', 'text': message_text}]], []))),
)
adapter.config = {'app_type': 'self'}
LarkAdapter.get_app_access_token = lambda self: None
LarkAdapter.get_tenant_access_token = lambda self, tenant_key: None
adapter.card_id_dict = {'msg-1': 'card-1'}
adapter.card_streaming_text = {'card-1': message_text}
adapter.card_pre_pause_text = {'card-1': ''}
adapter.card_resume_transitioned = {'card-1'}
adapter.card_sequence_dict = {}
adapter.card_last_accessed = {}
adapter.card_cleanup_at = 0.0
adapter.card_id_to_source_ids = {}
adapter.reply_message_card_ids = {}
adapter.card_form_content = {}
adapter.card_form_input_defs = {}
adapter.card_form_inputs = {}
adapter._update_card_layout = AsyncMock()
return adapter
@pytest.mark.asyncio
async def test_reply_message_chunk_resume_final_with_empty_pre_pause_keeps_main_empty():
"""End-to-end regression via reply_message_chunk: Dify paused before any
text, so the pre-pause cache is ''. The final card update must render the
resumed answer only once (empty main text + resume placeholder), not
twice as ('resumed answer', 'resumed answer')."""
adapter = _build_resume_final_chunk_adapter('resumed answer')
bot_message = MagicMock(
resp_message_id='msg-1',
msg_sequence=1,
spec=['resp_message_id', 'msg_sequence', '_resume_from_form'],
)
bot_message._resume_from_form = True
message_source = MagicMock(source_platform_object=None)
await adapter.reply_message_chunk(
message_source,
bot_message,
MagicMock(),
is_final=True,
)
adapter._update_card_layout.assert_awaited_once()
layout_kwargs = adapter._update_card_layout.await_args.kwargs
assert layout_kwargs['text_message'] == ''
assert layout_kwargs['resume_placeholder_text'] == 'resumed answer'
@@ -3,25 +3,18 @@ from __future__ import annotations
import pytest
from unittest.mock import MagicMock
from linebot.v3.webhooks import TextMessageContent, UserMentionee, AllMentionee
from linebot.v3.webhooks import TextMessageContent
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
):
def _make_event(*, source_type: str, user_id, group_id=None, room_id=None, message_id: str, text: str = 'hi'):
event = MagicMock()
event.timestamp = 1700000000000
message = MagicMock(spec=TextMessageContent)
message.id = message_id
message.text = text
message.mention = mention
event.message = message
event.message = MagicMock(spec=TextMessageContent)
event.message.id = message_id
event.message.text = text
event.message.webhook_event_id = f'webhook-{message_id}'
event.message.timestamp = event.timestamp
@@ -37,21 +30,16 @@ def _make_event(
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)
result1 = await line.LINEEventConverter.target2yiri(event1, bot_client=None)
result2 = await line.LINEEventConverter.target2yiri(event2, bot_client=None)
assert result1.sender.id == 'U-stable-user'
assert result1.sender.id == result2.sender.id
@@ -60,12 +48,11 @@ async def test_user_message_launcher_id_stable_across_messages() -> None:
@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)
result1 = await line.LINEEventConverter.target2yiri(event1, bot_client=None)
result2 = await line.LINEEventConverter.target2yiri(event2, bot_client=None)
assert result1.sender.group.id == 'G-stable-group'
assert result1.sender.group.id == result2.sender.group.id
@@ -74,186 +61,9 @@ async def test_group_message_uses_group_id_not_message_id() -> None:
@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)
result = await line.LINEEventConverter.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,59 +0,0 @@
"""Tests for WecomAdapter.send_message content-key handling."""
import pytest
import langbot_plugin.api.entities.builtin.platform.message as platform_message
from langbot.pkg.platform.sources.wecom import WecomAdapter
class StubWecomClient:
def __init__(self):
self.calls = []
async def get_media_id(self, msg):
return 'MEDIA_ID_123'
async def send_private_msg(self, user_id, agent_id, text):
self.calls.append(('text', user_id, agent_id, text))
async def send_image(self, user_id, agent_id, media_id):
self.calls.append(('image', user_id, agent_id, media_id))
async def send_voice(self, user_id, agent_id, media_id):
self.calls.append(('voice', user_id, agent_id, media_id))
async def send_file(self, user_id, agent_id, media_id):
self.calls.append(('file', user_id, agent_id, media_id))
def _make_adapter():
adapter = WecomAdapter.model_construct(bot=StubWecomClient())
return adapter
@pytest.mark.asyncio
@pytest.mark.parametrize(
('part', 'expected_type'),
[
(platform_message.Image(url='https://example.com/x.jpg'), 'image'),
(platform_message.Voice(url='https://example.com/x.amr'), 'voice'),
(platform_message.File(url='https://example.com/x.pdf', name='x.pdf'), 'file'),
],
)
async def test_send_message_dispatches_media_by_id(part, expected_type):
adapter = _make_adapter()
chain = platform_message.MessageChain([part])
await adapter.send_message('person', 'USER1|1000001', chain)
assert adapter.bot.calls == [(expected_type, 'USER1', 1000001, 'MEDIA_ID_123')]
@pytest.mark.asyncio
async def test_send_message_text_still_works():
adapter = _make_adapter()
chain = platform_message.MessageChain([platform_message.Plain(text='hello')])
await adapter.send_message('person', 'USER1|1000001', chain)
assert adapter.bot.calls == [('text', 'USER1', 1000001, 'hello')]
@@ -44,86 +44,6 @@ def test_webhook_dispatch_tasks_are_bounded():
assert len(client._dispatch_tasks) == 100
@pytest.mark.asyncio
async def test_ws_initial_stream_frame_precedes_pipeline_dispatch(monkeypatch):
from langbot.libs.wecom_ai_bot_api import ws_client as ws_client_module
order = []
logger = types.SimpleNamespace(
debug=Mock(),
error=Mock(),
warning=Mock(),
)
client = WecomBotWsClient('bot-id', 'secret', logger)
async def parse_message(*args, **kwargs):
del args, kwargs
return {'msgid': 'msg-1', 'type': 'single', 'userid': 'user-1'}
async def reply_stream(*args, **kwargs):
del args, kwargs
order.append('initial-frame')
return {}
async def dispatch_event(event):
del event
order.append('pipeline-dispatch')
monkeypatch.setattr(ws_client_module, 'parse_wecom_bot_message', parse_message)
monkeypatch.setattr(ws_client_module.wecombotevent, 'WecomBotEvent', lambda data: data)
client.reply_stream = reply_stream
client._dispatch_event = dispatch_event
await client._handle_message_callback({'headers': {'req_id': 'req-1'}, 'body': {}})
assert order == ['initial-frame', 'pipeline-dispatch']
@pytest.mark.asyncio
async def test_ws_initial_stream_failure_still_dispatches_message(monkeypatch):
from langbot.libs.wecom_ai_bot_api import ws_client as ws_client_module
dispatched = []
class Logger:
def __init__(self):
self.warnings = []
async def debug(self, message):
del message
async def error(self, message):
raise AssertionError(message)
async def warning(self, message):
self.warnings.append(message)
logger = Logger()
client = WecomBotWsClient('bot-id', 'secret', logger)
async def parse_message(*args, **kwargs):
del args, kwargs
return {'msgid': 'msg-1', 'type': 'single', 'userid': 'user-1'}
async def reply_stream(*args, **kwargs):
del args, kwargs
raise ConnectionError('simulated reply failure')
async def dispatch_event(event):
dispatched.append(event)
monkeypatch.setattr(ws_client_module, 'parse_wecom_bot_message', parse_message)
monkeypatch.setattr(ws_client_module.wecombotevent, 'WecomBotEvent', lambda data: data)
client.reply_stream = reply_stream
client._dispatch_event = dispatch_event
await client._handle_message_callback({'headers': {'req_id': 'req-1'}, 'body': {}})
assert len(dispatched) == 1
assert len(logger.warnings) == 1
assert 'simulated reply failure' in logger.warnings[0]
def test_extract_template_card_action_supports_nested_button_key():
task_id, event_key, card_type = extract_template_card_action(
{
@@ -1,4 +1,3 @@
import uuid
from types import SimpleNamespace
from unittest.mock import AsyncMock
@@ -50,29 +49,7 @@ 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 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
assert kwargs['msgid'].startswith('langbot_')
@pytest.mark.asyncio

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