diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 14919cc2d..5a167c64b 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -160,6 +160,7 @@ In this repo: - `pkg/plugin/handler.py` exposes LangBot actions to the runtime and calls runtime actions for plugin operations. - `pkg/provider/tools/loaders/plugin.py` exposes plugin Tool components to LLM runners. - Pipeline handlers emit SDK events such as normal-message events and prompt-processing events. +- [Certified plugin policy](docs/architecture/certified-plugins.md) defines Core's archive-fact, admission, and tenant-log-visibility boundary; the SDK remains responsible for certificate verification. In `langbot-plugin-sdk`: diff --git a/README.md b/README.md index 5de92d3ad..81277fd54 100644 --- a/README.md +++ b/README.md @@ -22,7 +22,7 @@ English / [简体中文](README_CN.md) / [繁體中文](README_TW.md) / [日本 Features | Docs | API | -Cloud | +Cloud | Plugin Market | Roadmap @@ -65,7 +65,11 @@ Click the Star and Watch buttons in the top-right corner of the repository to ge ### ☁️ LangBot Cloud (Recommended) -**[LangBot Cloud](https://space.langbot.app/cloud)** — Zero deployment, ready to use. +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +Zero deployment, ready to use. ### One-Line Launch diff --git a/README_CN.md b/README_CN.md index aad5c6efa..d49799cf6 100644 --- a/README_CN.md +++ b/README_CN.md @@ -24,7 +24,7 @@ 特性 | 文档 | API | -Cloud | +Cloud | 扩展市场 | 路线图 @@ -65,7 +65,11 @@ LangBot 是一个**开源的生产级平台**,用于构建 AI 驱动的即时 ### ☁️ LangBot Cloud(推荐) -**[LangBot Cloud](https://space.langbot.app/cloud)** — 免部署,开箱即用。 +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +免部署,开箱即用。 ### 一键启动 diff --git a/README_ES.md b/README_ES.md index 04ebc78a3..d6749d5bb 100644 --- a/README_ES.md +++ b/README_ES.md @@ -64,7 +64,11 @@ Haga clic en los botones Star y Watch en la esquina superior derecha del reposit ### ☁️ LangBot Cloud (Recomendado) -**[LangBot Cloud](https://space.langbot.app/cloud)** — Sin despliegue, listo para usar. +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +Sin despliegue, listo para usar. ### Lanzamiento en una línea diff --git a/README_FR.md b/README_FR.md index 78d99c692..71cda2887 100644 --- a/README_FR.md +++ b/README_FR.md @@ -64,7 +64,11 @@ Cliquez sur les boutons Star et Watch dans le coin supérieur droit du dépôt p ### ☁️ LangBot Cloud (Recommandé) -**[LangBot Cloud](https://space.langbot.app/cloud)** — Sans déploiement, prêt à utiliser. +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +Sans déploiement, prêt à utiliser. ### Lancement en une ligne diff --git a/README_JP.md b/README_JP.md index b876bd4ee..d410e6efa 100644 --- a/README_JP.md +++ b/README_JP.md @@ -64,7 +64,11 @@ LangBot は、AI搭載のインスタントメッセージングボットを構 ### ☁️ LangBot Cloud(推奨) -**[LangBot Cloud](https://space.langbot.app/cloud)** — デプロイ不要、すぐに使えます。 +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +デプロイ不要、すぐに使えます。 ### ワンライン起動 diff --git a/README_KO.md b/README_KO.md index b28d3ec2c..df06da4ab 100644 --- a/README_KO.md +++ b/README_KO.md @@ -64,7 +64,11 @@ LangBot은 AI 기반 인스턴트 메시징 봇을 구축하기 위한 **오픈 ### ☁️ LangBot Cloud (추천) -**[LangBot Cloud](https://space.langbot.app/cloud)** — 배포 없이 바로 사용. +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +배포 없이 바로 사용. ### 원라인 실행 diff --git a/README_RU.md b/README_RU.md index f6c1f8bce..ca65b6556 100644 --- a/README_RU.md +++ b/README_RU.md @@ -64,7 +64,11 @@ LangBot — это **платформа с открытым исходным к ### ☁️ LangBot Cloud (Рекомендуется) -**[LangBot Cloud](https://space.langbot.app/cloud)** — Без развёртывания, готово к использованию. +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +Без развёртывания, готово к использованию. ### Запуск одной командой diff --git a/README_TW.md b/README_TW.md index 140515650..36d3c765a 100644 --- a/README_TW.md +++ b/README_TW.md @@ -66,7 +66,11 @@ LangBot 是一個**開源的生產級平台**,用於建構 AI 驅動的即時 ### ☁️ LangBot Cloud(推薦) -**[LangBot Cloud](https://space.langbot.app/cloud)** — 免部署,開箱即用。 +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +免部署,開箱即用。 ### 一鍵啟動 diff --git a/README_VI.md b/README_VI.md index 356f577b3..2a412d992 100644 --- a/README_VI.md +++ b/README_VI.md @@ -64,7 +64,11 @@ Nhấp vào các nút Star và Watch ở góc trên bên phải của kho lưu t ### ☁️ LangBot Cloud (Khuyên dùng) -**[LangBot Cloud](https://space.langbot.app/cloud)** — Không cần triển khai, sẵn sàng sử dụng. +[![Deploy on LangBot Cloud](res/langbot-cloud.svg)](https://cloud.langbot.app) + +[cloud.langbot.app](https://cloud.langbot.app) + +Không cần triển khai, sẵn sàng sử dụng. ### Khởi chạy một dòng diff --git a/docs/architecture/certified-plugins.md b/docs/architecture/certified-plugins.md new file mode 100644 index 000000000..f4c9400ac --- /dev/null +++ b/docs/architecture/certified-plugins.md @@ -0,0 +1,74 @@ +# Certified Plugins + +## Admission boundary + +Core verifies a plugin archive **before** artifact storage, `PluginSetting` +persistence, or a Plugin Runtime apply request. It calls the SDK public +`langbot_plugin.certification.verify_archive()` API, which reads the strict +certificate envelope from the ZIP comment and verifies the signed normalized +ZIP digest without extracting the payload. + +Core retains the normalized digest (`normalized_zip_digest()`), verification +state, declared shared-runtime profile, key ID, selected admission profile, and +stable admission code in the durable plugin `install_info._certification` +record. The record belongs to the installation row; no schema migration is +needed for this additive JSON metadata. + +## Trusted issuer configuration + +Configure the non-secret Ed25519 public-key ring in `data/config.yaml`: + +```yaml +plugin: + certification: + trusted_public_keys: + issuer-2026-q3: "" +``` + +Key IDs must match the SDK envelope. Values are standard base64 raw public +keys, not private/signing keys. An invalid key-ring configuration is rejected +rather than weakening verification. Keep active issuer keys during a rotation +until archives signed by retired IDs are no longer installed. + +## Admission matrix + +| Deployment | SDK verification | Explicit `administrator_force` | Result | +| --- | --- | --- | --- | +| Cloud | valid envelope declaring `shared-runtime-v1` | any | admitted to the shared profile | +| Cloud | absent | any | reject before storage with `CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_REQUIRED` | +| Cloud | malformed, untrusted, invalid, or non-shared | any | reject before storage with `CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_INVALID` | +| OSS | absent legacy envelope | any | admitted to the dedicated profile | +| OSS | valid envelope declaring `shared-runtime-v1` | any | selected shared profile | +| OSS | malformed or invalid declaration | false | reject with `CERTIFIED_PLUGIN_OSS_FORCE_REQUIRED` | +| OSS | malformed or invalid declaration | true | admitted to the dedicated profile | + +`administrator_force` is deliberately strict: it is recognized only when the +install request carries boolean `true`. The local upload endpoint accepts the +multipart field `administrator_force=true`; GitHub and marketplace install +payloads carry the same field. The existing resource-manage authorization fence +protects those endpoints. A force never creates a Cloud dedicated fallback. + +## Runtime and logs + +The current Plugin Runtime control protocol has one process-wide runtime profile +per Core instance. In Cloud that existing profile is `shared`; Cloud admission +therefore prevents an archive that did not select `shared-runtime-v1` from +reaching its apply API. In OSS the existing `oss_dev` runtime remains the +dedicated compatibility profile. Core records the selected profile for every +installation so a future multi-runtime control protocol can consume it without +re-verifying an already persisted archive. + +The existing public plugin-log boundary already applies the immutable +installation binding (including workspace UUID) through +`RuntimeConnectionHandler.installation_scope()` before requesting logs. This is +the actual tenant exposure boundary, so valid shared certificates use that +binding-scoped transport; Core does not invent a second log stream or expose +process-wide log output. Dedicated and invalid/legacy installations use the +same existing installation scope. + +## SDK versioning + +Core intentionally continues to declare `langbot-plugin==0.5.8` until the SDK +beta containing this public certification API is released. Local development +and the integration tests may install the SDK source checkout, but this Core +change does not publish or pin a prerelease. diff --git a/packaging/fnos/README.md b/packaging/fnos/README.md index 1bd672b67..193a184f1 100644 --- a/packaging/fnos/README.md +++ b/packaging/fnos/README.md @@ -2,6 +2,10 @@ This directory packages LangBot as a `.fpk` app for the fnOS App Store. It is a native deployment: no Docker involved — uv creates a Python virtual environment directly on the NAS, and Node.js v22 from the fnOS App Store provides the Box sandbox and npx MCP capabilities. +## Privilege Model + +Per the [fnOS privilege docs](https://developer.fnnas.com/docs/core-concepts/privilege/), the whole app runs as the dedicated package user (`run-as: package`; username auto-generated by fnOS from `appname`) — no root anywhere. The official App Store apps `nodejs_v22` and `python312` (declared in `manifest` via `install_dep_apps`) are borrowed from `/var/apps/` cross-app. `HOME` is pointed at `/.home` inside the persistent share so tool caches (uv, npm/npx MCP) stay writable regardless of the generated user's system home. + ## Directory Structure ``` @@ -10,10 +14,10 @@ packaging/fnos/ ├── build.sh # One-shot build script (shared by local and CI) ├── LICENSE ├── config/ -│ ├── privilege # Privilege config (run-as: root) +│ ├── privilege # Privilege config (run-as: package, no root) │ └── resource # Persistent data share declaration (langbot/data) ├── cmd/ # Lifecycle scripts (fnOS invokes them with TRIM_* env vars) -│ ├── main # Service start/stop manager (start/stop/status, owns PID/log) +│ ├── main # Service start/stop manager (start/stop/status, owns PID/log); started via setsid, stop kills the whole process group │ ├── install_init # Pre-install hook │ ├── install_callback # Post-install hook: create venv, uv sync deps, seed config.yaml port │ ├── upgrade_init # Pre-upgrade hook diff --git a/packaging/fnos/app/desktop/langbot.main.url b/packaging/fnos/app/desktop/langbot.main.url index 513870fe8..1717216cc 100644 --- a/packaging/fnos/app/desktop/langbot.main.url +++ b/packaging/fnos/app/desktop/langbot.main.url @@ -1,7 +1,7 @@ { "title": "LangBot", "icon": "images/icon-256.png", - "type": "url", + "type": "iframe", "protocol": "http", "port": "${wizard_port}", "url": "/", diff --git a/packaging/fnos/app/ui/config b/packaging/fnos/app/ui/config index d45765dd4..2908f45de 100644 --- a/packaging/fnos/app/ui/config +++ b/packaging/fnos/app/ui/config @@ -3,7 +3,7 @@ "langbot.main": { "title": "LangBot", "icon": "images/icon-{0}.png", - "type": "url", + "type": "iframe", "protocol": "http", "port": "${wizard_port}", "url": "/", diff --git a/packaging/fnos/build.sh b/packaging/fnos/build.sh index 74675e2a3..ef0f135fa 100644 --- a/packaging/fnos/build.sh +++ b/packaging/fnos/build.sh @@ -73,8 +73,11 @@ rsync -a \ [ -d "${FPK_DIR}/app/langbot/web/dist" ] || { echo "ERROR: web/dist missing after rsync!" >&2; exit 1; } echo " Source synced ($(du -sh "${FPK_DIR}/app/langbot" | cut -f1))" -# --- 2.5 Download bundled uv binaries (offline install on NAS) --- -echo "[2.5/5] Downloading bundled uv binaries..." +# --- 2.5 Download bundled uv binary (offline install on NAS) --- +# Python comes from the official python312 App Store app (see manifest +# install_dep_apps); only uv is carried in the package. The download is a +# hard requirement — the build fails without it (no fallback installs). +echo "[2.5/5] Downloading bundled uv binary..." UV_VERSION="0.12.9" mkdir -p "${FPK_DIR}/app/bin" for arch in x86_64 aarch64; do @@ -84,15 +87,13 @@ for arch in x86_64 aarch64; do continue fi tmp="$(mktemp -d)" - if curl -sSL -o "${tmp}/uv.tar.gz" \ + curl -fsSL -o "${tmp}/uv.tar.gz" \ "https://github.com/astral-sh/uv/releases/download/${UV_VERSION}/uv-${arch}-unknown-linux-gnu.tar.gz" \ && tar xzf "${tmp}/uv.tar.gz" -C "${tmp}" \ - && cp "${tmp}/uv-${arch}-unknown-linux-gnu/uv" "${out}"; then - chmod +x "${out}" - echo " uv-${arch} downloaded (${UV_VERSION})" - else - echo " WARNING: failed to download uv for ${arch}, install will fall back to online install" >&2 - fi + && cp "${tmp}/uv-${arch}-unknown-linux-gnu/uv" "${out}" \ + && chmod +x "${out}" \ + && echo " uv-${arch} downloaded (${UV_VERSION})" \ + || { echo "ERROR: failed to download uv for ${arch}" >&2; rm -rf "${tmp}"; exit 1; } rm -rf "${tmp}" done diff --git a/packaging/fnos/cmd/install_callback b/packaging/fnos/cmd/install_callback index 4eac06acc..96b7cff79 100755 --- a/packaging/fnos/cmd/install_callback +++ b/packaging/fnos/cmd/install_callback @@ -19,7 +19,30 @@ cd "${APP_DIR}" || { } # --- Ensure data directory exists --- -mkdir -p "${DATA_DIR}/plugins" "${DATA_DIR}/box" "${DATA_DIR}/logs" 2>/dev/null || true +mkdir -p "${DATA_DIR}/plugins" "${DATA_DIR}/box" "${DATA_DIR}/logs" || { + echo "Data share not writable: ${DATA_DIR}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} + +# --- Writable HOME and tool caches --- +# Under run-as: package the generated user's system HOME and the app install +# dir (TRIM_APPDEST) can be read-only; uv aborts if it cannot initialise its +# cache. Point HOME / UV_CACHE_DIR at the persistent data share and keep the +# venv there too (UV_PROJECT_ENVIRONMENT), instead of inside APP_DIR. +VENV_DIR="${DATA_DIR}/.venv" +export HOME="${DATA_DIR}/.home" +export UV_CACHE_DIR="${DATA_DIR}/.cache/uv" +# Never download a managed CPython (the download host is unreachable on many +# NAS networks); use the distro Python only — fail loudly if it is missing. +export UV_PYTHON_DOWNLOADS=never +mkdir -p "${HOME}" "${UV_CACHE_DIR}" || { + echo "Data share not writable: ${DATA_DIR}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} + +# Real uv/pip errors are captured here for diagnosis (the installer popup +# only shows the short message we write to TRIM_TEMP_LOGFILE). +DEBUG_LOG="${DATA_DIR}/logs/install-debug.log" # --- Pre-seed config.yaml with the user-selected web port --- # Data root points at the persistent share (LANGBOT_DATA_ROOT is exported by @@ -35,7 +58,10 @@ TEMPLATE_FILE="${APP_DIR}/src/langbot/templates/config.yaml" _patch_config() { local cfg_dir="$1" local cfg_file="${cfg_dir}/config.yaml" - mkdir -p "${cfg_dir}" 2>/dev/null || true + mkdir -p "${cfg_dir}" || { + echo "Cannot create config directory: ${cfg_dir}" > "${TRIM_TEMP_LOGFILE}" + exit 1 + } if [ ! -f "${cfg_file}" ] && [ -f "${TEMPLATE_FILE}" ]; then cp "${TEMPLATE_FILE}" "${cfg_file}" fi @@ -68,15 +94,27 @@ if [ ! -d "/var/apps/nodejs_v${NODE_VERSION}" ]; then exit 1 fi -# --- Python check --- -PYTHON_BIN="python3" -if ! command -v "${PYTHON_BIN}" >/dev/null 2>&1; then - PYTHON_BIN="python" -fi -if ! command -v "${PYTHON_BIN}" >/dev/null 2>&1; then - echo "Python not found on this system" > "${TRIM_TEMP_LOGFILE}" +# --- CPU architecture (must be resolved before locating bundled binaries) --- +ARCH=$(uname -m) +case "${ARCH}" in + x86_64|aarch64) ;; + *) + echo "Unsupported CPU architecture: ${ARCH}" > "${TRIM_TEMP_LOGFILE}" + exit 1 + ;; +esac + +# --- Python: official python312 App Store app (same borrow pattern as Node.js) --- +PYTHON_APP="python312" +PYTHON_BIN="/var/apps/${PYTHON_APP}/target/bin/python3" +if [ ! -d "/var/apps/${PYTHON_APP}" ]; then + echo "未找到官方 Python 环境:请先在应用中心安装 ${PYTHON_APP},再重新安装本应用。" > "${TRIM_TEMP_LOGFILE}" exit 1 fi +[ -x "${PYTHON_BIN}" ] || { + echo "Python 解释器缺失或不可执行:${PYTHON_BIN}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} PY_VER=$("${PYTHON_BIN}" -c 'import sys; print(f"{sys.version_info.major}.{sys.version_info.minor}")' 2>/dev/null) if [ -z "${PY_VER}" ]; then @@ -90,55 +128,32 @@ if [ "${PY_MAJOR}" -lt 3 ] || { [ "${PY_MAJOR}" -eq 3 ] && [ "${PY_MINOR}" -lt 1 exit 1 fi -# --- Resolve uv: bundled binary first, then online fallbacks --- -UV_BIN="" -ARCH=$(uname -m) -case "${ARCH}" in - x86_64) BUNDLED_UV="${TRIM_APPDEST}/bin/uv-x86_64" ;; - aarch64) BUNDLED_UV="${TRIM_APPDEST}/bin/uv-aarch64" ;; - *) BUNDLED_UV="" ;; -esac - -if [ -n "${BUNDLED_UV}" ] && [ -x "${BUNDLED_UV}" ]; then - mkdir -p "${TRIM_PKGVAR}/bin" - cp "${BUNDLED_UV}" "${TRIM_PKGVAR}/bin/uv" && chmod +x "${TRIM_PKGVAR}/bin/uv" - UV_BIN="${TRIM_PKGVAR}/bin/uv" -fi - -if [ -z "${UV_BIN}" ] && command -v uv >/dev/null 2>&1; then - UV_BIN="uv" -fi - -if [ -z "${UV_BIN}" ]; then - "${PYTHON_BIN}" -m pip install --user --no-cache-dir uv 2>/dev/null || \ - "${PYTHON_BIN}" -m pip install --no-cache-dir uv 2>/dev/null || \ - curl -LsSf https://astral.sh/uv/install.sh | sh 2>/dev/null || true - export PATH="${HOME}/.local/bin:${PATH}" - if command -v uv >/dev/null 2>&1; then - UV_BIN="uv" - elif [ -x "${HOME}/.local/bin/uv" ]; then - UV_BIN="${HOME}/.local/bin/uv" - fi -fi - -if [ -z "${UV_BIN}" ]; then - echo "无法获取 uv:内置二进制缺失且在线安装失败。请检查网络后重新安装。" > "${TRIM_TEMP_LOGFILE}" +# --- Resolve uv: run the bundled binary in place (single canonical path) --- +UV_BIN="${TRIM_APPDEST}/bin/uv-${ARCH}" +[ -x "${UV_BIN}" ] || { + echo "Bundled uv binary missing or not executable: ${UV_BIN}" > "${TRIM_TEMP_LOGFILE}" exit 1 -fi +} -# --- Create venv via uv --- -if [ ! -d ".venv" ]; then - "${UV_BIN}" venv .venv --python "${PYTHON_BIN}" || { - echo "Failed to create Python virtual environment via uv" > "${TRIM_TEMP_LOGFILE}" +# --- Create venv via uv (on the writable data share, not APP_DIR) --- +if [ ! -d "${VENV_DIR}" ]; then + { + echo "== whoami: $(id 2>&1)" + echo "== APP_DIR perms: $(ls -ld "${APP_DIR}" 2>&1)" + echo "== DATA_DIR perms: $(ls -ld "${DATA_DIR}" 2>&1)" + echo "== HOME=${HOME} UV_CACHE_DIR=${UV_CACHE_DIR}" + "${UV_BIN}" venv "${VENV_DIR}" --python "${PYTHON_BIN}" + } >> "${DEBUG_LOG}" 2>&1 || { + echo "Failed to create Python virtual environment via uv. See ${DEBUG_LOG}" > "${TRIM_TEMP_LOGFILE}" exit 1 } fi -# --- Sync dependencies --- -"${UV_BIN}" sync --extra seekdb || { - echo "Dependency sync failed. Check network connectivity." > "${TRIM_TEMP_LOGFILE}" +# --- Sync dependencies (install into the relocated venv) --- +if ! UV_PROJECT_ENVIRONMENT="${VENV_DIR}" "${UV_BIN}" sync --extra seekdb >> "${DEBUG_LOG}" 2>&1; then + echo "Dependency sync failed. Check network connectivity. See ${DEBUG_LOG}" > "${TRIM_TEMP_LOGFILE}" exit 1 -} +fi # --- Verify frontend dist --- if [ ! -d "web/dist" ] || [ -z "$(ls -A web/dist 2>/dev/null)" ]; then @@ -146,4 +161,9 @@ if [ ! -d "web/dist" ] || [ -z "$(ls -A web/dist 2>/dev/null)" ]; then exit 1 fi +# --- Prepare a writable HOME inside the data dir for tool caches (uv, +# npm/npx MCP). Everything already runs as the package user (run-as: +# package), so no chown is needed — files created here belong to it. +mkdir -p "${DATA_DIR}/.home" 2>/dev/null || true + exit 0 diff --git a/packaging/fnos/cmd/install_init b/packaging/fnos/cmd/install_init index 2e940ee03..22134488d 100755 --- a/packaging/fnos/cmd/install_init +++ b/packaging/fnos/cmd/install_init @@ -1,5 +1,12 @@ #!/bin/bash -# cmd/install_init - pre-install hook -# Nothing special to do before extraction. +# cmd/install_init - pre-install hook (runs before files are applied) +# Sweep stray processes from a previous failed/killed install: orphans whose +# parent was hard-killed keep running and hold the runtime/box ws ports, +# which breaks the new instance. No match is the normal case on a clean +# install — pkill exits 1. + +RUN_USER=$(id -un) +pkill -KILL -u "${RUN_USER}" -f "appcenter/langbot" 2>/dev/null || true +pkill -KILL -u "${RUN_USER}" -f "appshare/langbot" 2>/dev/null || true exit 0 diff --git a/packaging/fnos/cmd/main b/packaging/fnos/cmd/main index a3181c383..6dfdc6afd 100755 --- a/packaging/fnos/cmd/main +++ b/packaging/fnos/cmd/main @@ -27,27 +27,48 @@ if [ -z "${DATA_DIR}" ]; then DATA_DIR="${TRIM_PKGVAR}/data" fi export LANGBOT_DATA_ROOT="${DATA_DIR}" -mkdir -p "${DATA_DIR}" 2>/dev/null || true +mkdir -p "${DATA_DIR}" || { + echo "Data share not writable: ${DATA_DIR}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} -# --- Locate Python --- -PYTHON_BIN="python3" -! command -v "${PYTHON_BIN}" >/dev/null 2>&1 && PYTHON_BIN="python" +# --- Writable HOME / tool caches / venv on the data share --- +# Must match cmd/install_callback: under run-as: package neither the system +# HOME nor TRIM_APPDEST are guaranteed writable, and the venv lives at +# ${DATA_DIR}/.venv instead of inside the app dir. +export HOME="${DATA_DIR}/.home" +export UV_CACHE_DIR="${DATA_DIR}/.cache/uv" +# Never download a managed CPython — distro Python only (see install_callback) +export UV_PYTHON_DOWNLOADS=never +export UV_PROJECT_ENVIRONMENT="${DATA_DIR}/.venv" +mkdir -p "${HOME}" "${UV_CACHE_DIR}" || { + echo "Data share not writable: ${DATA_DIR}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} -# --- Locate uv --- -# install_callback puts the bundled uv binary at ${TRIM_PKGVAR}/bin/uv -UV_BIN="${TRIM_PKGVAR}/bin/uv" -if [ ! -x "${UV_BIN}" ]; then - UV_BIN="uv" -fi -if ! command -v "${UV_BIN}" >/dev/null 2>&1; then - UV_BIN="${HOME}/.local/bin/uv" -fi -if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then - UV_BIN="${HOME}/.cargo/bin/uv" -fi -if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then - UV_BIN="${APP_DIR}/.venv/bin/uv" -fi +# --- CPU architecture (must be resolved before locating bundled binaries) --- +ARCH=$(uname -m) +case "${ARCH}" in + x86_64|aarch64) ;; + *) + echo "Unsupported CPU architecture: ${ARCH}" > "${TRIM_TEMP_LOGFILE}" + exit 1 + ;; +esac + +# --- Locate Python: official python312 App Store app (same borrow pattern as Node.js) --- +PYTHON_BIN="/var/apps/python312/target/bin/python3" +[ -x "${PYTHON_BIN}" ] || { + echo "Python interpreter missing or not executable: ${PYTHON_BIN} (install the python312 app)" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} + +# --- Locate uv: bundled binary at its single canonical path in the app dir --- +UV_BIN="${TRIM_APPDEST}/bin/uv-${ARCH}" +[ -x "${UV_BIN}" ] || { + echo "Bundled uv binary missing or not executable: ${UV_BIN}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} case $1 in start) @@ -69,7 +90,7 @@ case $1 in exit 1 } - if [ ! -d ".venv" ]; then + if [ ! -d "${DATA_DIR}/.venv" ]; then echo "Python virtual environment not found. Please reinstall LangBot." > "${TRIM_TEMP_LOGFILE}" exit 1 fi @@ -79,6 +100,9 @@ case $1 in # fresh data/ with a default 5300 config gets recreated inside target/ on # every install/upgrade. The symlink keeps everything on the persistent # share; it is recreated here on each start (upgrades wipe target/). + # Requires the package user to have write permission on the app dir — if + # fnOS ever mounts it read-only, this fails loudly instead of silently + # running against an ephemeral data directory. APP_DATA="${APP_DIR}/data" if [ -L "${APP_DATA}" ]; then # already a symlink; re-point if the persistent dir changed @@ -87,7 +111,7 @@ case $1 in # legacy real dir (created by LangBot before this fix): merge into the # persistent dir without overwriting newer files already there mkdir -p "${DATA_DIR}" - cp -an "${APP_DATA}/." "${DATA_DIR}/" 2>/dev/null || cp -a "${APP_DATA}/." "${DATA_DIR}/" + cp -an "${APP_DATA}/." "${DATA_DIR}/." || cp -a "${APP_DATA}/." "${DATA_DIR}/" rm -rf "${APP_DATA}" ln -s "${DATA_DIR}" "${APP_DATA}" else @@ -146,7 +170,22 @@ except Exception: # (--standalone-runtime would require an external runtime at # ws://langbot_plugin_runtime:5400, which only exists in Docker Compose.) # --standalone-box omitted: Box sandbox defaults off, users enable via Web UI - nohup "${UV_BIN}" run --no-sync main.py \ + # + # Privilege model: the whole app runs as the generated package user + # (run-as: package, see config/privilege) — no root anywhere. HOME / + # UV_CACHE_DIR / UV_PROJECT_ENVIRONMENT are exported at the top and point + # at the persistent share so tool caches (uv, npm/npx) and the relocated + # venv stay writable regardless of the generated user's system home. + # + # Process-group lifecycle: setsid makes the main process a session/group + # leader, so PID == PGID. "stop" kills the whole group — stdio children + # (plugin runtime, Box) die with the parent and can never survive as + # orphans holding their ws ports after a crash, stop or upgrade. + command -v setsid >/dev/null 2>&1 || { + echo "setsid not found (util-linux required for process-group lifecycle)" > "${TRIM_TEMP_LOGFILE}" + exit 1 + } + setsid nohup "${UV_BIN}" run --no-sync main.py \ > "${LOG_FILE}" 2>&1 & echo $! > "${PID_FILE}" @@ -165,12 +204,18 @@ except Exception: if [ -f "${PID_FILE}" ]; then PID=$(cat "${PID_FILE}" | tr -d '[:space:]') if [ -n "${PID}" ]; then - kill "${PID}" 2>/dev/null + # Started with setsid, so PID == PGID: kill the whole group so + # stdio children (plugin runtime, Box) die with the parent. The + # plain-PID kill covers instances started before the setsid + # change (group kill is a no-op for them). + kill -TERM -- "-${PID}" 2>/dev/null + kill -TERM "${PID}" 2>/dev/null for _ in 1 2 3 4 5 6 7 8 9 10; do kill -0 "${PID}" 2>/dev/null || break sleep 1 done - kill -9 "${PID}" 2>/dev/null + kill -KILL -- "-${PID}" 2>/dev/null + kill -KILL "${PID}" 2>/dev/null fi rm -f "${PID_FILE}" fi diff --git a/packaging/fnos/cmd/uninstall_callback b/packaging/fnos/cmd/uninstall_callback index 2baee90d9..2c982078c 100755 --- a/packaging/fnos/cmd/uninstall_callback +++ b/packaging/fnos/cmd/uninstall_callback @@ -7,8 +7,7 @@ if [ "${wizard_keep_data:-yes}" = "no" ]; then # 应用运行数据(pid、日志等) if [ -n "${TRIM_PKGVAR}" ]; then rm -rf "${TRIM_PKGVAR:?}"/langbot.pid \ - "${TRIM_PKGVAR:?}"/langbot.log \ - "${TRIM_PKGVAR:?}"/bin 2>/dev/null || true + "${TRIM_PKGVAR:?}"/langbot.log 2>/dev/null || true fi # 共享数据目录(langbot/data) diff --git a/packaging/fnos/cmd/upgrade_callback b/packaging/fnos/cmd/upgrade_callback index 9288c866c..f5b51c508 100755 --- a/packaging/fnos/cmd/upgrade_callback +++ b/packaging/fnos/cmd/upgrade_callback @@ -8,6 +8,20 @@ APP_DIR="${TRIM_APPDEST}/langbot" DATA_DIR="${TRIM_DATA_SHARE_PATHS%%:*}" [ -z "${DATA_DIR}" ] && DATA_DIR="${TRIM_PKGVAR}/data" +# Writable HOME / caches / venv on the data share (must match +# cmd/install_callback — venv is relocated under DATA_DIR). +VENV_DIR="${DATA_DIR}/.venv" +export HOME="${DATA_DIR}/.home" +export UV_CACHE_DIR="${DATA_DIR}/.cache/uv" +# Never download a managed CPython — distro Python only (see install_callback) +export UV_PYTHON_DOWNLOADS=never +export UV_PROJECT_ENVIRONMENT="${VENV_DIR}" +mkdir -p "${HOME}" "${UV_CACHE_DIR}" "${DATA_DIR}/logs" || { + echo "Data share not writable: ${DATA_DIR}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} +DEBUG_LOG="${DATA_DIR}/logs/upgrade-debug.log" + # Apply port from upgrade wizard (config persists across upgrades; this # only rewrites it when the user changed the value in the upgrade wizard) CONFIG_FILE="${DATA_DIR}/config.yaml" @@ -26,52 +40,56 @@ cd "${APP_DIR}" || { exit 1 } -# Find uv (bundled first, then PATH / ~/.local/bin / ~/.cargo/bin) -UV_BIN="${TRIM_PKGVAR}/bin/uv" -if [ ! -x "${UV_BIN}" ]; then - UV_BIN="uv" -fi -if ! command -v "${UV_BIN}" >/dev/null 2>&1; then - UV_BIN="${HOME}/.local/bin/uv" -fi -if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then - UV_BIN="${HOME}/.cargo/bin/uv" -fi +# --- CPU architecture (must be resolved before locating bundled binaries) --- +ARCH=$(uname -m) +case "${ARCH}" in + x86_64|aarch64) ;; + *) + echo "Unsupported CPU architecture: ${ARCH}" > "${TRIM_TEMP_LOGFILE}" + exit 1 + ;; +esac -PYTHON_BIN="python3" -! command -v "${PYTHON_BIN}" >/dev/null 2>&1 && PYTHON_BIN="python" +# uv: bundled binary at its single canonical path in the app dir +UV_BIN="${TRIM_APPDEST}/bin/uv-${ARCH}" +[ -x "${UV_BIN}" ] || { + echo "Bundled uv binary missing or not executable: ${UV_BIN}" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} -# Re-sync deps -if [ -d ".venv" ]; then - "${UV_BIN}" sync --extra seekdb 2>/dev/null || { - echo "Dependency sync failed after upgrade" > "${TRIM_TEMP_LOGFILE}" +# Python: official python312 App Store app (same borrow pattern as Node.js) +PYTHON_BIN="/var/apps/python312/target/bin/python3" +[ -x "${PYTHON_BIN}" ] || { + echo "Python interpreter missing or not executable: ${PYTHON_BIN} (install the python312 app)" > "${TRIM_TEMP_LOGFILE}" + exit 1 +} + +# Re-sync deps (venv lives at ${DATA_DIR}/.venv, see install_callback) +if [ -d "${VENV_DIR}" ]; then + UV_PROJECT_ENVIRONMENT="${VENV_DIR}" "${UV_BIN}" sync --extra seekdb >> "${DEBUG_LOG}" 2>&1 || { + echo "Dependency sync failed after upgrade. See ${DEBUG_LOG}" > "${TRIM_TEMP_LOGFILE}" exit 1 } else - # Venv was lost, recreate via uv - if ! command -v "${UV_BIN}" >/dev/null 2>&1 && [ ! -x "${UV_BIN}" ]; then - "${PYTHON_BIN}" -m pip install --user --no-cache-dir uv 2>/dev/null || \ - "${PYTHON_BIN}" -m pip install --no-cache-dir uv 2>/dev/null || { - echo "Failed to install uv" > "${TRIM_TEMP_LOGFILE}" - exit 1 - } - export PATH="${HOME}/.local/bin:${PATH}" - UV_BIN="uv" - fi - "${UV_BIN}" venv .venv --python "${PYTHON_BIN}" || { - echo "Failed to recreate virtual environment" > "${TRIM_TEMP_LOGFILE}" + # Venv was lost, recreate it with the bundled uv on the data share + "${UV_BIN}" venv "${VENV_DIR}" --python "${PYTHON_BIN}" >> "${DEBUG_LOG}" 2>&1 || { + echo "Failed to recreate virtual environment. See ${DEBUG_LOG}" > "${TRIM_TEMP_LOGFILE}" exit 1 } - "${UV_BIN}" sync --extra seekdb || { - echo "Dependency sync failed" > "${TRIM_TEMP_LOGFILE}" + UV_PROJECT_ENVIRONMENT="${VENV_DIR}" "${UV_BIN}" sync --extra seekdb >> "${DEBUG_LOG}" 2>&1 || { + echo "Dependency sync failed. See ${DEBUG_LOG}" > "${TRIM_TEMP_LOGFILE}" exit 1 } fi # Verify frontend dist still present if [ ! -d "web/dist" ] || [ -z "$(ls -A web/dist 2>/dev/null)" ]; then - echo "Frontend dist missing after upgrade! Web UI will not be available." > "${TRIM_TEMP_LOGFILE}" + echo "Frontend dist missing! Web UI will not be available." > "${TRIM_TEMP_LOGFILE}" exit 1 fi +# Everything runs as the package user (run-as: package); the venv recreated +# above and the config rewritten via sed already belong to it. Cache HOME is +# prepared at the top of this script (see cmd/install_callback). + exit 0 diff --git a/packaging/fnos/cmd/upgrade_init b/packaging/fnos/cmd/upgrade_init index fdac71831..9fc46a9b9 100755 --- a/packaging/fnos/cmd/upgrade_init +++ b/packaging/fnos/cmd/upgrade_init @@ -1,20 +1,36 @@ #!/bin/bash -# cmd/upgrade_init - pre-upgrade hook -# Stop the running LangBot process before files are replaced. +# cmd/upgrade_init - pre-upgrade hook (runs before files are replaced) +# 1. Stop the running LangBot instance — whole process group, so stdio +# children (plugin runtime, Box) die with the parent. +# 2. Sweep stray processes from a previous crash/upgrade: orphans whose +# parent was hard-killed keep running and hold the runtime/box ws ports, +# which breaks the next start. PID_FILE="${TRIM_PKGVAR}/langbot.pid" if [ -f "${PID_FILE}" ]; then PID=$(cat "${PID_FILE}" | tr -d '[:space:]') if [ -n "${PID}" ] && kill -0 "${PID}" 2>/dev/null; then - kill "${PID}" 2>/dev/null + # Instances started with setsid have PID == PGID; the plain-PID kill + # covers instances started before that change. + kill -TERM -- "-${PID}" 2>/dev/null + kill -TERM "${PID}" 2>/dev/null for _ in 1 2 3 4 5 6 7 8 9 10; do kill -0 "${PID}" 2>/dev/null || break sleep 1 done - kill -9 "${PID}" 2>/dev/null + kill -KILL -- "-${PID}" 2>/dev/null + kill -KILL "${PID}" 2>/dev/null fi rm -f "${PID_FILE}" fi +# Stray sweep: match only our volume paths (venv/uv under the app dir on +# @appcenter, venv/HOME/caches under the data share on @appshare). This +# script itself runs from /var/apps//cmd/, so it never matches +# itself. No match is the normal case on a clean upgrade — pkill exits 1. +RUN_USER=$(id -un) +pkill -KILL -u "${RUN_USER}" -f "appcenter/langbot" 2>/dev/null || true +pkill -KILL -u "${RUN_USER}" -f "appshare/langbot" 2>/dev/null || true + exit 0 diff --git a/packaging/fnos/config/privilege b/packaging/fnos/config/privilege index e21db569b..2aafe2347 100644 --- a/packaging/fnos/config/privilege +++ b/packaging/fnos/config/privilege @@ -1,5 +1,5 @@ { "defaults": { - "run-as": "root" + "run-as": "package" } } diff --git a/packaging/fnos/manifest b/packaging/fnos/manifest index 98699b795..a05b5be1a 100644 --- a/packaging/fnos/manifest +++ b/packaging/fnos/manifest @@ -1,5 +1,5 @@ appname=langbot -version=4.10.10 +version=4.10.11 display_name=LangBot desc=基于 LLM 的多平台智能对话机器人,支持 QQ、微信、飞书、钉钉、Telegram 等十余种即时通讯平台,内置 Web 管理界面和 AI Agent 能力。 platform=all diff --git a/pyproject.toml b/pyproject.toml index 319ad8143..70b0a4196 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -71,7 +71,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.6.0b3", + "langbot-plugin==0.6.0b5", "asyncpg>=0.30.0", "line-bot-sdk>=3.19.0", "matrix-nio>=0.25.2", diff --git a/res/langbot-cloud.svg b/res/langbot-cloud.svg new file mode 100644 index 000000000..d114fb53b --- /dev/null +++ b/res/langbot-cloud.svg @@ -0,0 +1 @@ +Deploy on LangBot Cloud \ No newline at end of file diff --git a/skills/skills/langbot-testing/references/langrag-knowledge-base.md b/skills/skills/langbot-testing/references/langrag-knowledge-base.md index 00e0fb816..1d78dd82a 100644 --- a/skills/skills/langbot-testing/references/langrag-knowledge-base.md +++ b/skills/skills/langbot-testing/references/langrag-knowledge-base.md @@ -60,6 +60,19 @@ fixtures/rag/sentinel-doc.txt - Retrieve Test returns the uploaded document with the sentinel text. - Browser console has no unexpected errors. +## Document Lifecycle Acceptance + +- Keep the public Host file UUID from upload/listing when calling file deletion. Core stores the engine-returned `document_id` separately in the server-owned `engine_document_id` column and sends it to the engine; do not overwrite the Host UUID or put this mapping in creation settings. +- Wait for the ingestion task to finish before deletion. Pending/processing files reject deletion to avoid losing the tracking row while an upstream document is still being created. +- Verify both the Host file-list readback and absence of the sentinel upstream. Core only removes its row after an explicit engine `True`; `False`, missing configuration, runtime errors, and unconfirmed absence remain failures with the row retained. A connector's `False` is not proof that the upstream document is absent. +- Failed ingestion with an acknowledged engine ID retains it for cleanup. Historical failed files with no mapping still use the Host UUID fallback, subject to confirmed deletion. New unacknowledged/malformed responses are `interrupted`, not confirmed failures. +- `interrupted` means Core cannot establish the remote ingestion outcome (cancellation, disconnect, timeout, lost acknowledgement, or abandoned pending/processing work). It is not proof of failure, successful ingestion, or remote quiescence. Core preserves the tracking row, any known engine ID, and its source upload; retention cleanup also protects pending/processing/interrupted uploads. +- Runtime loading and normal file listing recover abandoned rows to `interrupted`. Reloading a KB object must preserve genuinely live tasks. Delayed old tasks cannot replay interrupted rows or overwrite them with completion; an observed late engine ID is still retained. +- File deletion and whole-KB deletion reject interrupted work with operator guidance. Do not bypass this guard merely because the task list is empty or the Host restarted: the old SDK/plugin action may still write. There is deliberately no automatic retry, force-delete, or claim of a remote cancellation fence. +- Recovery procedure: preserve a DB/storage backup; identify the exact Workspace, KB, Host file UUID, engine ID (if known), and installation; inspect that plugin/upstream operation; establish that the old operation has stopped or settled; then have an operator reconcile the confirmed upstream outcome and exact mapping before cleanup or re-upload. Never invent an opaque upstream ID or blindly set `failed` to unlock deletion. Lost historical uploads/IDs cannot be reconstructed by this change. Automatic cleanup of arbitrary connector orphans is outside this recovery contract. +- SDK ingest/delete action signatures and envelopes are unchanged; old SDKs use the same conservative interrupted fallback. Existing acknowledged asynchronous-engine responses retain the legacy Host completion semantics (Host ingestion acknowledgement, not a guarantee that the upstream index is ready). +- Migration `0025_rag_document_identity` leaves historical mappings null: it cannot reconstruct IDs previously discarded, nor restore already-deleted Host rows. Those require separately authorized investigation/recovery. Downgrading removes the mapping column and loses these identities; back up before rollback. + ## Local-Agent RAG Check After retrieval passes: diff --git a/src/langbot/pkg/api/http/controller/groups/plugins.py b/src/langbot/pkg/api/http/controller/groups/plugins.py index a338a3575..c7ed22efc 100644 --- a/src/langbot/pkg/api/http/controller/groups/plugins.py +++ b/src/langbot/pkg/api/http/controller/groups/plugins.py @@ -932,10 +932,13 @@ class PluginsRouterGroup(group.RouterGroup): return self.http_status(400, -1, 'file is required') file_bytes = file.read() + form = await quart.request.form + administrator_force = form.get('administrator_force', '').strip().lower() == 'true' execution_context = await self.ap.plugin_connector.require_workspace_context(request_context) data = { 'plugin_file': file_bytes, + 'administrator_force': administrator_force, } ctx = taskmgr.TaskContext.new() diff --git a/src/langbot/pkg/api/http/service/knowledge.py b/src/langbot/pkg/api/http/service/knowledge.py index 4f15ec1f6..e8c7f1dbe 100644 --- a/src/langbot/pkg/api/http/service/knowledge.py +++ b/src/langbot/pkg/api/http/service/knowledge.py @@ -339,6 +339,11 @@ class KnowledgeService: workspace_uuid = require_workspace_uuid(context) if await self.get_knowledge_base(context, kb_uuid) is None: raise WorkspaceNotFoundError('Knowledge base not found') + if isinstance(context, (RequestContext, ExecutionContext)): + execution_context = self._execution_context(context) + runtime_kb = await self.ap.rag_mgr.get_knowledge_base_by_uuid(execution_context, kb_uuid) + if runtime_kb is not None: + await runtime_kb.reconcile_interrupted_ingestions(execution_context) result = await self.ap.persistence_mgr.execute_async( sqlalchemy.select(persistence_rag.File) .where(persistence_rag.File.workspace_uuid == workspace_uuid) @@ -377,11 +382,20 @@ class KnowledgeService: kb_uuid: str, ) -> None: """删除知识库""" - workspace_uuid = require_workspace_uuid(context) + require_workspace_uuid(context) if await self.get_knowledge_base(context, kb_uuid) is None: raise WorkspaceNotFoundError('Knowledge base not found') - # delete files + execution_context = self._execution_context(context) + runtime_kb = await self.ap.rag_mgr.get_knowledge_base_by_uuid(execution_context, kb_uuid) + if runtime_kb is not None: + async with runtime_kb.ingestion_admission_lock: + await self._delete_knowledge_base_records(context, kb_uuid) + else: + await self._delete_knowledge_base_records(context, kb_uuid) + + async def _delete_knowledge_base_records(self, context: RequestContext | ExecutionContext, kb_uuid: str) -> None: + workspace_uuid = require_workspace_uuid(context) # NOTE: Chunk cleanup is for legacy (pre-plugin) KBs that stored chunks locally. # For plugin-based Knowledge Engines, the Chunk table is not populated, so this is a no-op. files = await self.ap.persistence_mgr.execute_async( @@ -389,6 +403,12 @@ class KnowledgeService: .where(persistence_rag.File.workspace_uuid == workspace_uuid) .where(persistence_rag.File.kb_id == kb_uuid) ) + files = files.all() + if any(file.status in {'pending', 'processing', 'interrupted'} for file in files): + raise RuntimeError( + 'Knowledge base has active or interrupted ingestion; retain its files and reconcile ' + 'plugin/upstream state before deleting the knowledge base.' + ) for file in files: # delete chunks await self.ap.persistence_mgr.execute_async( diff --git a/src/langbot/pkg/api/http/service/maintenance.py b/src/langbot/pkg/api/http/service/maintenance.py index 436670054..cdaf32091 100644 --- a/src/langbot/pkg/api/http/service/maintenance.py +++ b/src/langbot/pkg/api/http/service/maintenance.py @@ -13,6 +13,7 @@ import sqlalchemy from ....core import app from ....entity.persistence import bstorage as persistence_bstorage from ....entity.persistence import monitoring as persistence_monitoring +from ....entity.persistence import rag as persistence_rag from ..authz import WorkspaceRequiredError from ..context import ExecutionContext from .tenant import TenantContext, require_workspace_uuid @@ -330,6 +331,7 @@ class MaintenanceService: retention_days, True, ) + candidates = await self._exclude_ingestion_uploads(context, candidates) return await asyncio.to_thread( self._delete_local_candidates, candidates, @@ -340,6 +342,22 @@ class MaintenanceService: return 0 + async def _exclude_ingestion_uploads( + self, context: TenantContext, candidates: list[dict[str, Any]] + ) -> list[dict[str, Any]]: + """Keep source material for active or unacknowledged remote ingestion.""" + protected = set() + for offset in range(0, len(candidates), 500): + keys = [item['key'] for item in candidates[offset : offset + 500]] + result = await self.ap.persistence_mgr.execute_async( + sqlalchemy.select(persistence_rag.File.file_name) + .where(persistence_rag.File.workspace_uuid == require_workspace_uuid(context)) + .where(persistence_rag.File.status.in_(['pending', 'processing', 'interrupted'])) + .where(persistence_rag.File.file_name.in_(keys)) + ) + protected.update(result.scalars().all()) + return [item for item in candidates if item['key'] not in protected] + async def _expired_uploaded_candidates( self, context: TenantContext, @@ -363,6 +381,7 @@ class MaintenanceService: ) -> int: provider = self.ap.storage_mgr.storage_provider candidates = await self._expired_s3_upload_candidates(context, retention_days) + candidates = await self._exclude_ingestion_uploads(context, candidates) deleted = 0 for item in candidates: await provider.delete(item['key']) diff --git a/src/langbot/pkg/core/stages/load_config.py b/src/langbot/pkg/core/stages/load_config.py index f27cbe86f..bf76c3961 100644 --- a/src/langbot/pkg/core/stages/load_config.py +++ b/src/langbot/pkg/core/stages/load_config.py @@ -2,6 +2,7 @@ from __future__ import annotations import os import copy +import json from typing import Any from langbot.pkg.utils import bounded_executor, constants import yaml @@ -42,6 +43,7 @@ _RUNTIME_POLICY_DEFAULTS = { }, 'plugin': { 'connect_timeout_seconds': 180.0, + 'certification': {'trusted_public_keys': {}}, 'worker': { 'max_cpus': 1.0, 'max_memory_mb': 512, @@ -214,6 +216,30 @@ def _apply_env_overrides_to_config(cfg: dict) -> dict: return cfg +def _apply_certification_key_ring_env(cfg: dict) -> dict: + """Load the public certification key ring from one strict JSON env value. + + The generic environment override intentionally skips dictionaries. This + narrow exception keeps trusted issuer keys deployable without relying on a + mutable persisted config file, while rejecting malformed input instead of + silently running with an empty trust ring. + """ + raw = os.getenv('PLUGIN__CERTIFICATION__TRUSTED_PUBLIC_KEYS_JSON') + if raw is None: + return cfg + try: + key_ring = json.loads(raw) + except json.JSONDecodeError as exc: + raise ValueError('PLUGIN__CERTIFICATION__TRUSTED_PUBLIC_KEYS_JSON must be valid JSON') from exc + if not isinstance(key_ring, dict) or any( + not isinstance(key_id, str) or not key_id.strip() or not isinstance(key, str) or not key.strip() + for key_id, key in key_ring.items() + ): + raise ValueError('PLUGIN__CERTIFICATION__TRUSTED_PUBLIC_KEYS_JSON must be a non-empty string-to-string mapping') + cfg['plugin']['certification']['trusted_public_keys'] = key_ring + return cfg + + @stage.stage_class('LoadConfigStage') class LoadConfigStage(stage.BootingStage): """Load config file stage""" @@ -229,6 +255,7 @@ class LoadConfigStage(stage.BootingStage): # Apply environment variable overrides to data/config.yaml ap.instance_config.data = _apply_env_overrides_to_config(ap.instance_config.data) + ap.instance_config.data = _apply_certification_key_ring_env(ap.instance_config.data) blocking_config = ap.instance_config.data['system']['blocking_executor'] ap.blocking_executor = bounded_executor.configure_bounded_default_executor( diff --git a/src/langbot/pkg/entity/persistence/rag.py b/src/langbot/pkg/entity/persistence/rag.py index 29ce72177..8c1a6cf24 100644 --- a/src/langbot/pkg/entity/persistence/rag.py +++ b/src/langbot/pkg/entity/persistence/rag.py @@ -87,7 +87,11 @@ class File(Base): file_name = sqlalchemy.Column(sqlalchemy.String) extension = sqlalchemy.Column(sqlalchemy.String) created_at = sqlalchemy.Column(sqlalchemy.DateTime, default=sqlalchemy.func.now()) - status = sqlalchemy.Column(sqlalchemy.String, default='pending') # pending, processing, completed, failed + status = sqlalchemy.Column( + sqlalchemy.String, default='pending' + ) # pending, processing, completed, failed, interrupted + # Server-owned engine identity; the public Host file UUID never changes. + engine_document_id = sqlalchemy.Column(sqlalchemy.Text, nullable=True) __table_args__ = ( sqlalchemy.UniqueConstraint('workspace_uuid', 'uuid', name='uq_knowledge_base_files_workspace_uuid'), diff --git a/src/langbot/pkg/persistence/alembic/versions/0025_rag_document_identity.py b/src/langbot/pkg/persistence/alembic/versions/0025_rag_document_identity.py new file mode 100644 index 000000000..29f1d7959 --- /dev/null +++ b/src/langbot/pkg/persistence/alembic/versions/0025_rag_document_identity.py @@ -0,0 +1,28 @@ +"""Persist engine document identity separately from the public Host file UUID.""" + +from alembic import op +import sqlalchemy as sa + +revision = '0025_rag_document_identity' +down_revision = '0024_passkey_credentials' +branch_labels = None +depends_on = None + + +def upgrade() -> None: + inspector = sa.inspect(op.get_bind()) + if 'knowledge_base_files' not in inspector.get_table_names(): + return + columns = {column['name'] for column in inspector.get_columns('knowledge_base_files')} + if 'engine_document_id' not in columns: + # Legacy upstream IDs cannot be inferred from Host UUIDs or user config. + op.add_column('knowledge_base_files', sa.Column('engine_document_id', sa.Text(), nullable=True)) + + +def downgrade() -> None: + inspector = sa.inspect(op.get_bind()) + if 'knowledge_base_files' not in inspector.get_table_names(): + return + columns = {column['name'] for column in inspector.get_columns('knowledge_base_files')} + if 'engine_document_id' in columns: + op.drop_column('knowledge_base_files', 'engine_document_id') diff --git a/src/langbot/pkg/persistence/alembic/versions/0029_merge_rag_identity.py b/src/langbot/pkg/persistence/alembic/versions/0029_merge_rag_identity.py new file mode 100644 index 000000000..2405dfb38 --- /dev/null +++ b/src/langbot/pkg/persistence/alembic/versions/0029_merge_rag_identity.py @@ -0,0 +1,15 @@ +"""Join document identity with the 4.11 manual migration and draft branches.""" + +revision = '0029_merge_rag_identity' +down_revision = ('0028_merge_knowledge_drafts', '0025_rag_document_identity') +branch_labels = None +depends_on = None + + +def upgrade() -> None: + pass + + +def downgrade() -> None: + # Unmerge the graph only; retain both branches' schemas and user data. + pass diff --git a/src/langbot/pkg/plugin/archive.py b/src/langbot/pkg/plugin/archive.py index f3a8511a4..28b23cbff 100644 --- a/src/langbot/pkg/plugin/archive.py +++ b/src/langbot/pkg/plugin/archive.py @@ -1,7 +1,10 @@ from __future__ import annotations +import hashlib import io import zipfile +from dataclasses import dataclass +from enum import Enum import yaml @@ -14,6 +17,34 @@ _PLUGIN_METADATA_MAX_BYTES = 1024 * 1024 _PLUGIN_REQUIREMENTS_MAX_ENTRIES = 1000 +class ArchiveCertificateState(str, Enum): + """Syntactic certificate declaration state; this is not verification.""" + + ABSENT = 'absent' + MALFORMED = 'malformed' + DECLARED = 'declared' + + +@dataclass(frozen=True) +class ArchiveCertificateDeclaration: + """Bounded, verifier-facing certificate declaration from ``manifest.yaml``.""" + + state: ArchiveCertificateState + runtime_profile: str | None = None + payload: dict[str, object] | None = None + + +@dataclass(frozen=True) +class PluginArchiveInspection: + """Validated archive metadata plus unverified certificate declaration facts.""" + + manifest: dict + requirements: list[str] + names: list[str] + artifact_digest: str + certificate: ArchiveCertificateDeclaration + + def _read_plugin_archive_member( archive: zipfile.ZipFile, member: zipfile.ZipInfo, @@ -29,12 +60,30 @@ def _read_plugin_archive_member( return content -def inspect_plugin_archive_metadata( - file_bytes: bytes, - *, - require_manifest: bool = True, -) -> tuple[dict, list[str], list[str]]: - """Validate archive size metadata and read only bounded preview fields.""" +def _inspect_certificate_declaration(manifest: dict) -> ArchiveCertificateDeclaration: + declaration = manifest.get('certification') + if declaration is None: + return ArchiveCertificateDeclaration(ArchiveCertificateState.ABSENT) + if not isinstance(declaration, dict): + return ArchiveCertificateDeclaration(ArchiveCertificateState.MALFORMED) + + runtime_profile = declaration.get('runtime_profile') + payload = declaration.get('certificate') + if not isinstance(runtime_profile, str) or not runtime_profile.strip() or not isinstance(payload, dict): + return ArchiveCertificateDeclaration(ArchiveCertificateState.MALFORMED) + return ArchiveCertificateDeclaration( + ArchiveCertificateState.DECLARED, + runtime_profile=runtime_profile, + payload=payload, + ) + + +def inspect_plugin_archive(file_bytes: bytes, *, require_manifest: bool = True) -> PluginArchiveInspection: + """Validate an archive and expose certificate declaration facts for a verifier. + + Certificate signatures and issuer trust are deliberately not evaluated here; + callers must pass the declaration and artifact digest to an SDK verifier. + """ with zipfile.ZipFile(io.BytesIO(file_bytes)) as archive: members = archive.infolist() @@ -95,4 +144,25 @@ def inspect_plugin_archive_metadata( for line in content.splitlines() if line.strip() and not line.strip().startswith('#') ][:_PLUGIN_REQUIREMENTS_MAX_ENTRIES] - return manifest, requirements, names + + return PluginArchiveInspection( + manifest=manifest, + requirements=requirements, + names=names, + artifact_digest=hashlib.sha256(file_bytes).hexdigest(), + certificate=_inspect_certificate_declaration(manifest), + ) + + +def inspect_plugin_archive_metadata( + file_bytes: bytes, + *, + require_manifest: bool = True, +) -> tuple[dict, list[str], list[str]]: + """Legacy tuple API for archive metadata callers. + + Use ``inspect_plugin_archive`` when certificate declaration facts are needed. + """ + + inspection = inspect_plugin_archive(file_bytes, require_manifest=require_manifest) + return inspection.manifest, inspection.requirements, inspection.names diff --git a/src/langbot/pkg/plugin/certification.py b/src/langbot/pkg/plugin/certification.py new file mode 100644 index 000000000..d1a233847 --- /dev/null +++ b/src/langbot/pkg/plugin/certification.py @@ -0,0 +1,234 @@ +"""Pure certified-plugin facts and admission policies. + +This module intentionally does not verify signatures. An SDK-backed verifier +must produce ``CertificateFacts`` from an inspected archive before admission. +""" + +from __future__ import annotations + +import base64 +from collections.abc import Callable, Mapping +from dataclasses import dataclass +from enum import Enum + +from cryptography.exceptions import InvalidSignature +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PublicKey +from langbot_plugin.certification import normalized_zip_digest, verify_archive + + +SHARED_RUNTIME_V1 = 'shared-runtime-v1' +DEDICATED_RUNTIME = 'dedicated' + + +class CertificateVerification(str, Enum): + ABSENT = 'absent' + MALFORMED = 'malformed' + INVALID = 'invalid' + VALID = 'valid' + + +class DeploymentMode(str, Enum): + CLOUD = 'cloud' + OSS = 'oss' + + +class AdmissionDisposition(str, Enum): + SHARED_ELIGIBLE = 'shared_eligible' + DEDICATED_ALLOWED = 'dedicated_allowed' + REJECTED = 'rejected' + ADMINISTRATOR_FORCE_REQUIRED = 'administrator_force_required' + + +class AdmissionCode(str, Enum): + SHARED_ELIGIBLE = 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE' + CLOUD_CERTIFICATE_REQUIRED = 'CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_REQUIRED' + CLOUD_CERTIFICATE_INVALID = 'CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_INVALID' + OSS_LEGACY_DEDICATED = 'CERTIFIED_PLUGIN_OSS_LEGACY_DEDICATED' + OSS_FORCE_REQUIRED = 'CERTIFIED_PLUGIN_OSS_FORCE_REQUIRED' + OSS_FORCED_DEDICATED = 'CERTIFIED_PLUGIN_OSS_FORCED_DEDICATED' + OSS_CERTIFIED_DEDICATED = 'CERTIFIED_PLUGIN_OSS_CERTIFIED_DEDICATED' + + +class PluginLogVisibility(str, Enum): + TENANT_SCOPED = 'tenant_scoped' + DETAILED_PROCESS = 'detailed_process' + + +@dataclass(frozen=True) +class CertificateFacts: + """Certificate result supplied by an archive verifier. + + ``VALID`` means the verifier has validated both the certificate and its + binding to the immutable artifact digest in ``PluginCertificationFacts``. + """ + + verification: CertificateVerification + runtime_profile: str | None = None + certificate_id: str | None = None + + @property + def is_valid_shared_runtime(self) -> bool: + return self.verification is CertificateVerification.VALID and self.runtime_profile == SHARED_RUNTIME_V1 + + @property + def is_declared(self) -> bool: + return self.verification is not CertificateVerification.ABSENT + + +@dataclass(frozen=True) +class PluginCertificationFacts: + """Immutable Core-side facts for one plugin installation artifact.""" + + installation_uuid: str + artifact_digest: str + certificate: CertificateFacts + + def __post_init__(self) -> None: + if len(self.artifact_digest) != 64 or any( + character not in '0123456789abcdef' for character in self.artifact_digest.lower() + ): + raise ValueError('artifact_digest must be a lowercase-or-uppercase SHA-256 hex digest') + + +@dataclass(frozen=True) +class VerifiedArchiveCertificate: + """SDK verification facts bound to the comment-normalized ZIP digest.""" + + normalized_digest: str + certificate: CertificateFacts + + def for_installation(self, installation_uuid: str) -> PluginCertificationFacts: + return PluginCertificationFacts( + installation_uuid=installation_uuid, + artifact_digest=self.normalized_digest, + certificate=self.certificate, + ) + + +def trusted_public_key_ring(config: object) -> dict[str, Callable[[bytes, bytes], bool]]: + """Build the non-secret Ed25519 verifier ring from instance configuration. + + ``plugin.certification.trusted_public_keys`` is a mapping of key IDs to + standard base64-encoded 32-byte Ed25519 public keys. Configuration errors + are explicit so an operator never silently gets a weaker trust policy. + """ + + if config is None: + return {} + if not isinstance(config, Mapping): + raise ValueError('plugin.certification.trusted_public_keys must be a mapping') + + ring: dict[str, Callable[[bytes, bytes], bool]] = {} + for raw_key_id, raw_public_key in config.items(): + key_id = str(raw_key_id).strip() + if not key_id or not isinstance(raw_public_key, str): + raise ValueError('plugin.certification.trusted_public_keys entries must have string IDs and values') + try: + public_key_bytes = base64.b64decode(raw_public_key.encode('ascii'), validate=True) + public_key = Ed25519PublicKey.from_public_bytes(public_key_bytes) + except (UnicodeEncodeError, ValueError) as exc: + raise ValueError(f'plugin.certification trusted public key {key_id!r} is invalid') from exc + + def verify(payload: bytes, signature: bytes, *, verifier: Ed25519PublicKey = public_key) -> bool: + try: + verifier.verify(signature, payload) + except (InvalidSignature, TypeError, ValueError): + return False + return True + + ring[key_id] = verify + return ring + + +def verify_plugin_archive_certificate( + archive: bytes, + *, + trusted_public_keys: object, +) -> VerifiedArchiveCertificate: + """Use the SDK ZIP-comment API and retain its normalized-digest binding.""" + + verification = verify_archive(archive, trusted_public_key_ring(trusted_public_keys).get) + envelope = verification.envelope + runtime_profile = envelope.shared_runtime if envelope is not None else None + certificate_id = envelope.key_id if envelope is not None else None + state = { + 'absent': CertificateVerification.ABSENT, + 'malformed': CertificateVerification.MALFORMED, + 'valid': CertificateVerification.VALID, + }.get(verification.status, CertificateVerification.INVALID) + return VerifiedArchiveCertificate( + normalized_digest=normalized_zip_digest(archive), + certificate=CertificateFacts( + verification=state, + runtime_profile=runtime_profile, + certificate_id=certificate_id, + ), + ) + + +@dataclass(frozen=True) +class PluginAdmissionDecision: + disposition: AdmissionDisposition + code: AdmissionCode + runtime_profile: str + + +def decide_plugin_admission( + *, + deployment: DeploymentMode | str, + facts: PluginCertificationFacts, + administrator_force: bool = False, +) -> PluginAdmissionDecision: + """Apply Cloud fail-closed and OSS administrator-force admission rules.""" + + mode = DeploymentMode(deployment) + certificate = facts.certificate + if certificate.is_valid_shared_runtime: + return PluginAdmissionDecision( + AdmissionDisposition.SHARED_ELIGIBLE, + AdmissionCode.SHARED_ELIGIBLE, + SHARED_RUNTIME_V1, + ) + + if mode is DeploymentMode.CLOUD: + code = ( + AdmissionCode.CLOUD_CERTIFICATE_REQUIRED + if certificate.verification is CertificateVerification.ABSENT + else AdmissionCode.CLOUD_CERTIFICATE_INVALID + ) + return PluginAdmissionDecision(AdmissionDisposition.REJECTED, code, DEDICATED_RUNTIME) + + if certificate.verification is CertificateVerification.ABSENT: + return PluginAdmissionDecision( + AdmissionDisposition.DEDICATED_ALLOWED, + AdmissionCode.OSS_LEGACY_DEDICATED, + DEDICATED_RUNTIME, + ) + + if certificate.verification is CertificateVerification.VALID: + return PluginAdmissionDecision( + AdmissionDisposition.DEDICATED_ALLOWED, + AdmissionCode.OSS_CERTIFIED_DEDICATED, + DEDICATED_RUNTIME, + ) + + if administrator_force: + return PluginAdmissionDecision( + AdmissionDisposition.DEDICATED_ALLOWED, + AdmissionCode.OSS_FORCED_DEDICATED, + DEDICATED_RUNTIME, + ) + + return PluginAdmissionDecision( + AdmissionDisposition.ADMINISTRATOR_FORCE_REQUIRED, + AdmissionCode.OSS_FORCE_REQUIRED, + DEDICATED_RUNTIME, + ) + + +def decide_plugin_log_visibility(facts: PluginCertificationFacts) -> PluginLogVisibility: + """Select the minimum log visibility compatible with a verified shared runtime.""" + + if facts.certificate.is_valid_shared_runtime: + return PluginLogVisibility.TENANT_SCOPED + return PluginLogVisibility.DETAILED_PROCESS diff --git a/src/langbot/pkg/plugin/connector.py b/src/langbot/pkg/plugin/connector.py index ee8ce2dd5..99c072dc8 100644 --- a/src/langbot/pkg/plugin/connector.py +++ b/src/langbot/pkg/plugin/connector.py @@ -29,6 +29,13 @@ from .errors import ( MarketplacePluginVersionNotFoundError, ) from .archive import inspect_plugin_archive_metadata +from .certification import ( + AdmissionDisposition, + PluginCertificationFacts, + VerifiedArchiveCertificate, + decide_plugin_admission, + verify_plugin_archive_certificate, +) from .github import ( validate_github_plugin_install_info, validate_github_release_asset_url, @@ -159,6 +166,30 @@ def _decode_json_object(body: bytes, *, subject: str) -> dict[str, Any]: return payload +def _select_marketplace_plugin_version( + versions: Any, + *, + requested_version: str | None, + plugin_author: str, + plugin_name: str, +) -> str: + if not isinstance(versions, list) or not versions: + raise ValueError(f'Plugin {plugin_author}/{plugin_name} has no versions') + + if requested_version is None: + candidate = versions[0] + if not isinstance(candidate, dict) or not candidate.get('version'): + raise ValueError(f'Plugin {plugin_author}/{plugin_name} has no versions') + return str(candidate['version']) + + for candidate in versions: + if isinstance(candidate, dict) and str(candidate.get('version') or '') == requested_version: + return requested_version + raise ValueError( + f'Plugin {plugin_author}/{plugin_name} version {requested_version} is not available in marketplace' + ) + + class PluginRuntimeConnector(ManagedRuntimeConnector): """Plugin runtime connector""" @@ -1718,20 +1749,57 @@ class PluginRuntimeConnector(ManagedRuntimeConnector): subject='Marketplace plugin versions', ) versions = versions_payload.get('data', {}).get('versions', []) - if ( - not isinstance(versions, list) - or not versions - or not isinstance(versions[0], dict) - or not versions[0].get('version') - ): - raise ValueError(f'Plugin {plugin_author}/{plugin_name} has no versions') - latest_version = str(versions[0]['version']) + version = _select_marketplace_plugin_version( + versions, + requested_version=None, + plugin_author=plugin_author, + plugin_name=plugin_name, + ) _download_status, plugin_package = await _marketplace_get( client, - f'{space_url}/api/v1/marketplace/plugins/download/{plugin_author}/{plugin_name}/{latest_version}', + f'{space_url}/api/v1/marketplace/plugins/download/{plugin_author}/{plugin_name}/{version}', max_bytes=_MARKETPLACE_PLUGIN_DOWNLOAD_MAX_BYTES, ) - return plugin_package, latest_version + return plugin_package, version + + def _admit_plugin_archive( + self, + file_bytes: bytes, + install_info: dict[str, Any], + ) -> tuple[dict[str, Any], VerifiedArchiveCertificate]: + """Verify and admit one archive before it can reach durable storage or Runtime.""" + + certification_config = self.ap.instance_config.data.get('plugin', {}).get('certification', {}) + if not isinstance(certification_config, dict): + raise ValueError('plugin.certification must be a mapping') + verified = verify_plugin_archive_certificate( + file_bytes, + trusted_public_keys=certification_config.get('trusted_public_keys', {}), + ) + facts = PluginCertificationFacts( + installation_uuid='pending-installation', + artifact_digest=verified.normalized_digest, + certificate=verified.certificate, + ) + decision = decide_plugin_admission( + deployment=getattr(getattr(self.ap, 'deployment', None), 'mode', 'oss'), + facts=facts, + administrator_force=install_info.get('administrator_force') is True, + ) + if decision.disposition not in { + AdmissionDisposition.DEDICATED_ALLOWED, + AdmissionDisposition.SHARED_ELIGIBLE, + }: + raise ValueError(decision.code.value) + certification_info = { + 'normalized_digest': facts.artifact_digest, + 'verification': facts.certificate.verification.value, + 'certificate_runtime_profile': facts.certificate.runtime_profile, + 'certificate_id': facts.certificate.certificate_id, + 'runtime_profile': decision.runtime_profile, + 'admission_code': decision.code.value, + } + return {**install_info, '_certification': certification_info}, verified @diagnostics.observe('lifecycle', 'runtime.install_plugin', source='runtime', stage='execute') async def install_plugin( @@ -1778,6 +1846,7 @@ class PluginRuntimeConnector(ManagedRuntimeConnector): if task_context is not None: task_context.set_current_action('validating plugin package') task_context.metadata['progress_percent'] = 32 + install_info, verified_certificate = self._admit_plugin_archive(file_bytes, install_info) manifest_author, manifest_name = self._inspect_plugin_package(file_bytes, task_context) if not manifest_author or not manifest_name: raise ValueError('Plugin package manifest identity is missing') @@ -1818,6 +1887,9 @@ class PluginRuntimeConnector(ManagedRuntimeConnector): except Exception: await self._delete_artifact_if_unreferenced(execution_context, artifact_digest) raise + certification_facts = verified_certificate.for_installation(binding.installation_uuid) + if certification_facts.artifact_digest != install_info['_certification']['normalized_digest']: + raise RuntimeError('Plugin certification digest changed before Runtime apply') if not previous_was_durable and self.runtime_profile == 'oss_dev': bridge = self._legacy_oss_bridge_binding(execution_context) try: diff --git a/src/langbot/pkg/rag/knowledge/kbmgr.py b/src/langbot/pkg/rag/knowledge/kbmgr.py index 1c992103b..4ef1c9d0c 100644 --- a/src/langbot/pkg/rag/knowledge/kbmgr.py +++ b/src/langbot/pkg/rag/knowledge/kbmgr.py @@ -1,6 +1,7 @@ from __future__ import annotations import asyncio import io +import inspect import mimetypes import os.path import traceback @@ -43,9 +44,42 @@ class RuntimeKnowledgeBase(KnowledgeBaseInterface): super().__init__(ap) self.knowledge_base_entity = knowledge_base_entity self.execution_context = execution_context + # Shared across KB object reloads, but deliberately not a remote-work fence. + self._ingestion_tasks = ap.__dict__.setdefault('_knowledge_ingestion_tasks', {}) + locks = ap.__dict__.setdefault('_knowledge_ingestion_locks', {}) + self.ingestion_admission_lock = locks.setdefault( + (execution_context.workspace_uuid, self.get_uuid()), asyncio.Lock() + ) async def initialize(self): - pass + await self.reconcile_interrupted_ingestions(self.execution_context) + + async def reconcile_interrupted_ingestions(self, execution_context: ExecutionContext) -> None: + """Persist loss of Host observation, never infer remote quiescence. + + A restarted Host cannot observe the old SDK request's outcome. Retain its + row, upload and identity for operator reconciliation; do not retry/delete. + The shared admission lock protects newly queued tasks during KB reloads. + """ + async with self.ingestion_admission_lock: + live_ids = [ + key[2] + for key, task in self._ingestion_tasks.items() + if key[:2] == (execution_context.workspace_uuid, self.get_uuid()) and not task.done() + ] + + async def reconcile(): + await self._assert_execution_context(execution_context) + await self.ap.persistence_mgr.execute_async( + sqlalchemy.update(persistence_rag.File) + .where(persistence_rag.File.workspace_uuid == execution_context.workspace_uuid) + .where(persistence_rag.File.kb_id == self.get_uuid()) + .where(persistence_rag.File.status.in_(['pending', 'processing'])) + .where(persistence_rag.File.uuid.not_in(live_ids)) + .values(status='interrupted') + ) + + await run_in_workspace_uow(self.ap, execution_context.workspace_uuid, reconcile) async def _assert_execution_context(self, execution_context: ExecutionContext) -> None: """Reject stale or cross-Workspace runtime access.""" @@ -99,23 +133,38 @@ class RuntimeKnowledgeBase(KnowledgeBaseInterface): task_context: taskmgr.TaskContext, parser_plugin_id: str | None = None, ): - await run_in_workspace_uow( - self.ap, - execution_context.workspace_uuid, - lambda: self._assert_execution_context(execution_context), - ) - self._require_upload_object_key(execution_context, file.file_name) + key = (execution_context.workspace_uuid, self.get_uuid(), file.uuid) + current_task = asyncio.current_task() + existing = self._ingestion_tasks.get(key) + if existing is not None and existing is not current_task and not existing.done(): + raise RuntimeError('Knowledge file ingestion is already running') + self._ingestion_tasks[key] = current_task + engine_document_id = None + dispatched = False + confirmed_failure = False + cleanup_upload = False + status_visible = False try: - # set file status to processing + await run_in_workspace_uow( + self.ap, + execution_context.workspace_uuid, + lambda: self._assert_execution_context(execution_context), + ) + self._require_upload_object_key(execution_context, file.file_name) + # Claim only pending work. Recovered/terminal rows cannot be replayed. status_visible = False for retry_delay in (0.0, 0.01, 0.05, 0.1): if retry_delay: await asyncio.sleep(retry_delay) - if await self._set_file_status(execution_context, file.uuid, 'processing'): + async with self.ingestion_admission_lock: + claimed = await self._set_file_status( + execution_context, file.uuid, 'processing', expected_statuses=('pending',) + ) + if claimed: status_visible = True break if not status_visible: - raise WorkspaceNotFoundError('Knowledge file was not committed before its background task started') + raise WorkspaceNotFoundError('Knowledge file is missing, already claimed, or interrupted') task_context.set_current_action('Processing file') @@ -146,9 +195,12 @@ class RuntimeKnowledgeBase(KnowledgeBaseInterface): 'metadata': {}, } await self._require_plugin_runtime_context(execution_context) + dispatched = True parsed_content = await self.ap.plugin_connector.call_parser(parser_plugin_id, parse_context, file_bytes) + dispatched = False - # Call plugin to ingest document + # From dispatch until a valid response, failure is an unknown outcome. + dispatched = True result = await self._ingest_document( execution_context, { @@ -162,57 +214,109 @@ class RuntimeKnowledgeBase(KnowledgeBaseInterface): parsed_content=parsed_content, ) - # Check plugin result status + # Failed ingestion can still have created an upstream document (for + # example, upload succeeded but parsing failed). Retain that identity + # for cleanup too. Never coerce or normalize an opaque engine ID. + returned_id = result.get('document_id') + if isinstance(returned_id, str) and returned_id.strip(): + engine_document_id = returned_id if result.get('status') == 'failed': + confirmed_failure = engine_document_id is not None error_msg = result.get('error_message', 'Plugin ingestion returned failed status') raise Exception(error_msg) + if engine_document_id is None: + raise ValueError('Plugin ingestion must return a nonempty string document_id') - # set file status to completed - if not await self._set_file_status(execution_context, file.uuid, 'completed'): - raise WorkspaceNotFoundError('Knowledge file not found') + # Commit the identity and completion together, never a status-only + # success that loses the only way to address the upstream document. + if not await self._set_file_status( + execution_context, + file.uuid, + 'completed', + engine_document_id=engine_document_id, + expected_statuses=('processing',), + ): + raise WorkspaceNotFoundError('Knowledge file not found or ingestion interrupted') + cleanup_upload = True - except Exception as e: - self.ap.logger.error(f'Error storing file {file.uuid}: {e}') - traceback.print_exc() - # A stale placement is fenced from all writes, including failure - # status updates from an old background task. + except (Exception, asyncio.CancelledError) as e: + cancelled = isinstance(e, asyncio.CancelledError) + status = 'interrupted' if cancelled or (dispatched and not confirmed_failure) else 'failed' + self.ap.logger.warning(f'Knowledge ingestion {status} for file {file.uuid}') + # A cancelled RPC waiter or transport error cannot establish remote + # failure. Preserve any returned ID and never downgrade a committed + # completion after an ambiguous commit acknowledgement. try: - if not await self._set_file_status(execution_context, file.uuid, 'failed'): - raise WorkspaceNotFoundError('Knowledge file not found') + if status_visible or cancelled: + changed = await self._set_file_status( + execution_context, + file.uuid, + status, + engine_document_id=engine_document_id, + expected_statuses=('pending', 'processing') if cancelled else ('processing',), + ) + cleanup_upload = changed and status == 'failed' except Exception: self.ap.logger.warning(f'Skipping stale RAG task status update for file {file.uuid}') - raise finally: - # An old background task must not touch an upload after its - # placement generation has been fenced off. - try: - await self._assert_execution_context(execution_context) - await self.ap.storage_mgr.delete_scoped_object_key( - execution_context, - file.file_name, - expected_owner_type='upload_document', - ) - except (WorkspaceRequiredError, WorkspaceNotFoundError): - self.ap.logger.warning(f'Skipping stale RAG upload cleanup for file {file.uuid}') + if self._ingestion_tasks.get(key) is current_task: + self._ingestion_tasks.pop(key, None) + # Only release recovery material after an acknowledged terminal write. + if cleanup_upload: + try: + await run_in_workspace_uow( + self.ap, + execution_context.workspace_uuid, + lambda: self._assert_execution_context(execution_context), + ) + await self.ap.storage_mgr.delete_scoped_object_key( + execution_context, + file.file_name, + expected_owner_type='upload_document', + ) + except (WorkspaceRequiredError, WorkspaceNotFoundError): + self.ap.logger.warning(f'Skipping stale RAG upload cleanup for file {file.uuid}') async def _set_file_status( self, execution_context: ExecutionContext, file_uuid: str, status: str, + *, + engine_document_id: str | None = None, + expected_statuses: tuple[str, ...] | None = None, ) -> bool: - """Commit one detached-task status transition in its own tenant UoW.""" + """Commit one detached-task status/identity transition in its tenant UoW.""" async def update() -> bool: await self._assert_execution_context(execution_context) - result = await self.ap.persistence_mgr.execute_async( + values = {'status': status} + if engine_document_id is not None: + values['engine_document_id'] = engine_document_id + statement = ( sqlalchemy.update(persistence_rag.File) .where(persistence_rag.File.workspace_uuid == execution_context.workspace_uuid) + .where(persistence_rag.File.kb_id == self.knowledge_base_entity.uuid) .where(persistence_rag.File.uuid == file_uuid) - .values(status=status) + .values(**values) ) - return getattr(result, 'rowcount', 0) > 0 + if expected_statuses is not None: + statement = statement.where(persistence_rag.File.status.in_(expected_statuses)) + result = await self.ap.persistence_mgr.execute_async(statement) + changed = getattr(result, 'rowcount', 0) > 0 + if not changed and engine_document_id is not None: + # A different Host may already have recovered observation. Save + # the late identity without claiming that the attempt completed. + await self.ap.persistence_mgr.execute_async( + sqlalchemy.update(persistence_rag.File) + .where(persistence_rag.File.workspace_uuid == execution_context.workspace_uuid) + .where(persistence_rag.File.kb_id == self.get_uuid()) + .where(persistence_rag.File.uuid == file_uuid) + .where(persistence_rag.File.status == 'interrupted') + .values(engine_document_id=engine_document_id) + ) + return changed persistence_mgr = self.ap.persistence_mgr managed_mode = getattr(getattr(persistence_mgr, 'mode', None), 'value', None) in { @@ -264,33 +368,43 @@ class RuntimeKnowledgeBase(KnowledgeBaseInterface): file_obj = persistence_rag.File(**file_obj_data) - await self.ap.persistence_mgr.execute_async(sqlalchemy.insert(persistence_rag.File).values(file_obj_data)) + # Serialize admission with reconciliation, not with remote ingestion. + async with self.ingestion_admission_lock: + await self.ap.persistence_mgr.execute_async(sqlalchemy.insert(persistence_rag.File).values(file_obj_data)) + ctx = taskmgr.TaskContext.new() + coroutine = self._store_file_task( + execution_context, file_obj, task_context=ctx, parser_plugin_id=parser_plugin_id + ) + try: + wrapper = self.ap.task_mgr.create_user_task( + coroutine, + kind='knowledge-operation', + name=f'knowledge-store-file-{file_id}', + label=f'Store file {file_id}', + context=ctx, + instance_uuid=execution_context.instance_uuid, + workspace_uuid=execution_context.workspace_uuid, + placement_generation=execution_context.placement_generation, + ) + except taskmgr.TaskCapacityError: + await self.ap.persistence_mgr.execute_async( + sqlalchemy.delete(persistence_rag.File) + .where(persistence_rag.File.workspace_uuid == execution_context.workspace_uuid) + .where(persistence_rag.File.uuid == file_uuid) + ) + raise + key = (execution_context.workspace_uuid, kb_id, file_uuid) + self._ingestion_tasks[key] = wrapper.task - # run background task asynchronously - ctx = taskmgr.TaskContext.new() - try: - wrapper = self.ap.task_mgr.create_user_task( - self._store_file_task( - execution_context, - file_obj, - task_context=ctx, - parser_plugin_id=parser_plugin_id, - ), - kind='knowledge-operation', - name=f'knowledge-store-file-{file_id}', - label=f'Store file {file_id}', - context=ctx, - instance_uuid=execution_context.instance_uuid, - workspace_uuid=execution_context.workspace_uuid, - placement_generation=execution_context.placement_generation, - ) - except taskmgr.TaskCapacityError: - await self.ap.persistence_mgr.execute_async( - sqlalchemy.delete(persistence_rag.File) - .where(persistence_rag.File.workspace_uuid == execution_context.workspace_uuid) - .where(persistence_rag.File.uuid == file_uuid) - ) - raise + def forget_task(task): + # The task manager wraps the coroutine; pre-start cancellation + # need not enter that wrapper's finally block. + if inspect.getcoroutinestate(coroutine) == inspect.CORO_CREATED: + coroutine.close() + if self._ingestion_tasks.get(key) is task: + self._ingestion_tasks.pop(key, None) + + wrapper.task.add_done_callback(forget_task) return wrapper.id async def _store_zip_file( @@ -445,17 +559,37 @@ class RuntimeKnowledgeBase(KnowledgeBaseInterface): async def delete_file(self, execution_context: ExecutionContext, file_id: str): await self._assert_execution_context(execution_context) result = await self.ap.persistence_mgr.execute_async( - sqlalchemy.select(persistence_rag.File.uuid) + sqlalchemy.select(persistence_rag.File) .where(persistence_rag.File.workspace_uuid == execution_context.workspace_uuid) .where(persistence_rag.File.kb_id == self.knowledge_base_entity.uuid) .where(persistence_rag.File.uuid == file_id) .limit(1) ) - if result.first() is None: + file = result.first() + if file is None: raise WorkspaceNotFoundError('Knowledge file not found') - await self._delete_document(execution_context, file_id) + # Pending/processing tasks may still create an upstream document. Do not + # discard their tracking row or race a delete against that creation. + if file.status in {'pending', 'processing'}: + raise RuntimeError( + 'Cannot delete a file while ingestion is pending or processing; wait for the task to finish' + ) + if file.status == 'interrupted': + raise RuntimeError( + 'Knowledge ingestion was interrupted; its remote outcome is unknown. ' + 'The file, upload and known engine identity were retained. ' + 'Check plugin/upstream state and stop or settle the old ingestion before operator reconciliation; ' + 'automatic deletion or re-ingestion is unsafe.' + ) + document_id = file.engine_document_id if file.engine_document_id is not None else file_id + if await self._delete_document(execution_context, document_id) is not True: + raise RuntimeError( + 'Knowledge engine did not confirm document deletion; the file was retained. ' + 'Check the engine configuration and plugin logs before retrying.' + ) - # Also cleanup DB record + # The plugin call may outlive the original placement generation. + await self._assert_execution_context(execution_context) await self.ap.persistence_mgr.execute_async( sqlalchemy.delete(persistence_rag.File) .where(persistence_rag.File.workspace_uuid == execution_context.workspace_uuid) diff --git a/src/langbot/templates/config.yaml b/src/langbot/templates/config.yaml index 1129c08e8..e2ff82887 100644 --- a/src/langbot/templates/config.yaml +++ b/src/langbot/templates/config.yaml @@ -255,6 +255,12 @@ plugin: runtime_ws_url: 'ws://langbot_plugin_runtime:5400/control/ws' enable_marketplace: true display_plugin_debug_url: 'ws://localhost:5401/plugin/debug/ws' + certification: + # Non-secret Ed25519 issuer key ring used to verify the SDK ZIP-comment + # certification envelope. Values are standard base64-encoded raw public + # keys; add keys during issuer rotation and remove retired IDs only after + # every affected archive has been upgraded. + trusted_public_keys: {} worker: # Instance-wide maximum for every plugin installation. Plugin # manifests cannot raise or override these limits. diff --git a/tests/integration/api/test_plugins_security.py b/tests/integration/api/test_plugins_security.py index 5f79f716a..52442e371 100644 --- a/tests/integration/api/test_plugins_security.py +++ b/tests/integration/api/test_plugins_security.py @@ -3,11 +3,13 @@ from __future__ import annotations import copy +import io from types import SimpleNamespace from unittest.mock import AsyncMock, Mock, call import pytest import quart +from quart.datastructures import FileStorage pytestmark = pytest.mark.integration @@ -287,3 +289,34 @@ async def test_github_install_rejects_internal_asset_url_before_task_creation( assert response.status_code == 400 assert 'HTTPS GitHub release asset URL' in (await response.get_json())['msg'] application.task_mgr.create_user_task.assert_not_called() + + +@pytest.mark.asyncio +async def test_local_install_forwards_explicit_administrator_force(plugin_security_api): + application, client, _ = plugin_security_api + execution_context = SimpleNamespace( + instance_uuid='instance-test', + workspace_uuid=WORKSPACE_UUID, + placement_generation=1, + ) + application.persistence_mgr.tenant_scope = None + application.plugin_connector.require_workspace_context = AsyncMock(return_value=execution_context) + application.plugin_connector.install_plugin = AsyncMock() + application.task_mgr.create_user_task = Mock(return_value=SimpleNamespace(id='task-certification')) + + response = await client.post( + '/api/v1/plugins/install/local', + headers=_headers('manager-token'), + files={ + 'file': FileStorage(stream=io.BytesIO(b'archive'), filename='plugin.lbpkg'), + }, + form={'administrator_force': 'true'}, + ) + + assert response.status_code == 200 + operation = application.task_mgr.create_user_task.call_args.args[0] + await operation + assert application.plugin_connector.install_plugin.await_args.args[1] == { + 'plugin_file': b'archive', + 'administrator_force': True, + } diff --git a/tests/integration/persistence/test_migrations.py b/tests/integration/persistence/test_migrations.py index cbf282a71..9edfd5abb 100644 --- a/tests/integration/persistence/test_migrations.py +++ b/tests/integration/persistence/test_migrations.py @@ -71,7 +71,13 @@ def test_migration_graph_has_one_head_containing_both_released_branches(): heads = scripts.get_heads() assert len(heads) == 1, f'Release migrations must converge, found {heads}' ancestors = {revision.revision for revision in scripts.walk_revisions()} - assert {'0024_passkey_credentials', '0025_bot_plugin_processors'} <= ancestors + assert { + '0024_passkey_credentials', + '0025_bot_plugin_processors', + '0025_rag_document_identity', + '0027_pipeline_migration', + '0028_merge_knowledge_drafts', + } <= ancestors assert all(len(revision) <= 32 for revision in ancestors) diff --git a/tests/integration/persistence/test_rag_document_identity.py b/tests/integration/persistence/test_rag_document_identity.py new file mode 100644 index 000000000..ce209c63c --- /dev/null +++ b/tests/integration/persistence/test_rag_document_identity.py @@ -0,0 +1,421 @@ +"""Real database regressions for Host/engine identity; no live plugin or customer data.""" + +import asyncio +import os +import uuid +from types import SimpleNamespace +from unittest.mock import AsyncMock, Mock + +import pytest +import pytest_asyncio +import sqlalchemy as sa +from sqlalchemy.ext.asyncio import create_async_engine + +from langbot.pkg.api.http.context import ExecutionContext +from langbot.pkg.entity.persistence.base import Base +from langbot.pkg.entity.persistence.rag import File, KnowledgeBase +from langbot.pkg.entity.persistence.user import User +from langbot.pkg.entity.persistence.workspace import Workspace +from langbot.pkg.persistence.alembic_runner import ( + get_alembic_current, + run_alembic_downgrade, + run_alembic_stamp, + run_alembic_upgrade, +) +from langbot.pkg.persistence.mgr import PersistenceManager +from langbot.pkg.rag.knowledge.kbmgr import RuntimeKnowledgeBase +from langbot.pkg.workspace.errors import WorkspaceNotFoundError + +OLD_HEAD = '0024_passkey_credentials' + + +def current_head(): + from alembic.config import Config + from alembic.script import ScriptDirectory + from langbot.pkg.persistence.alembic_runner import _ALEMBIC_DIR + + config = Config() + config.set_main_option('script_location', _ALEMBIC_DIR) + return ScriptDirectory.from_config(config).get_current_head() + + +DOCUMENT_IDENTITY_REVISION = '0025_rag_document_identity' +CONTEXT = ExecutionContext(instance_uuid='instance-a', workspace_uuid='workspace-a', placement_generation=5) + + +@pytest_asyncio.fixture(params=['sqlite', 'postgres']) +async def database(request, tmp_path): + admin = None + schema = 'ke_identity_' + uuid.uuid4().hex + if request.param == 'postgres': + url = os.environ.get('TEST_POSTGRES_URL') + if not url: + pytest.skip('TEST_POSTGRES_URL is required for disposable PostgreSQL tests') + admin = create_async_engine(url) + async with admin.begin() as conn: + await conn.execute(sa.text(f'CREATE SCHEMA {schema}')) + engine = create_async_engine(url, connect_args={'server_settings': {'search_path': schema}}) + else: + engine = create_async_engine(f'sqlite+aiosqlite:///{tmp_path / "identity.db"}') + + @sa.event.listens_for(engine.sync_engine, 'connect') + def enable_foreign_keys(connection, _): + connection.execute('PRAGMA foreign_keys=ON') + + try: + yield engine + finally: + await engine.dispose() + if admin is not None: + async with admin.begin() as conn: + await conn.execute(sa.text(f'DROP SCHEMA {schema} CASCADE')) + await admin.dispose() + + +async def create_schema(engine, *, legacy=False): + # Only the fixture-owned dependency closure; never all imported application tables. + async with engine.begin() as conn: + await conn.run_sync( + lambda sync: Base.metadata.create_all( + sync, tables=[User.__table__, Workspace.__table__, KnowledgeBase.__table__] + ) + ) + if legacy: + # Exact pre-0025 File shape, NOT current metadata stamped with an old revision. + await conn.execute( + sa.text("""CREATE TABLE knowledge_base_files ( + uuid VARCHAR(255) PRIMARY KEY UNIQUE, + workspace_uuid VARCHAR(36) NOT NULL REFERENCES workspaces(uuid) ON DELETE CASCADE, + kb_id VARCHAR(255), file_name VARCHAR, extension VARCHAR, created_at TIMESTAMP, status VARCHAR, + CONSTRAINT uq_knowledge_base_files_workspace_uuid UNIQUE (workspace_uuid, uuid), + CONSTRAINT fk_knowledge_base_files_workspace_kb FOREIGN KEY (workspace_uuid, kb_id) + REFERENCES knowledge_bases(workspace_uuid, uuid) ON DELETE CASCADE + )""") + ) + await conn.execute( + sa.text( + 'CREATE INDEX ix_knowledge_base_files_workspace_kb ON knowledge_base_files (workspace_uuid, kb_id)' + ) + ) + else: + await conn.run_sync(lambda sync: File.__table__.create(sync)) + for workspace in ('workspace-a', 'workspace-b'): + await conn.execute( + sa.insert(Workspace).values( + uuid=workspace, + instance_uuid='instance-a', + name=workspace, + slug=workspace, + source='cloud_projection', + ) + ) + for kb, workspace in [('kb-a', 'workspace-a'), ('kb-other', 'workspace-a'), ('kb-b', 'workspace-b')]: + await conn.execute( + sa.insert(KnowledgeBase).values( + uuid=kb, + workspace_uuid=workspace, + name=kb, + knowledge_engine_plugin_id='author/engine', + collection_id=kb, + ) + ) + + +@pytest_asyncio.fixture +async def runtime(database): + await create_schema(database) + ap = SimpleNamespace( + logger=Mock(), + workspace_service=SimpleNamespace( + get_execution_binding=AsyncMock( + return_value=SimpleNamespace( + instance_uuid='instance-a', + placement_generation=5, + ) + ) + ), + storage_mgr=SimpleNamespace( + require_scoped_object_key=Mock(), + size_scoped_object_key=AsyncMock(return_value=12), + delete_scoped_object_key=AsyncMock(), + ), + plugin_connector=SimpleNamespace( + require_workspace_context=AsyncMock(side_effect=lambda context: context), + call_rag_ingest=AsyncMock(return_value={'document_id': 'upstream-id', 'status': 'processing'}), + call_rag_delete_document=AsyncMock(return_value=True), + ), + ) + ap.persistence_mgr = PersistenceManager(ap) + ap.persistence_mgr.db = SimpleNamespace(get_engine=lambda: database) + kb = KnowledgeBase(uuid='kb-a', workspace_uuid='workspace-a', name='kb', knowledge_engine_plugin_id='author/engine') + return RuntimeKnowledgeBase(ap, kb, CONTEXT) + + +async def seed(runtime, *, status='pending', file_id='host-id', workspace='workspace-a', kb='kb-a'): + values = dict( + uuid=file_id, workspace_uuid=workspace, kb_id=kb, file_name='upload.txt', extension='txt', status=status + ) + await runtime.ap.persistence_mgr.execute_async(sa.insert(File).values(**values)) + return File(**values) + + +async def read_file_row(runtime, file_id='host-id'): + # A fresh connection proves committed state rather than an identity-map/mock result. + async with runtime.ap.persistence_mgr.get_db_engine().connect() as conn: + row = (await conn.execute(sa.select(File).where(File.uuid == file_id))).first() + return None if row is None else dict(row._mapping) + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + 'engine_id', ['dify-upstream', 'fastgpt-collection', 'ragflow-upstream', 'host-id', ' opaque ID '] +) +async def test_ingest_persists_exact_engine_identity_and_restart_delete(runtime, engine_id): + file = await seed(runtime) + runtime.ap.plugin_connector.call_rag_ingest.return_value['document_id'] = engine_id + await runtime._store_file_task(CONTEXT, file, Mock()) + row = await read_file_row(runtime) + assert row['uuid'] == 'host-id' + assert row.get('engine_document_id') == engine_id + assert row['status'] == 'completed' + restarted = RuntimeKnowledgeBase(runtime.ap, runtime.knowledge_base_entity, CONTEXT) + await restarted.delete_file(CONTEXT, 'host-id') + runtime.ap.plugin_connector.call_rag_delete_document.assert_awaited_once_with('author/engine', engine_id, 'kb-a') + assert await read_file_row(runtime) is None + + +@pytest.mark.asyncio +@pytest.mark.parametrize('response', [False, None, 0, 1, 'true', {}]) +async def test_delete_without_explicit_confirmation_preserves_row(runtime, response): + await seed(runtime, status='completed') + runtime.ap.plugin_connector.call_rag_delete_document.return_value = response + with pytest.raises(RuntimeError, match='delet'): + await runtime.delete_file(CONTEXT, 'host-id') + assert await read_file_row(runtime) is not None + + +@pytest.mark.asyncio +@pytest.mark.parametrize('failure', ['exception', 'missing_plugin']) +async def test_delete_unavailable_retains_row(runtime, failure): + await seed(runtime, status='completed') + if failure == 'exception': + runtime.ap.plugin_connector.call_rag_delete_document.side_effect = RuntimeError('upstream offline') + else: + runtime.knowledge_base_entity.knowledge_engine_plugin_id = None + with pytest.raises(RuntimeError, match='delet'): + await runtime.delete_file(CONTEXT, 'host-id') + assert await read_file_row(runtime) is not None + + +@pytest.mark.asyncio +@pytest.mark.parametrize('status', ['pending', 'processing']) +async def test_delete_rejects_inflight_ingestion(runtime, status): + await seed(runtime, status=status) + with pytest.raises(RuntimeError, match='ingest'): + await runtime.delete_file(CONTEXT, 'host-id') + runtime.ap.plugin_connector.call_rag_delete_document.assert_not_awaited() + assert await read_file_row(runtime) is not None + + +@pytest.mark.asyncio +@pytest.mark.parametrize('document_id', [None, '', ' ', 123, [], {}]) +async def test_invalid_response_identity_cannot_complete(runtime, document_id): + file = await seed(runtime) + runtime.ap.plugin_connector.call_rag_ingest.return_value = {'status': 'completed', 'document_id': document_id} + with pytest.raises(ValueError, match='document_id'): + await runtime._store_file_task(CONTEXT, file, Mock()) + assert (await read_file_row(runtime))['status'] == 'interrupted' + + +@pytest.mark.asyncio +async def test_missing_response_identity_cannot_complete(runtime): + file = await seed(runtime) + runtime.ap.plugin_connector.call_rag_ingest.return_value = {'status': 'completed'} + with pytest.raises(ValueError, match='document_id'): + await runtime._store_file_task(CONTEXT, file, Mock()) + assert (await read_file_row(runtime))['status'] == 'interrupted' + + +@pytest.mark.asyncio +async def test_failed_ingestion_retains_returned_identity_for_cleanup(runtime): + file = await seed(runtime) + runtime.ap.plugin_connector.call_rag_ingest.return_value = { + 'document_id': 'uploaded-before-parsing-failed', + 'status': 'failed', + 'error_message': 'parsing failed', + } + with pytest.raises(Exception, match='parsing failed'): + await runtime._store_file_task(CONTEXT, file, Mock()) + row = await read_file_row(runtime) + assert row['status'] == 'failed' + assert row.get('engine_document_id') == 'uploaded-before-parsing-failed' + await runtime.delete_file(CONTEXT, 'host-id') + runtime.ap.plugin_connector.call_rag_delete_document.assert_awaited_once_with( + 'author/engine', 'uploaded-before-parsing-failed', 'kb-a' + ) + + +@pytest.mark.asyncio +@pytest.mark.parametrize('status', ['completed', 'failed']) +async def test_legacy_rows_use_host_id_only_with_confirmed_delete(runtime, status): + await seed(runtime, status=status) + await runtime.delete_file(CONTEXT, 'host-id') + runtime.ap.plugin_connector.call_rag_delete_document.assert_awaited_once_with('author/engine', 'host-id', 'kb-a') + assert await read_file_row(runtime) is None + + +@pytest.mark.asyncio +@pytest.mark.parametrize('workspace,kb', [('workspace-a', 'kb-other'), ('workspace-b', 'kb-b')]) +async def test_status_update_and_delete_do_not_touch_other_scope(runtime, workspace, kb): + await seed(runtime, workspace=workspace, kb=kb) + assert not await runtime._set_file_status(CONTEXT, 'host-id', 'processing') + with pytest.raises(WorkspaceNotFoundError): + await runtime.delete_file(CONTEXT, 'host-id') + assert (await read_file_row(runtime))['status'] == 'pending' + runtime.ap.plugin_connector.call_rag_delete_document.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_delete_rechecks_generation_after_plugin_response(runtime): + await seed(runtime, status='completed') + + async def delete(*_): + runtime.ap.workspace_service.get_execution_binding.side_effect = WorkspaceNotFoundError('stale placement') + return True + + runtime.ap.plugin_connector.call_rag_delete_document.side_effect = delete + with pytest.raises(WorkspaceNotFoundError): + await runtime.delete_file(CONTEXT, 'host-id') + assert await read_file_row(runtime) is not None + + +@pytest.mark.asyncio +async def test_ingest_rechecks_generation_before_mapping_write(runtime): + file = await seed(runtime) + + async def ingest(*_): + runtime.ap.workspace_service.get_execution_binding.side_effect = WorkspaceNotFoundError('stale placement') + return {'document_id': 'upstream-id', 'status': 'completed'} + + runtime.ap.plugin_connector.call_rag_ingest.side_effect = ingest + with pytest.raises(WorkspaceNotFoundError): + await runtime._store_file_task(CONTEXT, file, Mock()) + row = await read_file_row(runtime) + assert row['status'] == 'processing' + assert row.get('engine_document_id') is None + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_concurrent_delete_cannot_remove_ingestion_tracking(runtime): + file = await seed(runtime) + entered, release = asyncio.Event(), asyncio.Event() + + async def ingest(*_): + entered.set() + await release.wait() + return {'document_id': 'upstream-id', 'status': 'completed'} + + runtime.ap.plugin_connector.call_rag_ingest.side_effect = ingest + task = asyncio.create_task(runtime._store_file_task(CONTEXT, file, Mock())) + try: + await asyncio.wait_for(entered.wait(), 5) + with pytest.raises(RuntimeError, match='ingest'): + await runtime.delete_file(CONTEXT, 'host-id') + finally: + release.set() + await task + assert (await read_file_row(runtime)).get('engine_document_id') == 'upstream-id' + + +@pytest.mark.asyncio +async def test_identity_and_completion_are_one_atomic_write(runtime): + file = await seed(runtime) + engine = runtime.ap.persistence_mgr.get_db_engine() + attempts = [] + + def reject_completion(_conn, _cursor, statement, parameters, _context, _many): + if statement.startswith('UPDATE knowledge_base_files') and 'completed' in parameters: + attempts.append(statement) + raise RuntimeError('fixture mapping write failure') + + sa.event.listen(engine.sync_engine, 'before_cursor_execute', reject_completion) + try: + with pytest.raises(RuntimeError, match='mapping write failure'): + await runtime._store_file_task(CONTEXT, file, Mock()) + finally: + sa.event.remove(engine.sync_engine, 'before_cursor_execute', reject_completion) + assert len(attempts) == 1 + assert 'engine_document_id=' in attempts[0] + row = await read_file_row(runtime) + assert row['status'] != 'completed' + assert row.get('engine_document_id') == 'upstream-id' + + +@pytest.mark.asyncio +async def test_populated_legacy_migration_roundtrip(database): + await create_schema(database, legacy=True) + async with database.begin() as conn: + await conn.execute( + sa.text("""INSERT INTO knowledge_base_files + (uuid, workspace_uuid, kb_id, file_name, extension, status) + VALUES ('legacy', 'workspace-a', 'kb-a', 'original.txt', 'txt', 'completed')""") + ) + assert 'engine_document_id' not in await conn.run_sync( + lambda sync: {col['name'] for col in sa.inspect(sync).get_columns('knowledge_base_files')} + ) + await run_alembic_stamp(database, OLD_HEAD) + await run_alembic_upgrade(database, DOCUMENT_IDENTITY_REVISION) + async with database.connect() as conn: + columns = await conn.run_sync(lambda sync: sa.inspect(sync).get_columns('knowledge_base_files')) + assert 'engine_document_id' in {col['name'] for col in columns} + assert next(col for col in columns if col['name'] == 'engine_document_id')['nullable'] + row = (await conn.execute(sa.text('SELECT * FROM knowledge_base_files'))).mappings().one() + assert row['uuid'] == 'legacy' and row['status'] == 'completed' + assert row['engine_document_id'] is None + assert await get_alembic_current(database) == DOCUMENT_IDENTITY_REVISION + await run_alembic_upgrade(database, DOCUMENT_IDENTITY_REVISION) + await run_alembic_stamp(database, OLD_HEAD) + await run_alembic_upgrade(database, DOCUMENT_IDENTITY_REVISION) + await run_alembic_downgrade(database, OLD_HEAD) + async with database.connect() as conn: + assert 'engine_document_id' not in await conn.run_sync( + lambda sync: {col['name'] for col in sa.inspect(sync).get_columns('knowledge_base_files')} + ) + assert (await conn.execute(sa.text('SELECT uuid FROM knowledge_base_files'))).scalar_one() == 'legacy' + await run_alembic_upgrade(database) + + +@pytest.mark.asyncio +async def test_published_document_identity_branch_converges_to_release_head(database): + await create_schema(database, legacy=True) + await run_alembic_stamp(database, OLD_HEAD) + await run_alembic_upgrade(database, DOCUMENT_IDENTITY_REVISION) + async with database.begin() as conn: + await conn.execute( + sa.text("""INSERT INTO knowledge_base_files + (uuid, workspace_uuid, kb_id, file_name, extension, status, engine_document_id) + VALUES ('host-stable', 'workspace-a', 'kb-a', 'original.txt', 'txt', 'completed', 'opaque-engine-id')""") + ) + await run_alembic_upgrade(database) + await run_alembic_upgrade(database) + assert await get_alembic_current(database) == current_head() + async with database.connect() as conn: + row = (await conn.execute(sa.text('SELECT * FROM knowledge_base_files'))).mappings().one() + assert row['uuid'] == 'host-stable' + assert row['engine_document_id'] == 'opaque-engine-id' + assert row['status'] == 'completed' + tables = await conn.run_sync(lambda sync: sa.inspect(sync).get_table_names()) + assert 'pipeline_migration_snapshots' in tables + + +@pytest.mark.asyncio +async def test_fresh_metadata_then_migration_is_idempotent(database): + await create_schema(database) + await run_alembic_stamp(database, OLD_HEAD) + await run_alembic_upgrade(database) + assert await get_alembic_current(database) == current_head() + async with database.connect() as conn: + assert 'engine_document_id' in await conn.run_sync( + lambda sync: {col['name'] for col in sa.inspect(sync).get_columns('knowledge_base_files')} + ) diff --git a/tests/integration/persistence/test_rag_interrupted_ingestion.py b/tests/integration/persistence/test_rag_interrupted_ingestion.py new file mode 100644 index 000000000..aa6b53dc6 --- /dev/null +++ b/tests/integration/persistence/test_rag_interrupted_ingestion.py @@ -0,0 +1,406 @@ +"""B1 lifecycle regressions; isolated databases, no external engines or customer data.""" + +import asyncio +from types import SimpleNamespace +from unittest.mock import AsyncMock, Mock + +import pytest + +from tests.integration.persistence import test_rag_document_identity as identity +from tests.integration.persistence.test_rag_document_identity import CONTEXT, read_file_row, seed +from langbot.pkg.core.entities import LifecycleControlScope +from langbot.pkg.core.taskmgr import AsyncTaskManager, TaskContext +from langbot.pkg.rag.knowledge.kbmgr import RuntimeKnowledgeBase + +# Reuse the real isolated-database fixtures without duplicating their lifecycle. +database = identity.database +runtime = identity.runtime + + +def attach_task_manager(runtime): + runtime.ap.event_loop = asyncio.get_running_loop() + runtime.ap.instance_config = SimpleNamespace(data={}) + runtime.ap.task_mgr = AsyncTaskManager(runtime.ap) + return runtime.ap.task_mgr + + +@pytest.mark.asyncio +async def test_cancelled_ingestion_retains_recovery_resources(runtime): + """Local cancellation is interrupted observation, not remote quiescence.""" + file = await seed(runtime) + entered, terminal = asyncio.Event(), asyncio.Event() + + async def ingest(*_): + entered.set() + try: + await asyncio.Event().wait() + finally: + terminal.set() + + runtime.ap.plugin_connector.call_rag_ingest.side_effect = ingest + manager = attach_task_manager(runtime) + wrapper = manager.create_user_task( + runtime._store_file_task(CONTEXT, file, TaskContext.new()), + kind='knowledge-operation', + name=f'knowledge-store-file-{file.file_name}', + instance_uuid=CONTEXT.instance_uuid, + workspace_uuid=CONTEXT.workspace_uuid, + placement_generation=CONTEXT.placement_generation, + ) + try: + await asyncio.wait_for(entered.wait(), 5) + manager.cancel_by_scope(LifecycleControlScope.APPLICATION) + with pytest.raises(asyncio.CancelledError): + await wrapper.task + assert terminal.is_set() + assert wrapper.task.cancelled() + assert (await read_file_row(runtime))['status'] == 'interrupted' + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + with pytest.raises(RuntimeError, match='interrupted'): + await runtime.delete_file(CONTEXT, file.uuid) + runtime.ap.plugin_connector.call_rag_delete_document.assert_not_awaited() + assert await read_file_row(runtime) is not None + finally: + if not wrapper.task.done(): + wrapper.task.cancel() + await asyncio.gather(wrapper.task, return_exceptions=True) + + +@pytest.mark.asyncio +@pytest.mark.parametrize('status', ['pending', 'processing']) +async def test_recreated_host_marks_abandoned_rows_interrupted(runtime, status): + """Fresh Host state recovers observation without authorizing remote deletion.""" + await seed(runtime, status=status) + new_ap = SimpleNamespace(**vars(runtime.ap)) + recreated = RuntimeKnowledgeBase(new_ap, runtime.knowledge_base_entity, CONTEXT) + manager = attach_task_manager(recreated) + assert manager.get_all_tasks() == [] + await recreated.initialize() + assert (await read_file_row(recreated))['status'] == 'interrupted' + with pytest.raises(RuntimeError, match='interrupted'): + await recreated.delete_file(CONTEXT, 'host-id') + recreated.ap.plugin_connector.call_rag_delete_document.assert_not_awaited() + recreated.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_b1_live_task_still_blocks_delete_after_runtime_object_recreation(runtime): + file = await seed(runtime) + entered, release = asyncio.Event(), asyncio.Event() + + async def ingest(*_): + entered.set() + await release.wait() + return {'document_id': file.uuid, 'status': 'completed'} + + runtime.ap.plugin_connector.call_rag_ingest.side_effect = ingest + manager = attach_task_manager(runtime) + wrapper = manager.create_user_task( + runtime._store_file_task(CONTEXT, file, Mock()), + kind='knowledge-operation', + name=f'knowledge-store-file-{file.file_name}', + instance_uuid=CONTEXT.instance_uuid, + workspace_uuid=CONTEXT.workspace_uuid, + placement_generation=CONTEXT.placement_generation, + ) + try: + await asyncio.wait_for(entered.wait(), 5) + recreated = RuntimeKnowledgeBase(runtime.ap, runtime.knowledge_base_entity, CONTEXT) + await recreated.initialize() + assert not wrapper.task.done() + with pytest.raises(RuntimeError, match='ingest'): + await recreated.delete_file(CONTEXT, file.uuid) + recreated.ap.plugin_connector.call_rag_delete_document.assert_not_awaited() + assert (await read_file_row(recreated))['status'] == 'processing' + finally: + release.set() + await asyncio.wait_for(wrapper.task, 5) + await recreated.delete_file(CONTEXT, file.uuid) + assert await read_file_row(recreated) is None + + +@pytest.mark.asyncio +@pytest.mark.parametrize('failure', [TimeoutError('timeout'), ConnectionError('disconnect')]) +async def test_transport_failure_is_interrupted_not_failed(runtime, failure): + file = await seed(runtime) + runtime.ap.plugin_connector.call_rag_ingest.side_effect = failure + with pytest.raises(type(failure)): + await runtime._store_file_task(CONTEXT, file, Mock()) + assert (await read_file_row(runtime))['status'] == 'interrupted' + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_recovered_row_cannot_be_replayed_by_delayed_old_task(runtime): + file = await seed(runtime) + await runtime.initialize() + with pytest.raises(Exception): + await runtime._store_file_task(CONTEXT, file, Mock()) + assert (await read_file_row(runtime))['status'] == 'interrupted' + runtime.ap.plugin_connector.call_rag_ingest.assert_not_awaited() + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_reconciliation_preserves_other_scopes_and_terminal_states(runtime): + await seed(runtime, status='processing') + await seed(runtime, file_id='other-kb', kb='kb-other', status='processing') + await seed(runtime, file_id='other-workspace', workspace='workspace-b', kb='kb-b', status='processing') + for status in ('completed', 'failed', 'interrupted'): + await seed(runtime, file_id=status, status=status) + await runtime.initialize() + assert (await read_file_row(runtime))['status'] == 'interrupted' + for file_id in ('other-kb', 'other-workspace'): + assert (await read_file_row(runtime, file_id))['status'] == 'processing' + for status in ('completed', 'failed', 'interrupted'): + assert (await read_file_row(runtime, status))['status'] == status + + +@pytest.mark.asyncio +async def test_cancellation_after_result_retains_identity(runtime): + file = await seed(runtime) + original = runtime._set_file_status + + async def cancel_completion(context, uuid, status, **kwargs): + if status == 'completed': + raise asyncio.CancelledError() + return await original(context, uuid, status, **kwargs) + + runtime._set_file_status = cancel_completion + with pytest.raises(asyncio.CancelledError): + await runtime._store_file_task(CONTEXT, file, Mock()) + row = await read_file_row(runtime) + assert row['status'] == 'interrupted' + assert row['engine_document_id'] == 'upstream-id' + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_cancellation_after_completion_commit_does_not_downgrade(runtime): + file = await seed(runtime) + original = runtime._set_file_status + + async def cancel_after_commit(context, uuid, status, **kwargs): + result = await original(context, uuid, status, **kwargs) + if status == 'completed': + raise asyncio.CancelledError() + return result + + runtime._set_file_status = cancel_after_commit + with pytest.raises(asyncio.CancelledError): + await runtime._store_file_task(CONTEXT, file, Mock()) + row = await read_file_row(runtime) + assert row['status'] == 'completed' + assert row['engine_document_id'] == 'upstream-id' + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + + +def make_service(runtime): + from langbot.pkg.api.http.service.knowledge import KnowledgeService + + runtime.ap.rag_mgr = SimpleNamespace( + get_knowledge_base_details=AsyncMock(return_value={'uuid': 'kb-a'}), + get_knowledge_base_by_uuid=AsyncMock(return_value=runtime), + delete_knowledge_base=AsyncMock(), + ) + return KnowledgeService(runtime.ap) + + +@pytest.mark.asyncio +async def test_file_list_recovers_abandoned_task_without_restart(runtime): + await seed(runtime, status='processing') + files = await make_service(runtime).get_files_by_knowledge_base(CONTEXT, 'kb-a') + assert files[0]['status'] == 'interrupted' + assert (await read_file_row(runtime))['status'] == 'interrupted' + + +@pytest.mark.asyncio +@pytest.mark.parametrize('status', ['pending', 'processing', 'interrupted']) +async def test_bulk_kb_delete_cannot_discard_unsettled_ingestion(runtime, status): + from langbot.pkg.entity.persistence.rag import Chunk + + async with runtime.ap.persistence_mgr.get_db_engine().begin() as conn: + await conn.run_sync(lambda sync: Chunk.__table__.create(sync)) + await seed(runtime, status=status) + service = make_service(runtime) + with pytest.raises(RuntimeError, match='ingestion'): + await service.delete_knowledge_base(CONTEXT, 'kb-a') + assert await read_file_row(runtime) is not None + runtime.ap.rag_mgr.delete_knowledge_base.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_queued_task_cancelled_before_start_recovers_on_file_list(runtime): + runtime.ap.storage_mgr.exists_scoped_object_key = AsyncMock(return_value=True) + manager = attach_task_manager(runtime) + await runtime.store_file(CONTEXT, 'upload.txt') + wrapper = manager.get_all_tasks()[0] + wrapper.task.cancel() + await asyncio.gather(wrapper.task, return_exceptions=True) + assert wrapper.task.cancelled() + files = await make_service(runtime).get_files_by_knowledge_base(CONTEXT, 'kb-a') + assert files[0]['status'] == 'interrupted' + runtime.ap.plugin_connector.call_rag_ingest.assert_not_awaited() + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + + +@pytest.mark.asyncio +@pytest.mark.parametrize('disconnect', [False, True], ids=['cancel-waiter', 'close-host-connection']) +async def test_actual_sdk_late_write_preserves_interrupted_state(runtime, tmp_path, monkeypatch, disconnect): + from langbot.pkg.api.http.context import ExecutionContext + from langbot_plugin.entities.io.actions.enums import RuntimeToPluginAction + from langbot_plugin.runtime.io.handler import ActionResponse + from tests.integration.plugin.test_rag_file_transfer_protocol import BINDING, protocol_stack + + context = ExecutionContext(instance_uuid='instance-a', workspace_uuid='workspace-a', placement_generation=7) + runtime.execution_context = context + file = await seed(runtime) + async with protocol_stack(tmp_path, monkeypatch, 'shared', BINDING) as stack: + entered, release, finished = asyncio.Event(), asyncio.Event(), asyncio.Event() + selected = SimpleNamespace(_runtime_plugin_handler=stack.bridge) + monkeypatch.setattr(stack.runtime.plugin_mgr, '_get_connected_rag_plugin', lambda *_: (selected, 'engine')) + vectors = set() + + @stack.plugin.action(RuntimeToPluginAction.INGEST_DOCUMENT) + async def ingest(data): + entered.set() + await release.wait() + vectors.add('opaque-upstream-id') + finished.set() + return ActionResponse.success({'document_id': 'opaque-upstream-id', 'status': 'completed'}) + + async def send_ingest(_plugin_id, data): + with stack.core.installation_scope(BINDING): + return await stack.core.rag_ingest_document('tester', 'engine', data) + + runtime.ap.plugin_connector.call_rag_ingest = send_ingest + task = asyncio.create_task(runtime._store_file_task(context, file, Mock())) + try: + await asyncio.wait_for(entered.wait(), 5) + if disconnect: + await stack.core.close() + else: + task.cancel() + outcome = await asyncio.gather(task, return_exceptions=True) + assert isinstance(outcome[0], BaseException) + assert (await read_file_row(runtime))['status'] == 'interrupted' + assert not finished.is_set() + with pytest.raises(RuntimeError, match='interrupted'): + await runtime.delete_file(context, file.uuid) + runtime.ap.plugin_connector.call_rag_delete_document.assert_not_awaited() + release.set() + await asyncio.wait_for(finished.wait(), 5) + assert vectors == {'opaque-upstream-id'} + assert (await read_file_row(runtime))['status'] == 'interrupted' + assert (await read_file_row(runtime))['engine_document_id'] is None + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + finally: + release.set() + if not task.done(): + task.cancel() + await asyncio.gather(task, return_exceptions=True) + + +@pytest.mark.asyncio +async def test_late_observed_identity_survives_recovery_without_false_completion(runtime): + file = await seed(runtime) + entered, release = asyncio.Event(), asyncio.Event() + + async def ingest(*_): + entered.set() + await release.wait() + return {'document_id': 'late-id', 'status': 'completed'} + + runtime.ap.plugin_connector.call_rag_ingest.side_effect = ingest + task = asyncio.create_task(runtime._store_file_task(CONTEXT, file, Mock())) + try: + await asyncio.wait_for(entered.wait(), 5) + # Independent Host observation, not shared process state. + new_ap = SimpleNamespace(**{k: v for k, v in vars(runtime.ap).items() if not k.startswith('_knowledge_')}) + recreated = RuntimeKnowledgeBase(new_ap, runtime.knowledge_base_entity, CONTEXT) + await recreated.initialize() + release.set() + await asyncio.gather(task, return_exceptions=True) + row = await read_file_row(runtime) + assert row['status'] == 'interrupted' + assert row['engine_document_id'] == 'late-id' + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + finally: + release.set() + await asyncio.gather(task, return_exceptions=True) + + +@pytest.mark.asyncio +async def test_parser_transport_failure_retains_upload(runtime): + file = await seed(runtime) + runtime.ap.storage_mgr.load_scoped_object_key = AsyncMock(return_value=b'file') + runtime.ap.plugin_connector.call_parser = AsyncMock(side_effect=TimeoutError()) + with pytest.raises(TimeoutError): + await runtime._store_file_task(CONTEXT, file, Mock(), parser_plugin_id='author/parser') + assert (await read_file_row(runtime))['status'] == 'interrupted' + runtime.ap.storage_mgr.delete_scoped_object_key.assert_not_awaited() + runtime.ap.plugin_connector.call_rag_ingest.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_kb_delete_serializes_with_new_upload_admission(runtime): + from langbot.pkg.entity.persistence.rag import Chunk + + async with runtime.ap.persistence_mgr.get_db_engine().begin() as conn: + await conn.run_sync(lambda sync: Chunk.__table__.create(sync)) + runtime.ap.storage_mgr.exists_scoped_object_key = AsyncMock(return_value=True) + attach_task_manager(runtime) + service = make_service(runtime) + entered, release = asyncio.Event(), asyncio.Event() + + async def pause_delete(*_): + entered.set() + await release.wait() + + runtime.ap.rag_mgr.delete_knowledge_base.side_effect = pause_delete + deletion = asyncio.create_task(service.delete_knowledge_base(CONTEXT, 'kb-a')) + upload = None + try: + await asyncio.wait_for(entered.wait(), 5) + upload = asyncio.create_task(runtime.store_file(CONTEXT, 'upload.txt')) + await asyncio.sleep(0.1) + assert not upload.done(), 'upload admission must wait for the in-progress KB deletion' + runtime.ap.plugin_connector.call_rag_ingest.assert_not_awaited() + release.set() + await deletion + with pytest.raises(Exception): + await upload # FK rejects admission after the KB was deleted. + runtime.ap.plugin_connector.call_rag_ingest.assert_not_awaited() + finally: + release.set() + await asyncio.gather(deletion, *([upload] if upload else []), return_exceptions=True) + await asyncio.gather(*(wrapper.task for wrapper in runtime.ap.task_mgr.get_all_tasks()), return_exceptions=True) + + +@pytest.mark.asyncio +@pytest.mark.parametrize('provider_name', ['LocalStorageProvider', 'S3StorageProvider']) +async def test_retention_cleanup_preserves_ingestion_recovery_upload(runtime, tmp_path, provider_name): + from langbot.pkg.api.http.service.maintenance import MaintenanceService + + await seed(runtime, status='interrupted') + retained = tmp_path / 'retained.txt' + expired = tmp_path / 'expired.txt' + retained.write_text('recovery source') + expired.write_text('unreferenced upload') + provider = type(provider_name, (), {})() + provider.delete = AsyncMock() + runtime.ap.storage_mgr.storage_provider = provider + candidates = [ + {'key': 'upload.txt', 'path': str(retained)}, + {'key': 'unreferenced.txt', 'path': str(expired)}, + ] + service = MaintenanceService(runtime.ap) + service._expired_local_upload_candidates = Mock(return_value=candidates) + service._expired_s3_upload_candidates = AsyncMock(return_value=candidates) + count = await service._cleanup_expired_uploaded_files(CONTEXT, 7) + assert count == 1 + if provider_name == 'LocalStorageProvider': + assert retained.read_text() == 'recovery source' + assert not expired.exists() + else: + provider.delete.assert_awaited_once_with('unreferenced.txt') diff --git a/tests/integration/plugin/test_certified_plugin_admission.py b/tests/integration/plugin/test_certified_plugin_admission.py new file mode 100644 index 000000000..7da555e86 --- /dev/null +++ b/tests/integration/plugin/test_certified_plugin_admission.py @@ -0,0 +1,223 @@ +"""Certified archive admission through the public Core installation API.""" + +from __future__ import annotations + +import base64 +import hashlib +import io +import zipfile +from types import SimpleNamespace +from unittest.mock import AsyncMock, Mock + +import pytest +import yaml +from cryptography.hazmat.primitives import serialization +from cryptography.hazmat.primitives.asymmetric.ed25519 import Ed25519PrivateKey + +from langbot.pkg.api.http.context import ExecutionContext +from langbot.pkg.plugin.connector import PluginRuntimeConnector +from langbot_plugin.entities.io.context import InstallationBinding +from langbot_plugin.runtime.plugin.mgr import PluginInstallSource + + +pytestmark = pytest.mark.integration + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + ('deployment', 'archive_kind', 'administrator_force', 'expected_profile'), + [ + ('cloud', 'signed_shared', False, 'shared-runtime-v1'), + ('oss', 'signed_shared', False, 'shared-runtime-v1'), + ('oss', 'legacy', False, 'dedicated'), + ('oss', 'invalid_shared', True, 'dedicated'), + ], +) +async def test_install_plugin_admits_archive_before_persistence_and_applies_selected_profile( + deployment: str, + archive_kind: str, + administrator_force: bool, + expected_profile: str, +) -> None: + package, trusted_public_keys = _archive(archive_kind) + connector, execution_context, binding = _connector(deployment, trusted_public_keys) + + await connector.install_plugin( + PluginInstallSource.LOCAL, + { + 'plugin_file': package, + 'administrator_force': administrator_force, + }, + ) + + connector._store_artifact_package.assert_awaited_once_with( + execution_context, + hashlib.sha256(package).hexdigest(), + package, + ) + persisted_info = connector._persist_installation_package.await_args.kwargs['install_info'] + assert persisted_info['_certification']['runtime_profile'] == expected_profile + assert persisted_info['_certification']['normalized_digest'] == _normalized_digest(package) + connector.handler.apply_plugin_installation.assert_awaited_once_with( + binding, + artifact_package=package, + enabled=True, + ) + + +@pytest.mark.asyncio +@pytest.mark.parametrize('archive_kind', ['legacy', 'invalid_shared']) +async def test_cloud_rejects_untrusted_archive_before_storage_persistence_or_runtime_apply(archive_kind: str) -> None: + package, trusted_public_keys = _archive(archive_kind) + connector, _execution_context, _binding = _connector('cloud', trusted_public_keys) + + with pytest.raises(ValueError, match='CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_'): + await connector.install_plugin(PluginInstallSource.LOCAL, {'plugin_file': package}) + + connector._store_artifact_package.assert_not_awaited() + connector._persist_installation_package.assert_not_awaited() + connector.handler.apply_plugin_installation.assert_not_awaited() + + +@pytest.mark.asyncio +async def test_oss_requires_explicit_administrator_force_for_declared_invalid_archive() -> None: + package, trusted_public_keys = _archive('invalid_shared') + connector, _execution_context, _binding = _connector('oss', trusted_public_keys) + + with pytest.raises(ValueError, match='CERTIFIED_PLUGIN_OSS_FORCE_REQUIRED'): + await connector.install_plugin(PluginInstallSource.LOCAL, {'plugin_file': package}) + + connector._store_artifact_package.assert_not_awaited() + connector._persist_installation_package.assert_not_awaited() + connector.handler.apply_plugin_installation.assert_not_awaited() + + +@pytest.mark.asyncio +@pytest.mark.parametrize('requested_version', [None, '1.0.0']) +@pytest.mark.parametrize('archive_kind', ['signed_shared', 'legacy']) +async def test_marketplace_version_selection_keeps_certificate_gate_and_single_apply( + monkeypatch, requested_version, archive_kind +): + import json + + import langbot.pkg.plugin.connector as connector_module + from langbot.pkg.core.taskmgr import TaskContext + + package, trusted_public_keys = _archive(archive_kind) + connector, _execution_context, binding = _connector('cloud', trusted_public_keys) + connector._refresh_runner_registry = AsyncMock() + requests = [] + + async def marketplace_get(_client, url, **kwargs): + requests.append(url) + if '/plugins/download/' in url: + assert url.endswith('/certified/example/1.0.0') + return 200, package + if url.endswith('/versions'): + return 200, json.dumps({'data': {'versions': [{'version': '1.0.0'}]}}).encode() + return 404, b'{}' + + monkeypatch.setattr(connector_module, '_marketplace_get', marketplace_get) + info = {'plugin_author': 'certified', 'plugin_name': 'example'} + if requested_version is not None: + info['plugin_version'] = requested_version + task_context = TaskContext.new() + if archive_kind == 'legacy': + with pytest.raises(ValueError, match='CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_REQUIRED'): + await connector.install_plugin(PluginInstallSource.MARKETPLACE, info, task_context) + connector._store_artifact_package.assert_not_awaited() + connector._persist_installation_package.assert_not_awaited() + connector.handler.apply_plugin_installation.assert_not_awaited() + connector._refresh_runner_registry.assert_not_awaited() + else: + await connector.install_plugin(PluginInstallSource.MARKETPLACE, info, task_context) + persisted_info = connector._persist_installation_package.await_args.kwargs['install_info'] + assert persisted_info['plugin_version'] == '1.0.0' + assert persisted_info['_certification']['runtime_profile'] == 'shared-runtime-v1' + connector.handler.apply_plugin_installation.assert_awaited_once_with( + binding, artifact_package=package, enabled=True + ) + connector._refresh_runner_registry.assert_awaited_once() + assert task_context.metadata['progress_percent'] == 100 + if requested_version is not None: + # Confirmed migrations must fetch the reviewed release, never latest/MCP/skill. + assert len(requests) == 1 + else: + assert any(url.endswith('/versions') for url in requests) + + +def _connector(deployment: str, trusted_public_keys: dict[str, str]): + package_digest = 'a' * 64 + execution_context = ExecutionContext( + instance_uuid='instance-a', + workspace_uuid='workspace-a', + placement_generation=1, + ) + binding = InstallationBinding( + instance_uuid='instance-a', + workspace_uuid='workspace-a', + placement_generation=1, + installation_uuid='00000000-0000-4000-8000-000000000001', + runtime_revision=1, + artifact_digest=package_digest, + ) + app = SimpleNamespace( + instance_config=SimpleNamespace( + data={ + 'plugin': { + 'enable': True, + 'certification': {'trusted_public_keys': trusted_public_keys}, + } + } + ), + deployment=SimpleNamespace(mode=deployment), + logger=Mock(), + ) + connector = PluginRuntimeConnector(app, AsyncMock()) + connector.handler = SimpleNamespace( + register_installation_binding=Mock(), + apply_plugin_installation=AsyncMock(return_value={'state': 'running'}), + ) + connector._current_execution_context = AsyncMock(return_value=execution_context) + connector._store_artifact_package = AsyncMock() + connector._persist_installation_package = AsyncMock(return_value=(binding, None, False)) + connector._wait_for_installed_plugin_ready = AsyncMock() + return connector, execution_context, binding + + +def _archive(kind: str) -> tuple[bytes, dict[str, str]]: + manifest = { + 'metadata': {'author': 'certified', 'name': 'example', 'version': '1.0.0'}, + 'execution': {'sharedRuntime': 'shared-runtime-v1'}, + } + if kind == 'legacy': + manifest.pop('execution') + archive = io.BytesIO() + with zipfile.ZipFile(archive, 'w') as package: + package.writestr('manifest.yaml', yaml.safe_dump(manifest)) + raw_archive = archive.getvalue() + if kind == 'legacy': + return raw_archive, {} + + from langbot_plugin.certification import create_envelope, write_envelope + + signing_key = Ed25519PrivateKey.generate() + signed_archive = write_envelope( + raw_archive, + create_envelope( + raw_archive, + 'wrong-key' if kind == 'invalid_shared' else 'ephemeral', + signing_key.sign, + ), + ) + trusted_key = signing_key.public_key().public_bytes( + serialization.Encoding.Raw, + serialization.PublicFormat.Raw, + ) + return signed_archive, {'ephemeral': base64.b64encode(trusted_key).decode('ascii')} + + +def _normalized_digest(archive: bytes) -> str: + from langbot_plugin.certification import normalized_zip_digest + + return normalized_zip_digest(archive) diff --git a/tests/integration/plugin/test_rag_cancellation_protocol.py b/tests/integration/plugin/test_rag_cancellation_protocol.py new file mode 100644 index 000000000..9ae3623a0 --- /dev/null +++ b/tests/integration/plugin/test_rag_cancellation_protocol.py @@ -0,0 +1,64 @@ +"""Exercise actual installed SDK RPC cancellation, not a mocked cancel contract.""" + +import asyncio +from types import SimpleNamespace + +import pytest + +from langbot_plugin.entities.io.actions.enums import RuntimeToPluginAction +from langbot_plugin.runtime.io.handler import ActionResponse +from tests.integration.plugin.test_rag_file_transfer_protocol import BINDING, protocol_stack + + +@pytest.mark.asyncio +@pytest.mark.parametrize('disconnect', [False, True], ids=['cancel-waiter', 'close-host-connection']) +async def test_b1_host_cancellation_does_not_fence_remote_ingest(tmp_path, monkeypatch, disconnect): + async with protocol_stack(tmp_path, monkeypatch, 'shared', BINDING) as stack: + entered, release, finished = asyncio.Event(), asyncio.Event(), asyncio.Event() + vectors = set() + selected = SimpleNamespace(_runtime_plugin_handler=stack.bridge) + monkeypatch.setattr(stack.runtime.plugin_mgr, '_get_connected_rag_plugin', lambda *_: (selected, 'engine')) + + @stack.plugin.action(RuntimeToPluginAction.INGEST_DOCUMENT) + async def ingest(data): + entered.set() + await release.wait() + vectors.add('host-id') + finished.set() + return ActionResponse.success({'document_id': 'host-id', 'status': 'completed'}) + + @stack.plugin.action(RuntimeToPluginAction.DELETE_DOCUMENT) + async def delete(data): + vectors.discard(data['document_id']) + return ActionResponse.success({'success': True}) + + async def send_ingest(): + with stack.core.installation_scope(BINDING): + return await stack.core.rag_ingest_document('tester', 'engine', {}) + + local = asyncio.create_task(send_ingest()) + try: + await asyncio.wait_for(entered.wait(), 5) + if disconnect: + await stack.core.close() + result = await asyncio.gather(local, return_exceptions=True) + assert isinstance(result[0], Exception) + else: + local.cancel() + with pytest.raises(asyncio.CancelledError): + await local + assert not finished.is_set() + # The SDK has no ingest/delete barrier: a new caller can receive a + # confirmed deletion while the old plugin action can still write. + # Use the surviving Runtime->plugin connection after Host disconnect. + result = await stack.bridge.rag_delete_document('kb-a', 'host-id') + assert result == {'success': True} + assert vectors == set() + release.set() + await asyncio.wait_for(finished.wait(), 5) + assert vectors == {'host-id'} + finally: + release.set() + if not local.done(): + local.cancel() + await asyncio.gather(local, return_exceptions=True) diff --git a/tests/unit_tests/core/test_load_config.py b/tests/unit_tests/core/test_load_config.py index 03dd7cc5d..051597875 100644 --- a/tests/unit_tests/core/test_load_config.py +++ b/tests/unit_tests/core/test_load_config.py @@ -427,3 +427,27 @@ class TestApplyEnvOverridesToConfig: result = load_config._apply_env_overrides_to_config(cfg) assert result['api']['extra_webhook_prefix'] == 'https://extra.example.com' + + +class TestCertificationKeyRingEnv: + def test_applies_string_key_mapping_from_strict_json(self): + load_config = get_load_config_module() + cfg = load_config._complete_runtime_policy_defaults({}) + with patch.dict( + os.environ, + {'PLUGIN__CERTIFICATION__TRUSTED_PUBLIC_KEYS_JSON': '{"ed25519:issuer":"YWJj"}'}, + clear=True, + ): + result = load_config._apply_certification_key_ring_env(cfg) + assert result['plugin']['certification']['trusted_public_keys'] == {'ed25519:issuer': 'YWJj'} + + def test_rejects_malformed_or_non_mapping_key_ring(self): + load_config = get_load_config_module() + for value in ('not-json', '[]', '{"":"YWJj"}', '{"ed25519:issuer": 1}'): + cfg = load_config._complete_runtime_policy_defaults({}) + with patch.dict(os.environ, {'PLUGIN__CERTIFICATION__TRUSTED_PUBLIC_KEYS_JSON': value}, clear=True): + try: + load_config._apply_certification_key_ring_env(cfg) + except ValueError: + continue + raise AssertionError(f'invalid certification key ring was accepted: {value!r}') diff --git a/tests/unit_tests/plugin/test_certified_plugin_policy.py b/tests/unit_tests/plugin/test_certified_plugin_policy.py new file mode 100644 index 000000000..9e1314c24 --- /dev/null +++ b/tests/unit_tests/plugin/test_certified_plugin_policy.py @@ -0,0 +1,159 @@ +from __future__ import annotations + +import hashlib +import io +import zipfile + +import pytest + + +@pytest.mark.parametrize( + ('deployment', 'certificate', 'force', 'expected_disposition', 'expected_code'), + [ + ('cloud', ('valid', 'shared-runtime-v1'), False, 'shared_eligible', 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE'), + ('cloud', ('absent', None), False, 'rejected', 'CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_REQUIRED'), + ('cloud', ('malformed', None), False, 'rejected', 'CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_INVALID'), + ('cloud', ('invalid', 'shared-runtime-v1'), True, 'rejected', 'CERTIFIED_PLUGIN_CLOUD_CERTIFICATE_INVALID'), + ('oss', ('absent', None), False, 'dedicated_allowed', 'CERTIFIED_PLUGIN_OSS_LEGACY_DEDICATED'), + ('oss', ('valid', 'shared-runtime-v1'), False, 'shared_eligible', 'CERTIFIED_PLUGIN_SHARED_ELIGIBLE'), + ( + 'oss', + ('invalid', 'shared-runtime-v1'), + False, + 'administrator_force_required', + 'CERTIFIED_PLUGIN_OSS_FORCE_REQUIRED', + ), + ('oss', ('invalid', 'shared-runtime-v1'), True, 'dedicated_allowed', 'CERTIFIED_PLUGIN_OSS_FORCED_DEDICATED'), + ], +) +def test_admission_policy_enforces_certification_matrix( + deployment: str, + certificate: tuple[str, str | None], + force: bool, + expected_disposition: str, + expected_code: str, +) -> None: + from langbot.pkg.plugin.certification import ( + CertificateFacts, + CertificateVerification, + PluginCertificationFacts, + decide_plugin_admission, + ) + + verification, runtime_profile = certificate + facts = PluginCertificationFacts( + installation_uuid='00000000-0000-4000-8000-000000000001', + artifact_digest='a' * 64, + certificate=CertificateFacts( + verification=CertificateVerification(verification), + runtime_profile=runtime_profile, + ), + ) + + decision = decide_plugin_admission( + deployment=deployment, + facts=facts, + administrator_force=force, + ) + + assert decision.disposition.value == expected_disposition + assert decision.code.value == expected_code + assert decision.runtime_profile == ( + 'shared-runtime-v1' if expected_disposition == 'shared_eligible' else 'dedicated' + ) + + +def test_archive_inspection_preserves_legacy_tuple_and_exposes_certificate_facts() -> None: + from langbot.pkg.plugin.archive import ( + ArchiveCertificateState, + inspect_plugin_archive, + inspect_plugin_archive_metadata, + ) + + manifest = { + 'kind': 'Plugin', + 'metadata': {'name': 'example'}, + 'certification': { + 'runtime_profile': 'shared-runtime-v1', + 'certificate': {'issuer': 'sdk-test', 'signature': 'not-verified-by-core'}, + }, + } + archive_bytes = _archive_bytes(manifest) + + inspection = inspect_plugin_archive(archive_bytes) + + assert inspection.artifact_digest == hashlib.sha256(archive_bytes).hexdigest() + assert inspection.certificate.state is ArchiveCertificateState.DECLARED + assert inspection.certificate.runtime_profile == 'shared-runtime-v1' + assert inspection.certificate.payload == manifest['certification']['certificate'] + assert inspect_plugin_archive_metadata(archive_bytes) == ( + inspection.manifest, + inspection.requirements, + inspection.names, + ) + + +@pytest.mark.parametrize( + ('certification', 'expected_state'), + [ + (None, 'absent'), + ({'runtime_profile': 123, 'certificate': {}}, 'malformed'), + ({'runtime_profile': 'shared-runtime-v1'}, 'malformed'), + ], +) +def test_archive_inspection_reports_nonverifying_certificate_states( + certification: object, + expected_state: str, +) -> None: + from langbot.pkg.plugin.archive import inspect_plugin_archive + + manifest: dict[str, object] = {'kind': 'Plugin', 'metadata': {'name': 'example'}} + if certification is not None: + manifest['certification'] = certification + + inspection = inspect_plugin_archive(_archive_bytes(manifest)) + + assert inspection.certificate.state.value == expected_state + + +@pytest.mark.parametrize( + ('verification', 'expected_visibility'), + [ + ('valid', 'tenant_scoped'), + ('absent', 'detailed_process'), + ('malformed', 'detailed_process'), + ('invalid', 'detailed_process'), + ], +) +def test_log_visibility_policy_only_scopes_valid_shared_certifications( + verification: str, + expected_visibility: str, +) -> None: + from langbot.pkg.plugin.certification import ( + CertificateFacts, + CertificateVerification, + PluginCertificationFacts, + decide_plugin_log_visibility, + ) + + facts = PluginCertificationFacts( + installation_uuid='00000000-0000-4000-8000-000000000001', + artifact_digest='a' * 64, + certificate=CertificateFacts( + verification=CertificateVerification(verification), + runtime_profile='shared-runtime-v1', + ), + ) + + visibility = decide_plugin_log_visibility(facts) + + assert visibility.value == expected_visibility + + +def _archive_bytes(manifest: dict[str, object]) -> bytes: + buffer = io.BytesIO() + with zipfile.ZipFile(buffer, 'w') as archive: + import yaml + + archive.writestr('manifest.yaml', yaml.safe_dump(manifest)) + return buffer.getvalue() diff --git a/tests/unit_tests/plugin/test_connector_reconcile.py b/tests/unit_tests/plugin/test_connector_reconcile.py index 605bcba88..d5292f8f3 100644 --- a/tests/unit_tests/plugin/test_connector_reconcile.py +++ b/tests/unit_tests/plugin/test_connector_reconcile.py @@ -40,6 +40,17 @@ def execution_binding(workspace_uuid: str, generation: int = 1) -> SimpleNamespa ) +def mock_archive_admission(connector: PluginRuntimeConnector, digest: str) -> None: + # These lifecycle tests use opaque package bytes. Real certificate admission + # is exercised by integration/plugin/test_certified_plugin_admission.py. + connector._admit_plugin_archive = Mock( + side_effect=lambda _package, info: ( + {**info, '_certification': {'normalized_digest': digest}}, + SimpleNamespace(for_installation=lambda _uuid: SimpleNamespace(artifact_digest=digest)), + ) + ) + + def plugin_setting( workspace_suffix: str, artifact_digest: str, @@ -275,6 +286,12 @@ async def test_local_install_persists_verified_package_before_runtime_apply(): connector._inspect_plugin_package = Mock(return_value=('author', 'plugin')) connector._store_artifact_package = AsyncMock() connector._persist_installation_package = AsyncMock(return_value=(binding, None, False)) + connector._admit_plugin_archive = Mock( + return_value=( + {'_certification': {'normalized_digest': digest}}, + SimpleNamespace(for_installation=lambda _installation_uuid: SimpleNamespace(artifact_digest=digest)), + ) + ) connector._wait_for_installed_plugin_ready = AsyncMock() await connector.install_plugin( @@ -351,6 +368,7 @@ async def test_marketplace_upgrade_reports_multistep_progress(): observed_actions.append(task_context.current_action) connector._download_marketplace_package = AsyncMock(side_effect=download) + mock_archive_admission(connector, digest) connector._inspect_plugin_package = Mock(side_effect=inspect) connector._store_artifact_package = AsyncMock(side_effect=store) connector._persist_installation_package = AsyncMock(side_effect=persist) @@ -403,6 +421,7 @@ async def test_workspace_reads_do_not_wait_for_an_installation_apply(): connector.handler = runtime_handler() connector._current_execution_context = AsyncMock(return_value=execution_context) connector._validate_execution_context = AsyncMock(return_value=execution_context) + mock_archive_admission(connector, digest) connector._inspect_plugin_package = Mock(return_value=('author', 'plugin')) connector._store_artifact_package = AsyncMock() connector._persist_installation_package = AsyncMock(return_value=(binding, None, False)) @@ -471,6 +490,7 @@ async def test_local_install_cleans_untracked_legacy_plugin_before_runtime_apply connector._persist_installation_package = AsyncMock(return_value=(binding, None, False)) connector._wait_for_installed_plugin_ready = AsyncMock() events: list[str] = [] + mock_archive_admission(connector, digest) async def delete_legacy_plugin(plugin_author: str, plugin_name: str): assert (plugin_author, plugin_name) == ('author', 'plugin') diff --git a/tests/unit_tests/plugin/test_marketplace_plugin_version.py b/tests/unit_tests/plugin/test_marketplace_plugin_version.py new file mode 100644 index 000000000..bddfe8c8d --- /dev/null +++ b/tests/unit_tests/plugin/test_marketplace_plugin_version.py @@ -0,0 +1,43 @@ +import pytest + +from langbot.pkg.plugin.connector import _select_marketplace_plugin_version + + +@pytest.mark.parametrize( + ('requested_version', 'expected'), + [ + (None, '0.1.4'), + ('0.1.3', '0.1.3'), + ], +) +def test_select_marketplace_plugin_version(requested_version, expected): + assert ( + _select_marketplace_plugin_version( + [{'version': '0.1.4'}, {'version': '0.1.3'}], + requested_version=requested_version, + plugin_author='langbot-team', + plugin_name='RunnerDemo', + ) + == expected + ) + + +def test_select_marketplace_plugin_version_rejects_requested_missing_version(): + with pytest.raises(ValueError, match='version 0.1.2 is not available'): + _select_marketplace_plugin_version( + [{'version': '0.1.4'}], + requested_version='0.1.2', + plugin_author='langbot-team', + plugin_name='RunnerDemo', + ) + + +@pytest.mark.parametrize('versions', [[], [{'unexpected': 'value'}], 'not-a-list']) +def test_select_marketplace_plugin_version_rejects_invalid_latest(versions): + with pytest.raises(ValueError, match='has no versions'): + _select_marketplace_plugin_version( + versions, + requested_version=None, + plugin_author='langbot-team', + plugin_name='RunnerDemo', + ) diff --git a/tests/unit_tests/rag/test_file_storage.py b/tests/unit_tests/rag/test_file_storage.py index 7b464d502..8886e221a 100644 --- a/tests/unit_tests/rag/test_file_storage.py +++ b/tests/unit_tests/rag/test_file_storage.py @@ -90,7 +90,7 @@ class TestStoreFile: def create_user_task(coro, **kwargs): coro.close() - return SimpleNamespace(id='task-1', kwargs=kwargs) + return SimpleNamespace(id='task-1', kwargs=kwargs, task=Mock()) kb.ap.task_mgr.create_user_task = Mock(side_effect=create_user_task) @@ -279,7 +279,7 @@ class TestStoreFileTask: kb._assert_execution_context = AsyncMock(side_effect=assert_execution_context) kb._set_file_status = AsyncMock(side_effect=[True, True]) - kb._ingest_document = AsyncMock(return_value={'status': 'completed'}) + kb._ingest_document = AsyncMock(return_value={'status': 'completed', 'document_id': 'file-uuid'}) object_key = _upload_key('scoped.pdf') file_obj = SimpleNamespace(uuid='file-uuid', file_name=object_key, extension='pdf') @@ -290,7 +290,7 @@ class TestStoreFileTask: @pytest.mark.asyncio async def test_store_file_task_marks_completed_and_cleans_storage(self): kb = _make_kb() - kb._ingest_document = AsyncMock(return_value={'status': 'completed'}) + kb._ingest_document = AsyncMock(return_value={'status': 'completed', 'document_id': 'file-uuid'}) object_key = _upload_key('test.pdf') file_obj = SimpleNamespace(uuid='file-uuid', file_name=object_key, extension='pdf') task_context = Mock() @@ -306,7 +306,9 @@ class TestStoreFileTask: @pytest.mark.asyncio async def test_store_file_task_marks_failed_and_cleans_storage(self): kb = _make_kb() - kb._ingest_document = AsyncMock(return_value={'status': 'failed', 'error_message': 'parser failed'}) + kb._ingest_document = AsyncMock( + return_value={'status': 'failed', 'error_message': 'parser failed', 'document_id': 'file-uuid'} + ) object_key = _upload_key('bad.pdf') file_obj = SimpleNamespace(uuid='file-uuid', file_name=object_key, extension='pdf') task_context = Mock() diff --git a/tests/unit_tests/rag/test_kbmgr.py b/tests/unit_tests/rag/test_kbmgr.py index f43814870..9ea0c54c9 100644 --- a/tests/unit_tests/rag/test_kbmgr.py +++ b/tests/unit_tests/rag/test_kbmgr.py @@ -288,7 +288,9 @@ async def test_ingestion_payload_uses_host_owned_kb_collection(): async def test_delete_file_checks_workspace_and_parent_before_plugin_call(): app = _app() runtime = RuntimeKnowledgeBase(app, _entity(), CONTEXT_A) - app.persistence_mgr.execute_async.return_value = _Result(first=('file-a',)) + app.persistence_mgr.execute_async.return_value = _Result( + first=SimpleNamespace(uuid='file-a', status='completed', engine_document_id=None) + ) await runtime.delete_file(CONTEXT_A, 'file-a') app.plugin_connector.call_rag_delete_document.assert_awaited_once_with( @@ -404,7 +406,8 @@ class TestRAGManagerCreateKnowledgeBase: ) assert manager.knowledge_bases == {} - assert app.persistence_mgr.execute_async.await_count == 2 + # Insert, interrupted-ingestion reconciliation, rollback delete. + assert app.persistence_mgr.execute_async.await_count == 3 @pytest.mark.asyncio async def test_sets_default_retrieval_settings(self): @@ -476,7 +479,9 @@ class TestRuntimeKnowledgeBaseDeleteFile: @pytest.mark.asyncio async def test_delete_file_calls_plugin_and_db(self): app = _app() - app.persistence_mgr.execute_async.return_value = _Result(first=('file-uuid',)) + app.persistence_mgr.execute_async.return_value = _Result( + first=SimpleNamespace(uuid='file-uuid', status='completed', engine_document_id=None) + ) await RuntimeKnowledgeBase(app, _entity(), CONTEXT_A).delete_file( CONTEXT_A, @@ -534,7 +539,7 @@ class TestRAGManagerLoadKnowledgeBasesFromDB: } @pytest.mark.asyncio - async def test_cloud_startup_reuses_validated_binding(self): + async def test_cloud_startup_revalidates_binding_before_recovery_write(self): class TenantUow: async def __aenter__(self): return self @@ -554,15 +559,13 @@ class TestRAGManagerLoadKnowledgeBasesFromDB: app.persistence_mgr.tenant_uow = lambda _workspace_uuid: TenantUow() app.persistence_mgr.execute_async.return_value = _Result([_entity()]) app.workspace_service.list_active_execution_bindings = AsyncMock(return_value=[binding]) - app.workspace_service.get_execution_binding = AsyncMock( - side_effect=AssertionError('startup RAG loader repeated a validated binding lookup') - ) + app.workspace_service.get_execution_binding = AsyncMock(return_value=binding) manager = RAGManager(app) await manager.load_knowledge_bases_from_db() assert set(manager.knowledge_bases) == {('workspace-a', 'kb-a')} - app.workspace_service.get_execution_binding.assert_not_awaited() + app.workspace_service.get_execution_binding.assert_awaited_once_with('workspace-a', expected_generation=5) @pytest.mark.asyncio async def test_handles_load_error_gracefully(self): diff --git a/uv.lock b/uv.lock index a1c69bdcb..65cf79f79 100644 --- a/uv.lock +++ b/uv.lock @@ -2180,7 +2180,7 @@ requires-dist = [ { name = "ebooklib", specifier = ">=0.18" }, { name = "gewechat-client", specifier = ">=0.1.5" }, { name = "html2text", specifier = ">=2024.2.26" }, - { name = "langbot-plugin", specifier = "==0.6.0b3" }, + { name = "langbot-plugin", specifier = "==0.6.0b5" }, { name = "langchain", specifier = ">=1.3.9" }, { name = "langchain-core", specifier = ">=1.3.3" }, { name = "langchain-text-splitters", specifier = ">=1.1.2" }, @@ -2250,7 +2250,7 @@ dev = [ [[package]] name = "langbot-plugin" -version = "0.6.0b3" +version = "0.6.0b5" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "aiofiles" }, @@ -2271,9 +2271,9 @@ dependencies = [ { name = "watchdog" }, { name = "websockets" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/dc/6c/198eee31ef8b1ada209cc3c1588637ae3649f16fb4c0c41acc8296559cfb/langbot_plugin-0.6.0b3.tar.gz", hash = "sha256:2becc1de2b2543c29609c88b7cd2de9267b94d11d2d055dfd46ecfb9e9bbd496", size = 607092, upload-time = "2026-09-18T16:20:37.564Z" } +sdist = { url = "https://files.pythonhosted.org/packages/7d/27/c23c5bab4137e755096f7ad32afcec834cb9b2e08129beb8bf797c1b8c5b/langbot_plugin-0.6.0b5.tar.gz", hash = "sha256:8cd1024ec1a8a131b8afe48dd3ff6c002626d6ab885ab99b8078985e4a8f53be", size = 613773, upload-time = "2026-09-20T10:18:18.171Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/7c/f0/41ffb773c114924d95cc78824d218d76a7d7973c79dc3b43a12d4e93992b/langbot_plugin-0.6.0b3-py3-none-any.whl", hash = "sha256:a4aa0c7427267685fed8b96c8f02457dad81728758f75373816f0146c4f72c8c", size = 404923, upload-time = "2026-09-18T16:20:35.972Z" }, + { url = "https://files.pythonhosted.org/packages/64/fc/b0c7c009e166650d82bcb58784e89a64ec04de5156f95da7fb7d59d2c4a3/langbot_plugin-0.6.0b5-py3-none-any.whl", hash = "sha256:15c5a60765db34e5a0216eab205178fb06c6447d286f648008aa89207e75d667", size = 408564, upload-time = "2026-09-20T10:18:16.967Z" }, ] [[package]]