Compare commits

...

2 Commits

22 changed files with 1656 additions and 64 deletions

View File

@ -5,6 +5,11 @@
> 所以不是每个版本号都有条目。条目格式 `## <版本> — <日期>`,新条目加在最上面。
> 工程口径的完整记录见 `PROGRESS.md` / git log。
## 0.64.0 — 2026-08-12
- 已发布产物现在拥有稳定身份:重命名和移动后,历史对话中的产物仍可正常预览和下载;复制产物会创建可独立管理的副本。
- 删除已发布产物时,文件会进入平台隐藏回收区而不是立即永久丢失;普通工作文件仍按原方式删除。
## 0.63.11 — 2026-08-12
- 修复 PDF 和 PPT 连续预览无法正常滚动的问题,现在可直接上下滚动浏览全部页面。

View File

@ -173,7 +173,7 @@ Eval 与生产 core 解耦,通过现有 `/v1` API 创建专用任务、监听
**新对话入口(0.60)**:登录未选 task 与左栏「+ 新对话」共用同一前端草稿页,先选择已有 working_dir 或输入新目录名,再直接写消息;草稿不落 DB首发时才 `POST /v1/tasks`,避免空 task 堆积。创建请求省略/留空 name 时必须显式给 working_dir后端据此判定自动命名以「新对话」占位并置一次性 `auto_title_pending`;显式 name 的旧调用继续视为人工标题working_dir 仍可省略并 fallback 到 name`auto_title` 字段只作兼容保留。首条消息并行触发短标题调用,结果只改 `tasks.name`、绝不改 working_dir人工 PATCH name 同时清 pending条件 UPDATE 保证在途标题也不能覆盖用户命名。原完整创建表单保留为「自定义」入口UI 同样要求明确选择 working_dirname 可选,并可预设 description/skill/model。标题是 UI 元数据辅助调用,记 `usage_events.kind="task_title"`;模型调用失败时以首条消息第一行生成本地兜底标题,不阻塞主 run也不把 pending 留给后续消息误命名。
**对话产物引用(0025)**:真实文件仍是事实源,不建 artifacts 表;`messages.artifact_refs` 只保存可重建的轻量 UI 元数据,规范路径以该 task 的**当前 working_dir 为根**,形如 `{version:1, scope:"working_dir", path:"reports/a.pdf", label?:"最终报告"}`。预览/下载走 task-scoped 文件 API服务端用 task 当前 `working_dir` 解析因此顶层工作目录改名后历史卡片仍有效。普通源码树、中间文件和配套资源只留文件面板agent 仅用 `publish_artifacts` 显式提升少量最终文件,单条消息最多 10 个,图像/视频/Office 转 PDF 等成品工具可自动提升。`NULL` 表示迁移前旧消息,前端继续使用正文路径抽取,并在 task-scoped API 上启用只读兼容链(旧 user-root 含义→原样 task-relative→去掉旧目录前缀新消息写 `[]` 或结构化列表,停止启发式抽取,避免重复卡片与误识别。文件在 working_dir 内再次移动或删除后引用可失效,这是 FS 事实源语义,不复制文件、不引不可变对象存储
**对话产物与生命周期(0025/0028)**:真实文件仍是内容事实源;`artifacts` 表记录已发布产物的稳定身份和生命周期,包含 user-root 相对当前路径、来源 task、复制来源、哈希/大小及 active/deleted、回收路径。新 `messages.artifact_refs` 使用 `{version:2, artifact_id, scope:"working_dir", path:"reports/a.pdf", label?:"最终报告"}``path` 是兼容快照,预览/下载优先按 `artifact_id` 找当前路径因此移动或重命名后历史卡片仍有效。version 1 和 `NULL` 旧消息继续走原 task-scoped 兼容链。普通源码树、中间文件和配套资源不登记agent 仅用 `publish_artifacts` 显式提升少量最终文件。移动保持身份,复制为每个副本创建新身份并记录直接来源;删除将文件移入 `.zcbot_artifact_trash/` 并软删记录,普通文件仍物理删除
### 7.2 资源模型(/v1)
@ -285,7 +285,7 @@ scheduled_jobs(§8.5) channel_bindings(§8.7,判别列+JSONB)
- **path-as-identity 而非 folder_id**:folder 真实存在于 FS,folder_id 是第二份 source of truth;rename 走 DB-aware 同事务 cascade。
- **files API 单一 mutation 入口**(2026-05-18):"顶层目录分支"从数据状态派生而非客户端意图,放服务端才有强制力;双命名空间(/folders vs /files)把分支搬给 client,失强制力且端点翻倍。
- **task 软删除(2026-06-17 推翻 hard cascade)**:公测后对话轨迹是训练/研究语料,`deleted_at` 置位 + restore,避免用户误删立即永久丢失。**当前实现仍无限期保留软删数据**,物理清理仅有管理员手段;后续生命周期已定为“软删除后保留 30 天再物理清理”(待容量信号实施,见 §8.5),届时恢复能力明确限于宽限期内。
- **文件留存(设计已定,实现待办)**:用户文件在 FS,删除/覆盖即字节丢。方案=① restic/borg 定时增量备份做地基(与应用解耦,新端点自动覆盖,捕获删除+覆盖+成品)+ ② 应用层 `data_events` 事件日志(补用户意图语义)。**不选**每个删除端点内联 copytree:横切关注点手写 N 处必漏。起步同盘(不防整盘损坏,已知边界)
- **文件留存**:普通用户文件仍是 FS 直接删除;已写入结构化 `messages.artifact_refs` 的已发布产物,在统一 files delete 入口删除时原子移动到用户根目录的平台隐藏区 `.zcbot_artifact_trash/`,原路径与历史卡片立即表现为已删除。递归目录只回收其中 artifact其他文件照常删除回收内容仍计入用户配额。该机制只防误删产物不防覆盖、agent/shell 绕过 files API 或整盘损坏。完整地基仍采用 restic/borg 定时增量备份(与应用解耦,捕获删除+覆盖+所有写入口),后续容量需要时再补 `data_events` 用户意图事件和回收区清理/恢复管理
- **0004 删 runs/usage_events 旧表**:只写不读的死代码;代价是失历史 run 元数据,真要细粒度审计再补(届时是新需求非技术债)。
- **本地也用 PG 不用 SQLite**:dogfood ≡ 真实路径;Docker 已是必然依赖;双 adapter 维护税 > 一次性配置。
- **API-only,UI 由 platform 实现**(2026-05-15):本仓库再维护一套 UI 是双套浪费;SSE payload 从 HTML 切 JSON;沉淀的 sink/broker/路径安全全保留。**dev SPA 留一份**作 dogfood 主路径(SSE 调试 curl/Swagger 都覆盖不了)。

View File

@ -2,7 +2,7 @@
> 配合 `DESIGN.md`。本文件只记 phase 状态、决策偏差、文件量、下一步。每条 1-2 句:做了啥 + 关键判断;细节查 `git log` / `git diff` / `DESIGN §7.9`
最后更新:2026-08-12(PDF/PPT 连续预览滚动布局修复,bump 0.63.11)
最后更新:2026-08-12(artifact 稳定身份与隐藏回收,bump 0.64.0)
---
@ -23,6 +23,8 @@
### 2026-08-12
- **08-12 / 0.64.0 / artifact 稳定身份 + 隐藏回收**:新增 0028 `artifacts` 生命周期表并回填存量结构化引用;新发布消息写带 `artifact_id` 的 v2 引用,移动/重命名保持身份,复制创建独立身份并记录来源,历史卡片按身份解析最新路径。删除已发布产物时移动到用户隐藏目录 `.zcbot_artifact_trash/` 并软删记录普通文件仍物理删除Python 35 项(测试库门控 1 skip、Node 前端 11 项、mypy、Alembic 单 head、编译、Ruff 致命规则及 diff 检查通过,未连接或写入生产 DB。
- **08-12 / 0.63.11 / PDF/PPT 连续预览滚动布局修复**:补齐 PDF 预览根容器的纵向 Flex 布局与可收缩高度约束,使内部连续页 viewport 获得实际可用高度和独立滚动区域,不再被外层 `overflow:hidden` 裁断回归测试同步锁定根容器、viewport 与页列表三层布局契约。Node 前端预览 11 项、JavaScript 语法及 diff 检查通过;无 schema、migration、HTTP API、依赖或运行方式变化未连接生产 DB。
### 2026-08-11

1
RUN.md
View File

@ -151,6 +151,7 @@
- **channel 长会话上下文(微信/企业微信通用,0019)**:常驻会话不再无限膨胀。① **自动分段**——入站时距上次消息超过 `config.json``channel.session_gap_hours`(默 **6** 小时,设 `<=0` 关闭)→ 软重置:只把「最后一条 user 消息起」喂模型(保留上一轮做续聊锚点),之前的历史仍全留 DB,网页端照旧翻完整记录;② **手动新话题**——用户在微信/企业微信里直接发「新话题 / 新会话 / `/new` / 清空上下文」→ 硬重置,彻底从零(回执提示已归档)。两者都**不删任何消息**,只移动「喂给模型的窗口起点」`tasks.context_base_idx`。网页端「清空对话」(`POST /v1/tasks/{id}/clear`)仍整清并把 base 归 0。需 `main.py db upgrade head` 带上 `0019`
- **PG**:`ZCBOT_DB_URL` 必填。本地 docker compose / 远端 dev / 生产任选;未设置时启动清晰报错,不引导 docker(§7.4)。
- **OpenAPI / MCP 外部系统**:① `.env` 配置独立的 `ZCBOT_CREDENTIAL_MASTER_KEY`,可选 `ZCBOT_CREDENTIAL_KEY_ID` 标识当前密钥;轮换时把旧 key 以 JSON 对象放入 `ZCBOT_CREDENTIAL_PREVIOUS_KEYS`,待用户凭据完成重写后再移除。② 执行 `main.py db upgrade head`。③ admin 进入管理后台「外部系统」,选择通用 OpenAPI 或通用 MCP具体 MES/ERP/LIMS 都作为数据库 definition 配置,不新增专用 provider。MCP 填写与登录 Base URL 同源的 Streamable HTTP URL可选填写期望 Server 名称;连接后以 `tools/list` 为事实源。④ 普通用户点击左栏 **「外部」**,页面按 definition 动态显示用户名密码、API Key 或 Bearer Token目标、Server 身份、登录、认证绑定或 TLS 变化后保留密文并暂停调用重新测试成功后恢复。OpenAPI spec 和 MCP tool catalog 只在进程内按连接身份有界缓存,登录与业务响应均限长,普通用户和模型不能传任意 URL。
- **Artifact 生命周期(0028)**:部署本版本必须先执行 `.venv/Scripts/python.exe main.py db upgrade head`。migration 会从存量 `messages.artifact_refs` 回填 active artifact 身份;删除后的已发布产物保存在用户根目录隐藏区 `.zcbot_artifact_trash/`,默认不自动清理且继续计入磁盘配额。普通文件删除语义不变。
- **旧 `factory_mes` definition 一次性转换**:新版代码不再识别 `factory_mes`;部署时保持旧服务进程运行,先从新代码目录执行数据脚本,转换成功后再重启到新版。脚本不加载 `.env`、不读取 `ZCBOT_DB_URL`,只认显式的 `ZCBOT_MIGRATION_DB_URL`;默认 dry-run检查同名冲突与配置合法性。确认输出后加 `--apply`,脚本把 definition 转为 `generic_openapi + query`、物化 JWT/提示/只读 POST 配置并同步 active connection revision不解密或改写用户凭据。
```powershell
$env:ZCBOT_MIGRATION_DB_URL="postgresql+psycopg://user:pass@host:5432/zcbot"

View File

@ -1,3 +1,3 @@
# zcbot 版本号单一事实源:web/app.py 的 FastAPI version、/healthz 返回、前端展示都引这里。
# 改版本只动这一行。
__version__ = "0.63.11"
__version__ = "0.64.0"

View File

@ -674,7 +674,7 @@ def build_agent(
agent = AgentLoop(
llm, executor, session, caps,
user_id=uid, working_dir=working_dir_path, sink=sink,
user_id=uid, working_dir=working_dir_path, user_root=ur_path, sink=sink,
skill_model_switch=_skill_model_switch,
deferred_actions=deferred_actions,
)

208
core/artifact_lifecycle.py Normal file
View File

@ -0,0 +1,208 @@
"""Database-backed lifecycle operations for published workspace artifacts."""
from __future__ import annotations
import hashlib
import mimetypes
import os
import shutil
from datetime import datetime, timezone
from pathlib import Path
from uuid import UUID, uuid4
from sqlalchemy import select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from .artifacts import ARTIFACT_TRASH_DIR, ArtifactRef
from .storage import session_scope
from .storage.models import Artifact
def _rel(root: Path, path: Path) -> str:
return Path(path).resolve().relative_to(Path(root).resolve()).as_posix()
def _hash_file(path: Path) -> str:
digest = hashlib.sha256()
with Path(path).open("rb") as stream:
for chunk in iter(lambda: stream.read(1024 * 1024), b""):
digest.update(chunk)
return digest.hexdigest()
def register_published_artifacts(
*,
user_id: UUID,
task_id: UUID,
user_root: Path,
working_dir: Path,
refs: tuple[dict, ...],
) -> tuple[dict, ...]:
"""Upsert active artifact identities and return version-2 message refs."""
root = Path(user_root).resolve()
wd = Path(working_dir).resolve()
output: list[dict] = []
with session_scope() as session:
for ref in refs:
task_path = str(ref.get("path") or "")
path = (wd / Path(task_path)).resolve()
path.relative_to(wd)
if not path.is_file():
continue
current_path = _rel(root, path)
label = str(ref.get("label") or "")
media_type = mimetypes.guess_type(path.name)[0]
size_bytes = path.stat().st_size
content_sha256 = _hash_file(path)
statement = pg_insert(Artifact).values(
user_id=user_id,
origin_task_id=task_id,
current_path=current_path,
label=label,
media_type=media_type,
size_bytes=size_bytes,
content_sha256=content_sha256,
).on_conflict_do_update(
index_elements=[Artifact.user_id, Artifact.current_path],
index_where=Artifact.status == "active",
set_={
"label": label,
"media_type": media_type,
"size_bytes": size_bytes,
"content_sha256": content_sha256,
"updated_at": datetime.now(timezone.utc),
},
).returning(Artifact.artifact_id)
artifact_id = session.execute(statement).scalar_one()
output.append(ArtifactRef(
path=task_path,
label=label,
artifact_id=artifact_id,
version=2,
).as_dict())
return tuple(output)
def rename_active_artifacts(
*,
user_id: UUID,
user_root: Path,
old_path: Path,
new_path: Path,
) -> int:
"""Rewrite active artifact paths for one file or directory subtree."""
root = Path(user_root).resolve()
old_rel = _rel(root, old_path)
new_rel = _rel(root, new_path)
with session_scope() as session:
rows = session.execute(
select(Artifact).where(
Artifact.user_id == user_id,
Artifact.status == "active",
)
).scalars().all()
changed = 0
for row in rows:
if row.current_path == old_rel:
suffix = ""
elif row.current_path.startswith(old_rel + "/"):
suffix = row.current_path[len(old_rel):]
else:
continue
row.current_path = new_rel + suffix
changed += 1
return changed
def copy_active_artifacts(
*,
user_id: UUID,
user_root: Path,
source: Path,
target: Path,
) -> int:
"""Create independent artifact identities for copied files."""
root = Path(user_root).resolve()
source_rel = _rel(root, source)
target_rel = _rel(root, target)
with session_scope() as session:
sources = session.execute(
select(Artifact).where(
Artifact.user_id == user_id,
Artifact.status == "active",
)
).scalars().all()
created = 0
for original in sources:
if original.current_path == source_rel:
suffix = ""
elif original.current_path.startswith(source_rel + "/"):
suffix = original.current_path[len(source_rel):]
else:
continue
copied_path = target_rel + suffix
session.add(Artifact(
user_id=user_id,
origin_task_id=original.origin_task_id,
copied_from_artifact_id=original.artifact_id,
current_path=copied_path,
label=(original.label + " 副本").strip(),
media_type=original.media_type,
size_bytes=original.size_bytes,
content_sha256=original.content_sha256,
))
created += 1
return created
def trash_active_artifacts(
*,
user_id: UUID,
user_root: Path,
target: Path,
) -> int:
"""Move matching active artifacts to hidden trash and mark them deleted."""
root = Path(user_root).resolve()
source = Path(target).resolve()
source_rel = _rel(root, source)
entry = (
root / ARTIFACT_TRASH_DIR
/ datetime.now(timezone.utc).strftime("%Y/%m/%d")
/ uuid4().hex
)
moved: list[tuple[Path, Path]] = []
try:
with session_scope() as session:
rows = session.execute(
select(Artifact).where(
Artifact.user_id == user_id,
Artifact.status == "active",
).with_for_update()
).scalars().all()
matches = [
row for row in rows
if row.current_path == source_rel
or row.current_path.startswith(source_rel + "/")
]
if not matches:
return 0
for row in matches:
original = root / Path(row.current_path)
if not original.is_file():
continue
destination = entry / "files" / Path(row.current_path)
destination.parent.mkdir(parents=True, exist_ok=True)
os.replace(original, destination)
moved.append((original, destination))
row.status = "deleted"
row.deleted_at = datetime.now(timezone.utc)
row.trash_path = _rel(root, destination)
return len(moved)
except Exception:
for original, destination in reversed(moved):
try:
original.parent.mkdir(parents=True, exist_ok=True)
os.replace(destination, original)
except OSError:
pass
shutil.rmtree(entry, ignore_errors=True)
raise

View File

@ -1,18 +1,21 @@
"""Task artifact references and working-dir scoped path resolution.
"""Artifact message references and working-dir scoped path resolution.
Files remain the source of truth. Artifact refs are small, rebuildable UI metadata:
they identify a user-facing deliverable relative to a task's mutable working_dir.
Files remain the content source of truth. Version-2 refs carry a stable database identity
plus a task-relative path snapshot; version-1 path-only refs remain readable.
"""
from __future__ import annotations
from collections.abc import Iterable
from dataclasses import dataclass
from pathlib import Path
from typing import Iterable, Optional
from typing import Optional
from uuid import UUID
ARTIFACT_REF_VERSION = 1
ARTIFACT_REF_CURRENT_VERSION = 2
MAX_ARTIFACTS_PER_MESSAGE = 10
_CONTAINER_ROOT = Path("/workspace")
ARTIFACT_TRASH_DIR = ".zcbot_artifact_trash"
class ArtifactPathError(ValueError):
@ -23,6 +26,7 @@ class ArtifactPathError(ValueError):
class ArtifactRef:
path: str
label: str = ""
artifact_id: Optional[UUID] = None
scope: str = "working_dir"
version: int = ARTIFACT_REF_VERSION
@ -32,6 +36,8 @@ class ArtifactRef:
"scope": self.scope,
"path": self.path,
}
if self.artifact_id:
out["artifact_id"] = str(self.artifact_id)
if self.label:
out["label"] = self.label
return out
@ -125,7 +131,10 @@ def normalize_artifact_refs(refs: Iterable[ArtifactRef]) -> list[dict]:
out: list[dict] = []
seen: set[tuple[str, str]] = set()
for ref in refs:
if ref.scope != "working_dir" or ref.version != ARTIFACT_REF_VERSION:
if ref.scope != "working_dir" or ref.version not in (
ARTIFACT_REF_VERSION,
ARTIFACT_REF_CURRENT_VERSION,
):
continue
key = (ref.scope, ref.path)
if key in seen:

View File

@ -221,6 +221,7 @@ class AgentLoop:
capabilities: ModelCapabilities,
user_id: UUID,
working_dir: Path,
user_root: Optional[Path] = None,
sink: Optional[Any] = None,
max_iterations: Optional[int] = None,
cancel_check: Optional[Callable[[], bool]] = None,
@ -235,6 +236,7 @@ class AgentLoop:
# ExecCtx 字段:user_id / task_id 已在,working_dir 单独传 —— 供 docker backend
# (Step 3)拼 `--workdir /workspace/<wd_name>` 与临时文件命名空间使用。
self.working_dir = working_dir
self.user_root = Path(user_root).resolve() if user_root else None
self.max_iterations = max_iterations or capabilities.max_iterations
self.sink = sink
# 协作式 cancel:web 层注入 `lambda: broker.is_cancelled(task_id)`;
@ -756,13 +758,31 @@ class AgentLoop:
def _remember_artifacts(self, refs: tuple[dict, ...]) -> None:
"""Accumulate a bounded, ordered set for the final assistant message."""
if refs and self.user_root is not None:
from .artifact_lifecycle import register_published_artifacts
refs = register_published_artifacts(
user_id=self.user_id,
task_id=self.session.task_id,
user_root=self.user_root,
working_dir=self.working_dir,
refs=refs,
)
seen = {
(str(ref.get("scope") or ""), str(ref.get("path") or ""))
(
str(ref.get("artifact_id") or ""),
str(ref.get("scope") or ""),
str(ref.get("path") or ""),
)
for ref in self._pending_artifact_refs
}
for ref in refs:
key = (str(ref.get("scope") or ""), str(ref.get("path") or ""))
if not key[1] or key in seen:
key = (
str(ref.get("artifact_id") or ""),
str(ref.get("scope") or ""),
str(ref.get("path") or ""),
)
if not key[2] or key in seen:
continue
self._pending_artifact_refs.append(dict(ref))
seen.add(key)

View File

@ -196,6 +196,54 @@ class Message(Base):
return sanitize_jsonb_nul(value)
class Artifact(Base):
"""Stable identity and lifecycle metadata for a published workspace file."""
__tablename__ = "artifacts"
__table_args__ = (
Index("ix_artifacts_user_status_path", "user_id", "status", "current_path"),
Index("ix_artifacts_origin_task", "origin_task_id"),
)
artifact_id: Mapped[UUID] = mapped_column(
PG_UUID(as_uuid=True), primary_key=True, default=uuid4
)
user_id: Mapped[UUID] = mapped_column(
PG_UUID(as_uuid=True), ForeignKey("users.user_id"), nullable=False
)
origin_task_id: Mapped[Optional[UUID]] = mapped_column(
PG_UUID(as_uuid=True),
ForeignKey("tasks.task_id", ondelete="SET NULL"),
nullable=True,
)
copied_from_artifact_id: Mapped[Optional[UUID]] = mapped_column(
PG_UUID(as_uuid=True),
ForeignKey("artifacts.artifact_id", ondelete="SET NULL"),
nullable=True,
)
current_path: Mapped[str] = mapped_column(Text, nullable=False)
label: Mapped[str] = mapped_column(Text, nullable=False, default="")
media_type: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
size_bytes: Mapped[Optional[int]] = mapped_column(BigInteger, nullable=True)
content_sha256: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
status: Mapped[str] = mapped_column(
Text, nullable=False, default="active", server_default="active"
)
trash_path: Mapped[Optional[str]] = mapped_column(Text, nullable=True)
created_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True), server_default=func.now(), nullable=False
)
updated_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True),
server_default=func.now(),
onupdate=func.now(),
nullable=False,
)
deleted_at: Mapped[Optional[datetime]] = mapped_column(
DateTime(timezone=True), nullable=True
)
class UsageEvent(Base):
"""per-event 用量记账(0006 v2 形态)。

View File

@ -13,7 +13,7 @@ from sqlalchemy import select, update
from .paths import to_db_path
from .storage import NoSubtaskError, check_no_subtask, session_scope
from .storage.models import Task
from .storage.models import Artifact, Task
class WorkingDirRenameError(RuntimeError):
@ -94,6 +94,22 @@ def rename_working_dir(
.where(Task.task_id.in_(tids))
.values(working_dir=new_db)
)
old_rel = old.relative_to(old.parent).as_posix()
new_rel = new.relative_to(new.parent).as_posix()
artifacts = s.execute(
select(Artifact).where(
Artifact.user_id == user_id,
Artifact.status == "active",
)
).scalars().all()
for artifact in artifacts:
if artifact.current_path == old_rel:
suffix = ""
elif artifact.current_path.startswith(old_rel + "/"):
suffix = artifact.current_path[len(old_rel):]
else:
continue
artifact.current_path = new_rel + suffix
try:
old.rename(new)
except OSError as e:

View File

@ -0,0 +1,143 @@
"""Add stable artifact identity and lifecycle metadata.
Revision ID: 0028
Revises: 0027
Create Date: 2026-08-12
"""
from typing import Sequence, Union
import sqlalchemy as sa
from alembic import op
from sqlalchemy.dialects import postgresql
revision: str = "0028"
down_revision: Union[str, None] = "0027"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.create_table(
"artifacts",
sa.Column(
"artifact_id",
postgresql.UUID(as_uuid=True),
nullable=False,
),
sa.Column(
"user_id",
postgresql.UUID(as_uuid=True),
nullable=False,
),
sa.Column(
"origin_task_id",
postgresql.UUID(as_uuid=True),
nullable=True,
),
sa.Column(
"copied_from_artifact_id",
postgresql.UUID(as_uuid=True),
nullable=True,
),
sa.Column("current_path", sa.Text(), nullable=False),
sa.Column("label", sa.Text(), server_default="", nullable=False),
sa.Column("media_type", sa.Text(), nullable=True),
sa.Column("size_bytes", sa.BigInteger(), nullable=True),
sa.Column("content_sha256", sa.Text(), nullable=True),
sa.Column("status", sa.Text(), server_default="active", nullable=False),
sa.Column("trash_path", sa.Text(), nullable=True),
sa.Column(
"created_at",
sa.DateTime(timezone=True),
server_default=sa.func.now(),
nullable=False,
),
sa.Column(
"updated_at",
sa.DateTime(timezone=True),
server_default=sa.func.now(),
nullable=False,
),
sa.Column("deleted_at", sa.DateTime(timezone=True), nullable=True),
sa.CheckConstraint(
"status IN ('active', 'deleted')",
name="ck_artifacts_status",
),
sa.ForeignKeyConstraint(["user_id"], ["users.user_id"]),
sa.ForeignKeyConstraint(
["origin_task_id"],
["tasks.task_id"],
ondelete="SET NULL",
),
sa.ForeignKeyConstraint(
["copied_from_artifact_id"],
["artifacts.artifact_id"],
ondelete="SET NULL",
),
sa.PrimaryKeyConstraint("artifact_id"),
)
op.create_index(
"ix_artifacts_user_status_path",
"artifacts",
["user_id", "status", "current_path"],
)
op.create_index(
"ix_artifacts_origin_task",
"artifacts",
["origin_task_id"],
)
op.create_index(
"uq_artifacts_active_user_path",
"artifacts",
["user_id", "current_path"],
unique=True,
postgresql_where=sa.text("status = 'active'"),
)
# Backfill legacy structured refs without rewriting append-only message metadata.
# working_dir is stored as workspace/users/<uuid>/<user-relative-dir>.
op.execute(
r"""
INSERT INTO artifacts (
artifact_id, user_id, origin_task_id, current_path, label, status
)
SELECT DISTINCT ON (t.user_id, current_path)
gen_random_uuid(),
t.user_id,
t.task_id,
current_path,
COALESCE(ref->>'label', ''),
'active'
FROM tasks AS t
JOIN messages AS m ON m.task_id = t.task_id
CROSS JOIN LATERAL jsonb_array_elements(
CASE
WHEN jsonb_typeof(m.artifact_refs) = 'array'
THEN m.artifact_refs
ELSE '[]'::jsonb
END
) AS ref
CROSS JOIN LATERAL (
SELECT
regexp_replace(
t.working_dir,
'^workspace/users/' || t.user_id::text || '/',
''
) || '/' || (ref->>'path') AS current_path
) AS resolved
WHERE m.artifact_refs IS NOT NULL
AND ref->>'scope' = 'working_dir'
AND COALESCE(ref->>'path', '') <> ''
AND ref->>'path' !~ '(^/|(^|/)\.\.(/|$))'
AND strpos(ref->>'path', chr(92)) = 0
ORDER BY t.user_id, current_path, m.created_at
ON CONFLICT DO NOTHING
"""
)
def downgrade() -> None:
op.drop_index("uq_artifacts_active_user_path", table_name="artifacts")
op.drop_index("ix_artifacts_origin_task", table_name="artifacts")
op.drop_index("ix_artifacts_user_status_path", table_name="artifacts")
op.drop_table("artifacts")

636
docs/windows-node-design.md Normal file
View File

@ -0,0 +1,636 @@
# zcbot Windows Node 技术方案
> 状态:方案设计,待云环境和 Origin 版本确认后实施
> 首期能力Origin 科研绘图
> 云端Ubuntu 上运行的 zcbot
> Windows 端:.NET 10 LTS / ASP.NET Core Worker Service
> 目标系统Windows 11 Enterprise 优先;经兼容性验证并具备安全更新的 Windows 10 亦可
## 1. 决策摘要
Windows 端采用 **zcbot Windows Node**,由节点主动连接云端 zcbot。云端是唯一的用户入口、智能决策者和任务控制面Windows Node 是受控执行节点,负责本机工业软件、交互式桌面、截图视频和产物回传。
Windows Node 不是完整的本地 zcbot
- 不直接与用户对话;
- 不保存主会话和用户长期记忆;
- 不持有云端模型密钥;
- 不自主扩大任务范围;
- 不接受任意命令、任意脚本或任意桌面动作;
- 只执行云端下发且本地能力清单允许的任务。
核心技术决策:
1. **连接方向**Windows Node 通过出站 WSS/HTTPS 主动连接云端Windows 不开放业务入站端口。
2. **Windows 主体**:使用 .NET 10 LTS、Worker Service、ASP.NET Core 本机诊断 API。
3. **双进程边界**Windows Service 管控制面;`DesktopRunner.exe` 在持续登录的交互式会话中管 GUI、COM 和截图。
4. **软件适配器**Origin/Fluent 使用固定 Python WorkerAspen/CAD 优先使用 C# COM/.NETGUI 仅作补充。
5. **状态分工**:云端 PG 保存任务所有权与期望状态;节点 SQLite/终态文件保存本机执行证据。二者通过租约和重连对账,不构成两个业务事实源。
6. **产物分工**Windows 本地文件是计算现场事实上传到云端后zcbot 用户工作目录中的文件是用户交付事实。
7. **首批 Origin**:主流程使用 Origin 官方推荐的外部 Python `originpro`Computer Use 只做视觉验收和异常处理。
## 2. 目标与非目标
### 2.1 目标
- 云端 zcbot 可以发现节点能力、容量、软件版本和在线状态。
- 用户可以提交、查询、取消长时间 Windows 计算任务。
- 网络中断、云端重启或节点重启后,任务可以确定性对账。
- 节点可以上报阶段、进度、结构化指标、日志摘要和事件截图。
- 中间产物和最终产物支持校验、断点上传与按需导入工作目录。
- 单台节点可以逐步增加 Origin、Fluent、Aspen、CAD 等适配器。
- Computer Use 在独立低权限桌面会话运行,并经过能力白名单和动作策略约束。
- 节点不暴露用户工作区、云端凭据或其他用户任务。
### 2.2 首期非目标
- 不把 Windows Node 做成第二个聊天机器人或子 agent。
- 不开放任意 PowerShell、Python、LabTalk、COM、URL 或 Windows 路径。
- 不实现节点之间迁移已经开始执行的软件任务。
- 不实现通用 HPC 调度器、工作流 DAG 或自动重跑整个计算。
- 不持续把完整桌面视频发送给模型。
- 不依赖 SMB 共享目录作为任务协议。
- 不在首期实现自动更新、WebRTC 和多节点智能负载均衡。
## 3. 总体架构
```mermaid
flowchart LR
User["用户"] --> Z["zcbot Web / 渠道"]
Z --> Agent["Agent Loop"]
Agent --> Tools["Compute Tools"]
Tools --> CP["云端 Node Control Plane<br/>PG + Connection Manager"]
Node["Windows Node Service<br/>.NET Worker"] -->|"出站 WSS + mTLS"| CP
CP -->|"任务 offer / cancel / config"| Node
Node -->|"heartbeat / event / terminal"| CP
Node --> Pipe["ACL Named Pipe"]
Pipe --> Desktop["DesktopRunner.exe<br/>交互式用户会话"]
Desktop --> Origin["Origin / OriginPro"]
Desktop --> Other["Fluent / Aspen / CAD"]
Node --> Local["本机 SQLite + Job Directories"]
Desktop --> Local
Node -->|"HTTPS 分块上传"| Cache["云端隐藏计算缓存"]
Cache --> Workspace["用户 working_dir"]
```
### 3.1 云端组件
| 组件 | 职责 |
|---|---|
| `NodeRegistry` | 节点注册、证书指纹、启停、能力和管理员标签 |
| `NodeConnectionManager` | WSS 连接、心跳、消息 ACK、同节点单活连接 |
| `ComputeJobService` | 用户授权、幂等提交、节点选择、租约、取消和终态 |
| `ComputeTransferService` | 输入下载凭证、分块上传、SHA-256、容量与保留期 |
| `ComputeBroker` | 把任务事件推送到 Web UI不承担持久化事实源 |
| `ComputeTools` | agent 可调用的能力发现、提交、查询、取消、产物导入工具 |
节点机制不复用现有 `generic_openapi` 外部系统定义。外部系统表达“用户凭据下的受控请求”Windows Node 表达“平台托管设备、持久连接、执行租约与本机容量”。两者可以复用认证、审计、大结果和 ActionPolicy 原则,但不共用运行状态模型。
### 3.2 Windows 组件
| 进程 | 运行身份 | 职责 | 禁止事项 |
|---|---|---|---|
| `Zcbot.WindowsNode.Service` | 低权限服务账号 | 云连接、任务台账、租约、调度、传输、日志、进程控制 | 点击桌面、持有用户密码、执行任意脚本 |
| `Zcbot.DesktopRunner.exe` | 专用持续登录账号 | 会话内启动软件、UI Automation、截图、视频、弹窗检测 | 对外监听、决定用户任务、访问其他任务 |
| `OriginAdapter.Worker` | 由 DesktopRunner 启动 | 读取固定 schema、调用 `originpro`、导出和验证 | 网络访问、解释自然语言、加载用户代码 |
| 后续软件 Worker | 按适配器定义 | PyFluent、Aspen COM、CAD .NET/COM | 绕过 Node 能力与路径边界 |
Windows Service 位于 Session 0不能可靠操作用户可见桌面。所有需要用户 profile、COM 桌面对象、可见窗口或截图的动作都由 `DesktopRunner.exe` 执行。Service 与 DesktopRunner 只使用本机命名管道通信,命名管道 ACL 仅允许两个指定账号。
## 4. Windows 技术栈
### 4.1 Node Service
- .NET 10 LTS
- `Microsoft.NET.Sdk.Worker`
- `BackgroundService`
- `ClientWebSocket` / `HttpClient`
- ASP.NET Core Minimal API仅绑定 `127.0.0.1` 提供诊断接口
- `System.Threading.Channels` 管理本机有界任务队列
- SQLite WAL 保存执行台账
- Serilog 写结构化滚动日志和 Windows Event Log
- Windows Certificate Store 保存节点私钥和客户端证书
- Windows Job Objects 或受控 PID 树负责取消和超时回收
- Named Pipes 连接 DesktopRunner
ASP.NET Core 在这里仍是正确选型,但职责已经变化:它不是等云端调用的公开业务网关,而是 Windows Node 的宿主、诊断端点和基础设施框架。跨机业务通信由 Node 主动发起。
### 4.2 本机诊断 API
仅监听 loopback
```text
GET /health/live
GET /health/ready
GET /diagnostics/node
GET /diagnostics/jobs
GET /diagnostics/jobs/{job_id}
POST /diagnostics/reconnect
```
诊断 API 不提供提交任务、任意命令或软件操作。远程运维通过 VPN/堡垒机登录 Windows 后在本机访问,避免形成第二个业务入口。
### 4.3 Python 适配器运行时
Origin 和 Fluent 使用独立、固定版本、离线部署的 Python 运行时。Node 不在任务执行期间安装依赖,不接受请求指定解释器或包版本。
```text
D:\ZcbotNode\runtimes\
├── origin\
│ ├── python.exe
│ └── adapter\
└── fluent\
├── python.exe
└── adapter\
```
Node 生成规范化 `request.json`,固定 Worker 读取请求;事件写 `events.jsonl`,产物写 `artifacts.json`,终态原子写 `terminal.json`。标准输出只用于诊断,不作为任务状态事实源。
## 5. 节点身份与连接
### 5.1 注册流程
```mermaid
sequenceDiagram
participant Admin as 管理员
participant Cloud as zcbot 云端
participant Node as Windows Node
Admin->>Cloud: 创建一次性 enrollment token
Node->>Node: 生成设备密钥对
Node->>Cloud: POST /v1/compute/nodes/enroll
Cloud->>Cloud: 消耗 token创建 node_id
Cloud-->>Node: 客户端证书、CA、云端地址
Node->>Node: 私钥写入 Windows Certificate Store
Node->>Cloud: mTLS 建立 WSS
Cloud-->>Node: 注册成功与当前配置 revision
```
约束:
- enrollment token 一次性、短期有效,并绑定预期节点名称;
- 私钥在 Windows 本机生成且不可导出;
- 云端只保存证书指纹和公钥身份;
- 节点证书支持轮换、吊销和管理员禁用;
- 重装节点必须重新注册,不复用复制出来的私钥。
### 5.2 长连接
节点连接:
```text
WSS /v1/compute/nodes/connect
```
统一消息 envelope
```json
{
"protocol_version": 1,
"message_id": "019...",
"type": "heartbeat",
"sent_at": "2026-08-12T08:00:00Z",
"payload": {}
}
```
| 方向 | 类型 | 语义 |
|---|---|---|
| Node → Cloud | `hello` | Node 版本、OS、boot ID、本机配置 revision |
| Node → Cloud | `heartbeat` | 容量、软件健康、运行任务摘要、磁盘 |
| Node → Cloud | `job_accept` / `job_reject` | 是否已原子接收 offer |
| Node → Cloud | `job_state` | 阶段、进度和结构化指标 |
| Node → Cloud | `job_event` | 重要事件、警告和截图引用 |
| Node → Cloud | `job_terminal` | 成功、失败、取消的唯一远端终态消息 |
| Node → Cloud | `lease_renew` | 续租仍在本机执行的任务 |
| Cloud → Node | `job_offer` | 带租约的任务候选,不等于已分派 |
| Cloud → Node | `job_cancel` | 请求协作式取消 |
| Cloud → Node | `config_changed` | 提示节点拉取新配置 |
| 双向 | `ack` / `ping` / `pong` | 去重、保活和连接质量 |
消息按 `message_id` 幂等。重要消息在发送方保存到收到 ACK心跳和高频进度允许覆盖不进入无限重放队列。
### 5.3 单活连接
同一 `node_id` 只允许一个活动连接。新连接带 `boot_id` 和递增 `connection_epoch`;云端原子替换旧连接。旧连接收到 fencing 后不得继续接受新任务,防止快照克隆或网络分区导致双执行。
## 6. 能力与容量模型
Node 在 `hello` 和心跳中声明由本机可信配置生成的能力:
```json
{
"node_id": "019...",
"node_version": "0.1.0",
"os": "windows-11-enterprise-24h2",
"capabilities": [
{
"name": "origin.plot",
"adapter_version": "0.1.0",
"software": "OriginPro",
"software_version": "2026",
"schema_version": 1,
"max_concurrency": 1,
"available_slots": 1,
"health": "ready"
},
{
"name": "computer.windows.capture",
"schema_version": 1,
"max_concurrency": 1,
"available_slots": 1,
"health": "ready"
}
],
"disk_free_bytes": 536870912000
}
```
能力清单来自管理员部署的适配器 manifest不能由模型或用户请求自行添加。云端 definition 约束允许哪些用户、角色或任务使用某项能力;节点再次校验 capability、schema version、文件配额和本机 slot形成双层防护。
## 7. 云端数据模型
建议新增三张表,不复用外部系统连接表:
```text
compute_nodes(
node_id pk, name, cert_fingerprint, status,
labels jsonb, capabilities jsonb, config_revision,
last_seen_at, disabled_at, created_at, updated_at
)
compute_jobs(
job_id pk, user_id fk, task_id fk, tool_call_id,
capability, schema_version, request jsonb,
idempotency_key, request_digest,
node_id fk null, lease_id, lease_expires_at,
status, stage, progress, metrics jsonb,
error jsonb, artifact_manifest jsonb,
created_at, started_at, terminal_at, updated_at
)
compute_job_events(
event_id pk, job_id fk, sequence,
kind, level, payload jsonb, created_at
)
```
约束:
- `(user_id, idempotency_key)` 唯一;
- request 只保存规范化业务参数和文件引用,不保存二进制内容或密钥;
- events 只保存阶段变化、警告、错误、用户可见事件和产物事件,不保存视频帧和全部原始日志;
- 原始日志、截图和上传中产物进入隐藏文件缓存;
- 事件按配置保留,终态任务和用户交付文件遵守 zcbot 通用保留策略。
## 8. 任务状态、租约与恢复
### 8.1 云端状态机
```mermaid
stateDiagram-v2
[*] --> queued
queued --> offered
offered --> dispatched: node accept
offered --> queued: reject / offer timeout
dispatched --> running
running --> succeeded
running --> failed
running --> cancelling
cancelling --> cancelled
queued --> cancelled
offered --> cancelled
dispatched --> unknown: lease expired
running --> unknown: lease expired
cancelling --> unknown: lease expired
unknown --> running: node reconnect + proof
unknown --> succeeded: terminal reconciliation
unknown --> failed: operator reconciliation
```
`unknown` 表示云端不知道现场状态,不等于失败,也不得立即把任务分派给另一台节点。
### 8.2 分派协议
1. 云端按 capability、schema、管理员授权、健康、空闲 slot 和标签选择候选节点。
2. 云端创建短期 offer lease发送 `job_offer`
3. Node 完成 schema、容量、磁盘、许可证和输入可用性预检。
4. Node 先把 job 和 lease 原子写入本地 SQLite再发送 `job_accept`
5. 云端以相同 lease ID 把状态改为 `dispatched`
6. Node 下载输入并启动 Worker持续续租。
节点收到重复 offer 时,根据本地 `job_id + request_digest` 返回原 accept/reject不重复执行。
### 8.3 断线与租约
- 运行任务每 15 秒续租,默认租约 60 秒;参数可配置。
- 连接断开后,节点继续运行已开始任务,并把重要事件保存在本地 outbox。
- 租约过期后云端标记 `unknown`,进入恢复宽限期,不自动重派。
- 节点重连时上报所有非终态任务的 job、request digest、lease、本机阶段、PID、心跳和终态文件摘要。
- 云端对账后续租、接收终态,或要求节点停止孤儿任务。
- 只有明确证明 Worker 尚未开始,任务才能安全回到 `queued`
- 已进入外部软件执行阶段的任务不自动跨节点重试。
### 8.4 节点重启
本地 SQLite 记录任务台账;任务目录的 `terminal.json` 是本机唯一终态文件。Node 启动时:
1. 扫描本地非终态台账;
2. 检查 Worker PID、boot ID、进程创建时间和任务目录
3. 读取 `terminal.json`
4. 能证明仍在运行则恢复监控;
5. 有终态文件则加入 outbox 等待回传;
6. 无进程无终态则标记 `NODE_RESTARTED_DURING_JOB`,不伪造成功。
## 9. 文件与产物传输
### 9.1 输入
Node 不访问整个用户 workspace。云端只为显式引用的文件创建短期、单任务、只读下载凭证。Node 下载后校验 SHA-256、声明大小、扩展名和内容类型。凭证不能列目录、不能换路径、不能用于其他 job。
### 9.2 上传
大产物不通过 WSS 消息传输。Node 使用 HTTPS 分块上传:
```text
POST /v1/compute/jobs/{job_id}/artifacts/upload-session
PUT /v1/compute/transfers/{transfer_id}/parts/{part_number}
POST /v1/compute/transfers/{transfer_id}/complete
```
上传校验 artifact ID、job 与租约、文件大小、分块摘要、最终 SHA-256、配额、允许类型及压缩包安全。
### 9.3 云端生命周期
Node 上传完成后先进入:
```text
<user_root>/.zcbot_cache/<task_id>/compute_jobs/<job_id>/
```
该目录默认隐藏且有 TTL。用户或 agent 明确导入后复制到:
```text
<working_dir>/materials/simulation/<job_id>/
```
最终图、PDF、视频和可编辑工程可通过 `publish_artifacts` 提升为对话产物。真实文件仍是事实源。
## 10. zcbot 工具面
```text
compute_capability_list
compute_job_submit
compute_job_status
compute_job_cancel
compute_job_artifact_import
```
- capability list 只返回用户有权使用且有健康节点承载的能力;
- submit 只接受 capability schema不接受 node ID、命令和路径
- 节点选择由平台完成;
- status 返回阶段、进度、指标、重要事件和产物 manifest
- cancel 按 ActionPolicy 分类并审计;
- artifact import 只接受 manifest 中的 artifact ID。
长任务不得让 agent 持续空等。完成时更新 Web 任务卡并通知用户,但不自动启动新 LLM run用户点击“分析结果”或回复继续后再查询并导入产物。
## 11. Computer Use 与屏幕采集
优先级:软件官方 API/SDK → Windows UI Automation → 截图视觉坐标 → 暂停人工处理。Computer Use 不是工业计算主协议,结构化进度和结果始终优先。
首期 DesktopRunner 内部能力:
```text
session.health
window.list
window.capture
window.wait
popup.detect
process.launch_adapter
process.cancel_adapter
```
后续经安全评审后才开放点击、输入、快捷键和视频控制。优先使用 `automation_id + control name + window identity`,坐标点击只作最后手段。
截图与视频规则:
- 默认只采集关键事件截图;
- 优先捕获目标窗口;
- 敏感信息上传前遮罩;
- 视频默认关闭,显式启用时建议 1280×720、510 FPS、H.264 分段;
- 采集失败不得使计算任务失败;
- 模型只按需读取关键帧,不持续消费完整视频。
## 12. Origin 首批适配器
### 12.1 接口与执行位置
首期使用 Origin 官方推荐的外部 Python `originpro`。它能读写数据、创建和修改图形、导出图形,并启动可见或隐藏的 Origin需要 Windows 本机安装并授权 Origin 2021 或更高版本。
- [Origin External Python](https://docs.originlab.com/externalpython/)
- [Origin Automation Server](https://docs.originlab.com/com/)
- [Origin 图形导出](https://docs.originlab.com/user-guide/publishing-and-export/)
C# COM 只作未覆盖能力或旧版本备用,并与实际 Origin 版本匹配。Origin Worker 由 DesktopRunner 在专用登录账号会话内启动,以统一用户 profile、许可证上下文、COM 会话和截图来源。默认隐藏执行;视觉验收时使用可见模式。
### 12.2 首期 capability
只开放 `origin.plot@v1`,接受 CSV、XLSX 或规范化 JSON不接受任意 Python、LabTalk、模板文件或宏。
```json
{
"schema_version": 1,
"input": {"input_id": "data-01", "sheet": "Sheet1"},
"plot": {
"type": "line_scatter",
"x": "Temperature",
"y": ["Strength_7d", "Strength_28d"],
"template": "publication_double_column",
"title": "温度对抗压强度的影响",
"x_axis": {"title": "温度", "unit": "°C", "scale": "linear"},
"y_axis": {"title": "抗压强度", "unit": "MPa", "scale": "linear"},
"legend": {"enabled": true, "position": "top_right"},
"error_bars": null
},
"output": {
"formats": ["opju", "png", "svg", "pdf"],
"dpi": 600,
"capture_screenshots": true,
"record_video": false
}
}
```
首期图形类型line、scatter、line_scatter、grouped_bar、box、histogram、heatmap。后续增加误差棒组合、三元图、等高线、三维曲面、XRD 堆叠图、热分析联图和多面板布局。
### 12.3 模板与产物
字体、线宽、配色、图幅和导出参数由管理员审核的模板定义。请求只引用模板 IDprovenance 记录模板、Origin、适配器版本和输入 SHA-256。
```text
D:\ZcbotNode\jobs\<job_id>\
├── request\request.json
├── input\
├── work\
├── logs\events.jsonl
├── screenshots\
├── recordings\
├── output\
│ ├── project.opju
│ ├── figure.png
│ ├── figure.svg
│ ├── figure.pdf
│ ├── plot-spec.json
│ └── provenance.json
├── artifacts.json
└── terminal.json
```
`terminal.json` 使用临时文件、flush 和原子 replace。成功前必须完成产物校验和 manifest 写入。
自动验收包括输出存在且非空、PNG 尺寸/DPI、SVG/PDF 可解析、系列数量、轴标题/单位/图例、OPJU 保存、provenance 完整、Origin 正常退出且无许可证泄漏。
## 13. 并发、资源与许可证
- 每项 capability 独立声明 `max_concurrency`
- Origin 默认并发 1
- Origin 实例不跨 job 共享;
- 启动前检查内存、磁盘、桌面会话和许可证;
- Node 使用有界队列,云端只向有空闲 slot 的节点发 offer
- 超时区分启动、无进展、总运行和取消宽限期;
- 取消先协作式退出,超时后回收精确进程树;
- PID、创建时间、可执行文件摘要与 job 绑定,避免误杀。
## 14. 安全与审计
- Node 只需出站 HTTPS/WSSWindows 不开放业务端口;
- RDP 通过 VPN/堡垒机,不暴露公网;
- mTLS 验证节点,禁止跳过服务端 TLS
- Service、DesktopRunner、运维账号分离且不授予管理员权限
- 节点私钥存 Windows Certificate Store 并限制 ACL
- capability schema 是唯一业务输入入口;
- 不接受命令、脚本、URL、绝对路径、UNC 和环境变量覆盖;
- 适配器与模板只由管理员部署;
- Worker 默认无外网,许可证服务器单独放行;
- 未知 capability 和未知 UI 动作 fail closed。
审计记录提交人、task、capability、node、软件/模板/适配器版本、分派、续租、取消、重连、终态、Computer Use 动作、输入输出摘要和导入位置。不得记录密码、Token、许可证密钥或未脱敏截图。
## 15. 可观测性与错误码
Node 心跳上报版本、OS、boot ID、connection epoch、CPU、内存、磁盘、上传队列、DesktopRunner 状态、软件版本、许可证健康、slot、运行 job 和 outbox 积压。
稳定错误码:
```text
NODE_OFFLINE
NODE_CAPABILITY_UNAVAILABLE
NODE_DISK_INSUFFICIENT
DESKTOP_SESSION_UNAVAILABLE
ADAPTER_VERSION_MISMATCH
INPUT_HASH_MISMATCH
ORIGIN_NOT_INSTALLED
ORIGIN_LICENSE_UNAVAILABLE
ORIGIN_START_TIMEOUT
ORIGIN_AUTOMATION_FAILED
OUTPUT_VALIDATION_FAILED
ARTIFACT_UPLOAD_FAILED
JOB_CANCELLED
NODE_RESTARTED_DURING_JOB
```
用户侧只返回可操作说明、是否可重试和 retry-after完整诊断留在节点和管理员日志。
## 16. 部署建议
```text
D:\ZcbotNode\
├── app\
├── config\node.json
├── state\node.db
├── logs\
├── jobs\
├── runtimes\
├── adapters\
├── templates\
└── updates\
```
首期云主机建议Windows 11 Enterprise 或经验证且有安全更新的 Windows 10 Enterprise、8 vCPU、32 GB RAM、200 GB 系统盘、500 GB 独立数据盘、CPU 实例。Origin 二维科研绘图不要求独立 GPU。
Windows 基线:固定版本并受控更新;禁止休眠和自动锁屏;固定分辨率与缩放;数据盘固定盘符;完成 Origin 人工基线任务创建三个独立账号DesktopRunner 登录后自动启动;仅允许必要出站和受控运维;稳定后制作镜像和恢复手册。
## 17. 实施阶段
### Phase 0环境探针
- 确认 Windows、Origin/OriginPro 精确版本和许可证;
- 手工导入并导出 OPJU、PNG、SVG、PDF
- 验证外部 `originpro` 可见/隐藏实例;
- 验证账号、RDP 断开、锁屏策略和截图;
- 连续运行 24 小时,确认许可证和残留进程;
- 测量典型任务时间、文件和磁盘占用。
### Phase 1最小闭环
- enrollment、节点表、WSS Connection Manager
- Node 注册、证书、心跳、能力上报;
- job 表、幂等提交、offer/accept 和基本租约;
- Origin line/scatter/line_scatter
- 输入下载、分块上传、SHA-256
- zcbot 查询、取消、导入和发布;
- 最终截图与稳定错误码。
### Phase 2可靠性与绘图可用性
- outbox、重连对账、unknown 状态和节点重启恢复;
- bar、box、histogram、heatmap、误差棒、多系列
- 模板版本、自动机检、视觉验收;
- Web 任务卡、完成通知和渠道通知;
- 事件截图和可选录像;
- Origin 内置 skill 和回归样例集。
### Phase 3平台化
- 经审核的 Computer Use 动作;
- 节点禁用、证书轮换、配置 revision、自动升级
- 多节点选择、标签和许可证容量;
- Fluent、Aspen、CAD 适配器;
- 根据真实负载决定 GPU 节点或统一调度器。
## 18. 首期验收标准
1. Windows 不开放业务入站端口Node 能通过 mTLS WSS 注册、重连和心跳。
2. 云端能发现 `origin.plot@v1`、软件版本、健康和 slot。
3. 同一幂等键不会创建两个任务。
4. 双连接时只有最新 epoch 能接受任务。
5. zcbot 能用 CSV 生成指定模板图并获得 `job_id`
6. OPJU、PNG、SVG、PDF 均可打开并通过自动检查。
7. zcbot 能校验、导入和发布结果。
8. 连续执行 50 个任务,无残留进程、跨任务污染或许可证泄漏。
9. WSS 断开时任务继续,重连后能对账。
10. 云端或 Node 重启不会重复执行或误报终态。
11. 取消能结束精确进程树,不误杀其他实例。
12. 截图/录像失败不影响主任务。
13. 非法路径、脚本、越权 job、超限文件和伪造 artifact ID 被拒绝。
14. 提交、分派、取消、Computer Use 和产物导入均有审计。
## 19. 已否决方案
- **云端直调 Windows OpenAPI**:首版简单,但入站网络、设备身份、租约和多节点仍需重做;只保留 loopback 诊断 API。
- **Windows 安装完整 zcbot**形成两套会话、模型、权限和状态Node 只做执行。
- **全程 Computer Use**:对窗口和会话敏感且不适合数值结果;官方 API/SDK 为主。
- **Go 编写 Node**Windows Service、COM/.NET、UI Automation、Named Pipes 更适合 .NETPython 仅作固定适配器。
- **SMB + WinRM**:任意命令、共享权限、恢复和审计边界较差。
## 20. 实施前待确认
- 云厂商、Windows 镜像、网络出口和授权;
- Origin/OriginPro 精确版本;
- 节点锁或浮动许可证及并发数;
- 典型 Origin 输入和期望结果;
- 首批图形类型与现有绘图模板;
- 输入、产物、任务容量和保留期;
- 云端域名、内部 CA/mTLS 和 enrollment 管理;
- 首期完成通知范围;
- Windows 维护责任人与操作窗口。

View File

@ -0,0 +1,262 @@
# Windows Node MVP 实施方案(内网版)
> **当前有效的第一阶段开发依据。**本文件取代 `windows-node-mvp.md` 中的 MVP 通信与注册方案。
> 长期演进边界见 `windows-node-design.md`
> 首批能力:`origin.plot@v1`。
## 1. 适用边界
本方案成立的前提:
- zcbot 与 Windows Node 位于同一受控内网、云 VPC 或专用 VPN
- 两端使用固定私网地址或内网 DNS
- zcbot 的 Node 接口只监听私网地址;
- 安全组仅允许指定 Node IP 访问指定 zcbot 端口;
- 通信不经过公网、访客网或不可信办公终端所在网络。
MVP 使用:
```text
注册、状态和文件HTTP
任务控制长连接WS
节点认证:每个 Node 独立的长期 Bearer Token
```
HTTP/WS 不加密传输Node Token、输入数据、截图和产物在链路上均为明文。如果上述网络边界发生变化必须先升级 HTTPS/WSS。
## 2. MVP 架构
```mermaid
flowchart LR
Z["内网 zcbot"] <-->|"HTTP + WS<br/>独立 Node Token"| N["Windows Node.exe<br/>登录会话自动启动"]
N --> P["Origin Python Worker"]
P --> O["Origin / OriginPro"]
P --> F["独立任务目录"]
N -->|"状态、截图、产物"| Z
```
第一版只有两个进程:
- `Zcbot.WindowsNode.exe`.NET 10负责注册、自动连接、任务状态、进程控制、截图和上传
- `OriginAdapter.Worker`:固定 Python 环境,使用 `originpro` 生成 OPJU、PNG、SVG 和 PDF。
Node 首期运行在持续登录的专用 Windows 用户会话,通过计划任务在登录后自动启动。暂不拆 Windows Service 与 DesktopRunner。
## 3. Node 注册
### 3.1 操作流程
1. 管理员在 zcbot 管理端创建一次性注册码,默认 10 分钟有效。
2. Windows Node 首次启动时输入 zcbot 内网地址、节点名称和注册码。
3. Node 调用内网注册接口。
4. zcbot 原子消费注册码并返回 `node_id + node_token`
5. Node 使用 Windows DPAPI 加密保存 Token。
6. 此后 Node 或 Windows 重启时自动建立 WS 连接,不再要求人工输入。
```mermaid
sequenceDiagram
participant A as 管理员
participant Z as zcbot
participant N as Windows Node
A->>Z: 创建一次性注册码
A->>N: 输入内网地址和注册码
N->>Z: HTTP POST /v1/compute/nodes/enroll
Z->>Z: 校验并原子消费注册码
Z-->>N: node_id + node_token + 配置
N->>N: DPAPI 加密保存 node_token
N->>Z: 携带 Bearer Token 建立 WS
Z-->>N: active
```
注册请求:
```http
POST http://zcbot.internal:8765/v1/compute/nodes/enroll
Content-Type: application/json
```
```json
{
"enrollment_code": "ZCN-7H4K-9P2M",
"node_name": "win-origin-01",
"install_id": "019...",
"node_version": "0.1.0",
"os_version": "Windows 11 Enterprise 24H2",
"capabilities": ["origin.plot@v1"]
}
```
注册响应:
```json
{
"node_id": "019...",
"node_token": "仅返回一次的高熵随机 Token",
"heartbeat_seconds": 15,
"max_concurrency": 1
}
```
### 3.2 注册与 Token 底线
- 注册码一次性、短期有效,云端只保存哈希;
- 注册码可以绑定预期节点名称和允许能力;
- 注册成功或达到失败尝试上限后立即失效;
- 每个 Node 使用不同 Token禁止共享全局永久 Token
- Token 至少包含 32 字节密码学安全随机数;
- zcbot 数据库只保存 Token 强哈希,明文仅在注册响应返回一次;
- Node 使用 Windows DPAPI `LocalMachine` 加密 Token并用文件 ACL 限制为 Node 运行账号可读;
- Token 只放 `Authorization` Header不进入 URL、查询参数或日志
- Node 不上传 Windows 密码、许可证密钥或完整硬件指纹;
- 管理员可以禁用 Node 或轮换 Token禁用后立即拒绝连接、任务和上传
- Windows 重装、本地身份丢失或克隆云盘后必须重新注册。
## 4. 自动连接与重连
### 4.1 WS 连接
```http
GET ws://zcbot.internal:8765/v1/compute/nodes/connect
Authorization: Bearer <node_token>
X-Node-Id: <node_id>
Upgrade: websocket
```
HTTP 注册、状态、文件上传下载和 WS 长连接使用相同的 Node Token不增加短期 connection token、HMAC、nonce、时间戳签名或证书体系。
连接后 Node 上报:
- Node、OS 和 Origin 版本;
- capability 和可用 slot
- 本机磁盘和桌面会话状态;
- 本地运行中或尚未确认终态的 job 摘要。
同一 `node_id` 只保留一个活动连接。新连接成功后关闭旧连接。MVP 单节点场景下,发现相同身份来自不同 `install_id` 时拒绝新连接并提示重新注册,避免云盘克隆产生双执行者。
### 4.2 重连策略
- 断线后按 1、2、5、10、30、60 秒并加随机抖动重连;
- 最大间隔 60 秒;
- 正在运行的 Origin 任务不因 WS 断开而终止;
- 重要终态与上传状态写入本地 SQLite重连后补报
- `401/403` 表示 Token 失效,停止高频重试并显示“需要重新注册”;
- 连接超时和 `5xx` 继续退避重试;
- Windows 网络恢复时立即触发连接尝试。
## 5. 内网安全配置
最低网络规则:
```text
zcbot Node API 绑定zcbot 私网 IP:8765示例
zcbot 入站安全组:只允许 Windows Node 私网 IP → TCP 8765
Windows Node 出站:只允许 zcbot 私网 IP → TCP 8765
Windows Node 入站:不开放 Node 业务端口
RDP不向公网开放使用 VPN、堡垒机或云安全登录
```
即使在内网,应用层仍拒绝:
- 任意 PowerShell、Python、LabTalk 和命令行;
- 任意 URL、绝对路径和 UNC 路径;
- 未声明 capability
- 越权 job 和伪造 artifact ID
- 超出大小、类型和配额限制的文件。
建议在日志中记录 Node、job、用户、动作和结果但不得记录注册码或 Token。
## 6. MVP 状态与数据
云端首期只增加:
```text
compute_nodes(
node_id pk, name, install_id, token_hash, status,
capabilities jsonb, last_seen_at,
created_at, updated_at
)
compute_jobs(
job_id pk, user_id fk, task_id fk,
idempotency_key, capability, request jsonb,
node_id fk, status, progress,
error jsonb, artifact_manifest jsonb,
created_at, terminal_at, updated_at
)
```
MVP 状态:
```text
queued
running
succeeded
failed
cancelled
disconnected
```
Node 断线且本地任务可能仍在执行时标记 `disconnected`不得自动重派。Node 重连后按 `job_id + request_digest` 和本地终态对账。
## 7. Origin 任务闭环
```text
用户上传 CSV/XLSX
→ zcbot 生成受控 plot spec
→ Node 接收 origin.plot@v1
→ 先持久化 job再启动 Origin Worker
→ originpro 生成 OPJU/PNG/SVG/PDF
→ terminal.json 原子记录终态
→ Node 通过 HTTP 上传最终截图和产物
→ zcbot 导入 working_dir 并发布
```
MVP 保留:
- 一次性注册与自动连接;
- 每 Node 独立 Token、DPAPI、禁用和轮换
- 幂等提交;
- 独立任务目录;
- 断线不终止计算;
- SQLite 和 `terminal.json`
- SHA-256 校验;
- 协作取消与精确进程回收;
- 最终截图和产物上传。
MVP 暂缓:
- HTTPS/WSS
- 短期连接令牌、HMAC、nonce
- 客户端证书和 mTLS
- 完整任务租约和多节点调度;
- 独立 DesktopRunner
- 分块断点上传;
- 自动更新;
- 任意 Computer Use
- 视频直播。
## 8. 验收标准
1. 新 Node 能凭一次性注册码完成注册。
2. Node 或 Windows 重启后无需人工操作即可自动连接。
3. zcbot 能显示在线状态、最后心跳、Origin 版本和 slot。
4. Token 不出现在 URL、日志或 zcbot 数据库明文中。
5. 禁用 Node 后,现有连接关闭且无法重新连接或上传。
6. 复制 Node 配置到不同 `install_id` 的机器不能形成两个活动执行者。
7. 同一幂等键不会创建两个 Origin 任务。
8. WS 中断时 Origin 继续运行,重连后能补报状态和产物。
9. Node 重启后能识别已有终态,不重复绘图。
10. zcbot 能接收并发布 OPJU、PNG、SVG 和 PDF。
11. 连续运行 50 个任务,无残留 Origin 进程或许可证泄漏。
12. 从非白名单内网 IP 访问 Node API 被安全组拒绝。
## 9. 升级触发条件
出现以下任一情况,先将通信升级到 HTTPS/WSS
- Node 与 zcbot 跨 VPC、跨安全域或经过公网
- 同一网络出现不受信任终端;
- 输入、截图或产物属于敏感数据并要求链路加密;
- 安全审计明确要求传输加密。
升级时保持 URL path、Node ID、Bearer Header 和任务协议不变,只把 `http/ws` scheme 改为 `https/wss` 并部署服务端证书。设备身份治理进一步提高时,再升级客户端证书和 mTLS。

View File

@ -106,7 +106,8 @@ test("assistant HTML artifacts render inline with lazy loading and an expand act
assert.match(chatJs, /Array\.isArray\(m\.artifact_refs\)/);
assert.match(chatJs, /renderArtifactBarHtml\(m\.artifact_refs, true, state\.taskId/);
assert.match(previewJs, /\/v1\/tasks\/\$\{encodeURIComponent\(taskId\)\}\/files\/download/);
assert.match(previewJs, /downloadFile\(_fpCurrentRel, _fpCurrentTaskId, _fpCurrentLegacy\)/);
assert.match(previewJs, /downloadFile\(_fpCurrentRel, _fpCurrentTaskId, _fpCurrentLegacy, _fpCurrentArtifactId\)/);
assert.match(mediaJs, /data-artifact-id/);
assert.match(chatJs, /dataset\.legacyPath === "1"/);
const clickHandler = chatJs.indexOf('$("chat-stream").addEventListener("click"');
const expandHandler = chatJs.indexOf('e.target.closest(".art-html-open[data-rel]")');

View File

@ -1,10 +1,20 @@
import tempfile
import unittest
from contextlib import contextmanager
from pathlib import Path
from types import SimpleNamespace
from unittest.mock import MagicMock, patch
from uuid import uuid4
from core.artifacts import ArtifactPathError, ToolExecutionResult, resolve_artifact_path
from core.artifacts import (
ArtifactRef,
ArtifactPathError,
ToolExecutionResult,
resolve_artifact_path,
)
from core.executor import ExecCtx
from core.executor_host import HostExecutor
from core.artifact_lifecycle import trash_active_artifacts
from tools.publish_artifacts import PublishArtifactsTool
from web.routers.files import _task_file_target
@ -35,6 +45,47 @@ class ArtifactPathTests(unittest.TestCase):
self.assertEqual(actual, expected.resolve())
self.assertEqual(rel, "report.pdf")
def test_version_two_ref_carries_stable_artifact_identity(self) -> None:
artifact_id = uuid4()
ref = ArtifactRef(
path="report.pdf",
label="最终报告",
artifact_id=artifact_id,
version=2,
).as_dict()
self.assertEqual(ref["version"], 2)
self.assertEqual(ref["artifact_id"], str(artifact_id))
def test_trash_moves_file_and_marks_lifecycle_row_deleted(self) -> None:
artifact = self.wd / "report.pdf"
row = SimpleNamespace(
current_path="技术讨论/report.pdf",
status="active",
deleted_at=None,
trash_path=None,
)
session = MagicMock()
session.execute.return_value.scalars.return_value.all.return_value = [row]
@contextmanager
def fake_scope():
yield session
with patch("core.artifact_lifecycle.session_scope", fake_scope):
count = trash_active_artifacts(
user_id=uuid4(),
user_root=self.root,
target=artifact,
)
self.assertEqual(count, 1)
self.assertFalse(artifact.exists())
self.assertEqual(row.status, "deleted")
self.assertIsNotNone(row.deleted_at)
trashed = self.root / row.trash_path
self.assertTrue(trashed.is_file())
self.assertEqual(trashed.read_bytes(), b"pdf")
def test_explicit_dot_slash_disambiguates_same_named_subdirectory(self) -> None:
nested = self.wd / "技术讨论" / "nested.html"
nested.parent.mkdir()

View File

@ -48,7 +48,7 @@ try:
if not _test_db_ready():
raise RuntimeError("ZCBOT_TEST_DB_URL 未设")
from core.storage import session_scope
from core.storage.models import Message, Task, UsageEvent, User
from core.storage.models import Artifact, Message, Task, UsageEvent, User
with session_scope() as _s:
_s.execute(__import__("sqlalchemy").select(1))
@ -57,6 +57,7 @@ except Exception:
_DB_OK = False
if _DB_OK:
from sqlalchemy import select
from starlette.testclient import TestClient
from web.app import create_app
from web.auth import AuthConfig, mint_token
@ -90,6 +91,7 @@ def tearDownModule() -> None:
tids = s.execute(select(Task.task_id).where(Task.user_id == _UID)).scalars().all()
if tids:
s.execute(delete(Message).where(Message.task_id.in_(tids)))
s.execute(delete(Artifact).where(Artifact.user_id == _UID))
s.execute(delete(UsageEvent).where(UsageEvent.user_id == _UID))
s.execute(delete(Task).where(Task.user_id == _UID))
s.execute(delete(User).where(User.user_id == _UID))
@ -338,6 +340,80 @@ class FilesDbAwareTests(unittest.TestCase):
self.assertEqual(r.status_code, 201, r.text)
return r.json()["task_id"]
def test_published_artifact_delete_moves_to_hidden_trash(self):
tid = self._mk_task("产物回收任务", "产物回收目录")
artifact = _user_root() / "产物回收目录" / "reports" / "result.pdf"
artifact.parent.mkdir(parents=True, exist_ok=True)
artifact.write_bytes(b"artifact")
with session_scope() as s:
s.add(Artifact(
user_id=_UID,
origin_task_id=uuid.UUID(tid),
current_path="产物回收目录/reports/result.pdf",
label="结果",
))
r = _client.post(
"/v1/files/delete",
json={"path": "产物回收目录/reports/result.pdf"},
headers=_AUTH,
)
self.assertEqual(r.status_code, 200, r.text)
self.assertEqual(r.json()["artifacts_trashed"], 1)
self.assertFalse(artifact.exists())
trash = _user_root() / ".zcbot_artifact_trash"
trashed = list(trash.rglob("result.pdf"))
self.assertEqual(len(trashed), 1)
self.assertEqual(trashed[0].read_bytes(), b"artifact")
def test_artifact_copy_gets_new_identity_and_move_keeps_it(self):
tid = self._mk_task("产物复制任务", "产物复制目录")
source = _user_root() / "产物复制目录" / "result.pdf"
source.write_bytes(b"artifact")
destination = _user_root() / "复制目标"
destination.mkdir()
with session_scope() as s:
original = Artifact(
user_id=_UID,
origin_task_id=uuid.UUID(tid),
current_path="产物复制目录/result.pdf",
label="结果",
)
s.add(original)
s.flush()
original_id = original.artifact_id
copied = _client.post(
"/v1/files/copy",
json={"paths": ["产物复制目录/result.pdf"], "dest_dir": "复制目标"},
headers=_AUTH,
)
self.assertEqual(copied.status_code, 200, copied.text)
self.assertEqual(copied.json()["transferred"][0]["artifacts_copied"], 1)
with session_scope() as s:
copy_row = s.execute(
select(Artifact).where(
Artifact.user_id == _UID,
Artifact.current_path == "复制目标/result.pdf",
)
).scalar_one()
copied_id = copy_row.artifact_id
self.assertNotEqual(copied_id, original_id)
self.assertEqual(copy_row.copied_from_artifact_id, original_id)
archive = _user_root() / "归档"
archive.mkdir()
moved = _client.post(
"/v1/files/move",
json={"paths": ["复制目标/result.pdf"], "dest_dir": "归档"},
headers=_AUTH,
)
self.assertEqual(moved.status_code, 200, moved.text)
with session_scope() as s:
moved_row = s.get(Artifact, copied_id)
self.assertEqual(moved_row.current_path, "归档/result.pdf")
def test_toplevel_rename_cascades_db(self):
tid = self._mk_task("改名任务", "改名前目录")
r = _client.post("/v1/files/rename",

View File

@ -384,16 +384,18 @@ class FilesRoutesTests(unittest.TestCase):
def test_rename_delete_copy_nontop(self):
(self.wd / "sub" / "c.txt").write_text("c", encoding="utf-8")
# 非顶层改名:纯 FS,tasks_updated=0
r = _client.post("/v1/files/rename",
json={"path": "route-test-wd/sub/c.txt", "new_name": "c2.txt"}, headers=_AUTH)
with patch("core.artifact_lifecycle.rename_active_artifacts", return_value=0):
r = _client.post("/v1/files/rename",
json={"path": "route-test-wd/sub/c.txt", "new_name": "c2.txt"}, headers=_AUTH)
self.assertEqual(r.status_code, 200)
self.assertEqual(r.json()["tasks_updated"], 0)
self.assertTrue((self.wd / "sub" / "c2.txt").is_file())
# 拷贝到子目录内(非顶层,无 DB 闸)
(self.wd / "dest").mkdir(exist_ok=True)
r = _client.post("/v1/files/copy",
json={"paths": ["route-test-wd/sub/c2.txt"], "dest_dir": "route-test-wd/dest"},
headers=_AUTH)
with patch("core.artifact_lifecycle.copy_active_artifacts", return_value=0):
r = _client.post("/v1/files/copy",
json={"paths": ["route-test-wd/sub/c2.txt"], "dest_dir": "route-test-wd/dest"},
headers=_AUTH)
self.assertEqual(r.status_code, 200)
self.assertTrue((self.wd / "dest" / "c2.txt").is_file())
# 目标已存在 → 409(预检整批 abort)
@ -402,11 +404,13 @@ class FilesRoutesTests(unittest.TestCase):
headers=_AUTH)
self.assertEqual(r.status_code, 409)
# 删文件
r = _client.post("/v1/files/delete", json={"path": "route-test-wd/dest/c2.txt"}, headers=_AUTH)
with patch("core.artifact_lifecycle.trash_active_artifacts", return_value=0):
r = _client.post("/v1/files/delete", json={"path": "route-test-wd/dest/c2.txt"}, headers=_AUTH)
self.assertEqual(r.status_code, 200)
self.assertFalse((self.wd / "dest" / "c2.txt").exists())
# 删非空目录不带 recursive → 400
r = _client.post("/v1/files/delete", json={"path": "route-test-wd/sub"}, headers=_AUTH)
with patch("core.artifact_lifecycle.trash_active_artifacts", return_value=0):
r = _client.post("/v1/files/delete", json={"path": "route-test-wd/sub"}, headers=_AUTH)
self.assertEqual(r.status_code, 400)

View File

@ -103,6 +103,30 @@ def _task_file_target(root: Path, working_dir: Path, path: str, legacy: bool) ->
return candidates[0]
def _artifact_target(
root: Path,
user_id: UUID,
artifact_id: str,
) -> Path:
try:
aid = UUID(artifact_id)
except ValueError:
raise HTTPException(404, "invalid artifact id")
from core.storage.models import Artifact
with session_scope() as s:
current_path = s.execute(
select(Artifact.current_path).where(
Artifact.artifact_id == aid,
Artifact.user_id == user_id,
Artifact.status == "active",
)
).scalar_one_or_none()
if not current_path:
raise HTTPException(404, "artifact not found")
return safe_join(root, current_path)
async def _pptx_preview_response(target: Path, display_path: str) -> FileResponse:
from ..pptx_render import (
PptxConvertError,
@ -221,12 +245,17 @@ def register_file_routes(app, *, require_user) -> None:
task_id: str,
path: str,
legacy: bool = False,
artifact_id: str = "",
user_id: UUID = Depends(require_user),
):
"""Download a file addressed relative to the task's current working_dir."""
root = load_user_root(user_id)
_tid, working_dir = _task_working_dir(task_id, user_id, root)
target = _task_file_target(root, working_dir, path, legacy)
target = (
_artifact_target(root, user_id, artifact_id)
if artifact_id
else _task_file_target(root, working_dir, path, legacy)
)
return _regular_file_response(target, path)
@app.get("/v1/files/preview_pdf", tags=["files"])
@ -248,12 +277,17 @@ def register_file_routes(app, *, require_user) -> None:
task_id: str,
path: str,
legacy: bool = False,
artifact_id: str = "",
user_id: UUID = Depends(require_user),
):
"""Preview a PPT addressed relative to the task's current working_dir."""
root = load_user_root(user_id)
_tid, working_dir = _task_working_dir(task_id, user_id, root)
target = _task_file_target(root, working_dir, path, legacy)
target = (
_artifact_target(root, user_id, artifact_id)
if artifact_id
else _task_file_target(root, working_dir, path, legacy)
)
return await _pptx_preview_response(target, path)
@app.post("/v1/files/upload", tags=["files"])
@ -353,6 +387,8 @@ def register_file_routes(app, *, require_user) -> None:
- 顶层空目录 / 子级空目录无论 recursive 与否都可删:task.working_dir 字段不动,
下次 build_agent 按需 mkdir 重建,FS 目录视为可重生
- root 400;不存在 404
- 已通过结构化 artifact_refs 发布的文件先移入平台隐藏回收区用户侧仍立即消失
普通文件仍物理删除递归目录仅回收其中的 artifact
"""
root = load_user_root(user_id)
target = safe_join(root, body.path)
@ -360,8 +396,15 @@ def register_file_routes(app, *, require_user) -> None:
raise HTTPException(400, "cannot delete user_root")
if not target.exists():
raise HTTPException(404, f"path not found: {body.path}")
target_is_dir = target.is_dir()
if target_is_dir and not body.recursive:
try:
if any(target.iterdir()):
raise HTTPException(400, "delete failed: directory is not empty")
except OSError as e:
raise HTTPException(400, f"delete failed: {e}")
if target.is_dir() and body.recursive:
if target_is_dir and body.recursive:
is_top_level = target.parent.resolve() == root.resolve()
if is_top_level:
db_form = to_db_path(target)
@ -381,17 +424,28 @@ def register_file_routes(app, *, require_user) -> None:
)
try:
if target.is_dir():
from core.artifact_lifecycle import trash_active_artifacts
trashed_artifacts = trash_active_artifacts(
user_id=user_id,
user_root=root,
target=target,
)
if target_is_dir:
if body.recursive:
import shutil
shutil.rmtree(target)
else:
target.rmdir() # 非空目录会触发 OSError
else:
target.rmdir()
elif not trashed_artifacts:
target.unlink()
except OSError as e:
raise HTTPException(400, f"delete failed: {e}")
return {"ok": True, "path": body.path}
return {
"ok": True,
"path": body.path,
"artifacts_trashed": trashed_artifacts,
}
@app.post("/v1/files/rename", tags=["files"])
def rename_path(
@ -443,8 +497,21 @@ def register_file_routes(app, *, require_user) -> None:
if not is_top_level_dir:
try:
target.rename(new_target)
from core.artifact_lifecycle import rename_active_artifacts
rename_active_artifacts(
user_id=user_id,
user_root=root,
old_path=target,
new_path=new_target,
)
except OSError as e:
raise HTTPException(400, f"rename failed: {e}")
except Exception as e:
try:
new_target.rename(target)
except OSError:
pass
raise HTTPException(500, f"artifact metadata update failed: {e}")
return {
"ok": True,
"old": body.path,
@ -481,7 +548,8 @@ def register_file_routes(app, *, require_user) -> None:
- 不覆盖(任一目标已存在 409)
- 不能拷到自己 / 自身子树
- 顶层目录(可能是某 task working_dir)可以拷:新副本无 task 关联,不动 DB
- 顶层目录(可能是某 task working_dir)可以拷:不创建 task其中 artifact
为副本创建独立身份并记录 copied_from_artifact_id
- 部分失败语义:任一 FS 拷贝抛错 HTTPException,**前面已成功的拷贝保留**
( FS 事务可回滚;预检通过后通常不会失败,失败也是磁盘满 / 权限这类不能恢复的)
"""
@ -496,15 +564,31 @@ def register_file_routes(app, *, require_user) -> None:
shutil.copytree(src, target)
else:
shutil.copy2(src, target)
from core.artifact_lifecycle import copy_active_artifacts
artifacts_copied = copy_active_artifacts(
user_id=user_id,
user_root=root,
source=src,
target=target,
)
except OSError as e:
raise HTTPException(
500,
f"copy failed at {src.name!r}: {e} "
f"(已成功 {len(transferred)} 项,剩余未处理)",
)
except Exception as e:
if target.is_dir():
shutil.rmtree(target, ignore_errors=True)
else:
target.unlink(missing_ok=True)
raise HTTPException(
500, f"copy metadata failed at {src.name!r}: {e}"
)
transferred.append({
"old": rel_to(root, src),
"new": rel_to(root, target),
"artifacts_copied": artifacts_copied,
})
return {"ok": True, "count": len(transferred), "transferred": transferred}
@ -519,7 +603,7 @@ def register_file_routes(app, *, require_user) -> None:
- **顶层目录是某 task working_dir 409**,维持 "working_dir = 顶层目录" invariant
(允许的话 task working_dir 沉到子目录会让 rename 顶层的 DB-aware 逻辑失效;
用户想归档: DELETE task)
- 拷贝(`/copy`)无此限制,因为新副本无 task 关联
- 拷贝(`/copy`)无此限制因为副本不创建 taskartifact 身份独立复制
- 部分失败: /copy,前面成功的不回滚(`shutil.move` 失败几乎只发生在
跨卷拷贝中断,workspace 都在同一磁盘下罕见)
"""
@ -562,14 +646,30 @@ def register_file_routes(app, *, require_user) -> None:
target = dest / src.name
try:
shutil.move(str(src), str(target))
from core.artifact_lifecycle import rename_active_artifacts
artifacts_moved = rename_active_artifacts(
user_id=user_id,
user_root=root,
old_path=src,
new_path=target,
)
except OSError as e:
raise HTTPException(
500,
f"move failed at {src.name!r}: {e} "
f"(已成功 {len(transferred)} 项,剩余未处理)",
)
except Exception as e:
try:
shutil.move(str(target), str(src))
except OSError:
pass
raise HTTPException(
500, f"move metadata failed at {src.name!r}: {e}"
)
transferred.append({
"old": rel_to(root, src),
"new": rel_to(root, target),
"artifacts_moved": artifacts_moved,
})
return {"ok": True, "count": len(transferred), "transferred": transferred}

View File

@ -2532,19 +2532,19 @@ $("chat-stream").addEventListener("click", (e) => {
const chip = e.target.closest && e.target.closest(".art-chip");
if (chip) {
const rel = chip.dataset.rel;
if (rel) openFilePreview(rel, chip.dataset.taskId || "", chip.dataset.legacyPath === "1");
if (rel) openFilePreview(rel, chip.dataset.taskId || "", chip.dataset.legacyPath === "1", chip.dataset.artifactId || "");
return;
}
const htmlOpen = e.target.closest && e.target.closest(".art-html-open[data-rel]");
if (htmlOpen) {
const rel = htmlOpen.dataset.rel;
if (rel) openFilePreview(rel, htmlOpen.dataset.taskId || "", htmlOpen.dataset.legacyPath === "1");
if (rel) openFilePreview(rel, htmlOpen.dataset.taskId || "", htmlOpen.dataset.legacyPath === "1", htmlOpen.dataset.artifactId || "");
return;
}
const inlineImg = e.target.closest && e.target.closest(".art-media-image[data-rel]");
if (inlineImg) {
const rel = inlineImg.dataset.rel;
if (rel) openFilePreview(rel, inlineImg.dataset.taskId || "", inlineImg.dataset.legacyPath === "1");
if (rel) openFilePreview(rel, inlineImg.dataset.taskId || "", inlineImg.dataset.legacyPath === "1", inlineImg.dataset.artifactId || "");
return;
}
// 正文里的 markdown 链接:模型常把工作区相对路径写成 [<rel>](<rel>),renderMd 出 <a>。

View File

@ -188,21 +188,24 @@ export function renderArtifactBarHtml(rels, inlineMode = true, taskId = "", lega
const items = rels.map((item) => {
const ref = (item && typeof item === "object") ? item : { path: item };
const rel = String(ref.path || "");
const artifactAttr = ref.artifact_id
? ` data-artifact-id="${escapeHtml(String(ref.artifact_id))}"`
: "";
if (!rel) return "";
const name = String(ref.label || rel.split("/").pop() || rel);
const cat = _categorize(rel);
if ((inlineMode === true || inlineMode === "html") && cat === "html") {
return `<section class="art-html" data-rel="${escapeHtml(rel)}"${taskAttr}${legacyAttr} title="${escapeHtml(rel)}">
<div class="art-html-head"><span>${escapeHtml(name)}</span><button type="button" class="art-html-open" data-rel="${escapeHtml(rel)}"${taskAttr}${legacyAttr}></button></div>
return `<section class="art-html" data-rel="${escapeHtml(rel)}"${taskAttr}${legacyAttr}${artifactAttr} title="${escapeHtml(rel)}">
<div class="art-html-head"><span>${escapeHtml(name)}</span><button type="button" class="art-html-open" data-rel="${escapeHtml(rel)}"${taskAttr}${legacyAttr}${artifactAttr}></button></div>
<div class="art-html-viewport"><span class="art-media-loading">进入可视区域后加载</span></div>
</section>`;
}
if (inlineMode === true && (cat === "image" || cat === "video")) {
// 占位元素;插入 DOM 后 upgradeMediaArtifacts 异步 fetch blob → 填 <img>/<video>。
// 不在这里发请求避免 string-build 阶段失控的并发;upgrade 走 DOM walk 一次。
return `<span class="art-media art-media-${cat}" data-rel="${escapeHtml(rel)}" data-cat="${cat}"${taskAttr}${legacyAttr} title="${escapeHtml(rel)}"><span class="art-media-loading">${escapeHtml(name)} 加载中…</span></span>`;
return `<span class="art-media art-media-${cat}" data-rel="${escapeHtml(rel)}" data-cat="${cat}"${taskAttr}${legacyAttr}${artifactAttr} title="${escapeHtml(rel)}"><span class="art-media-loading">${escapeHtml(name)} 加载中…</span></span>`;
}
return `<button type="button" class="art-chip" data-rel="${escapeHtml(rel)}"${taskAttr}${legacyAttr} title="${escapeHtml(rel)} · 点击预览(可下载)">${renderArtifactChipContent(name, rel)}</button>`;
return `<button type="button" class="art-chip" data-rel="${escapeHtml(rel)}"${taskAttr}${legacyAttr}${artifactAttr} title="${escapeHtml(rel)} · 点击预览(可下载)">${renderArtifactChipContent(name, rel)}</button>`;
}).join("");
return `<div class="artifact-bar">${items}</div>`;
}
@ -214,17 +217,18 @@ const _mediaArtifactCache = new Map();
const _htmlArtifactCache = new Map();
const INLINE_HTML_MAX = 2 * 1024 * 1024;
function _artifactDownloadUrl(rel, taskId = "", legacy = false) {
function _artifactDownloadUrl(rel, taskId = "", legacy = false, artifactId = "") {
const base = taskId
? `/v1/tasks/${encodeURIComponent(taskId)}/files/download`
: "/v1/files/download";
return base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : "");
return base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : "")
+ (artifactId ? "&artifact_id=" + encodeURIComponent(artifactId) : "");
}
function _fetchMediaBlobUrl(rel, taskId = "", legacy = false) {
const key = `${taskId}:${legacy ? "legacy:" : ""}${rel}`;
function _fetchMediaBlobUrl(rel, taskId = "", legacy = false, artifactId = "") {
const key = `${taskId}:${artifactId}:${legacy ? "legacy:" : ""}${rel}`;
if (_mediaArtifactCache.has(key)) return _mediaArtifactCache.get(key);
const p = fetch(_artifactDownloadUrl(rel, taskId, legacy), {
const p = fetch(_artifactDownloadUrl(rel, taskId, legacy, artifactId), {
headers: { "Authorization": "Bearer " + state.token },
}).then(async (r) => {
if (!r.ok) throw new Error("HTTP " + r.status);
@ -235,10 +239,10 @@ function _fetchMediaBlobUrl(rel, taskId = "", legacy = false) {
return p;
}
function _fetchHtmlArtifact(rel, taskId = "", legacy = false) {
const key = `${taskId}:${legacy ? "legacy:" : ""}${rel}`;
function _fetchHtmlArtifact(rel, taskId = "", legacy = false, artifactId = "") {
const key = `${taskId}:${artifactId}:${legacy ? "legacy:" : ""}${rel}`;
if (_htmlArtifactCache.has(key)) return _htmlArtifactCache.get(key);
const p = fetch(_artifactDownloadUrl(rel, taskId, legacy), {
const p = fetch(_artifactDownloadUrl(rel, taskId, legacy, artifactId), {
headers: { "Authorization": "Bearer " + state.token },
}).then(async (r) => {
if (!r.ok) throw new Error("HTTP " + r.status);
@ -258,7 +262,7 @@ function _loadInlineHtml(node) {
const taskId = node.dataset.taskId || "";
const legacy = node.dataset.legacyPath === "1";
const viewport = node.querySelector(".art-html-viewport");
_fetchHtmlArtifact(rel, taskId, legacy).then((source) => {
_fetchHtmlArtifact(rel, taskId, legacy, node.dataset.artifactId || "").then((source) => {
if (!node.isConnected || !viewport) return;
viewport.innerHTML = "";
const frame = document.createElement("iframe");
@ -308,7 +312,7 @@ export function upgradeMediaArtifacts(root) {
const taskId = node.dataset.taskId || "";
const legacy = node.dataset.legacyPath === "1";
const cat = node.dataset.cat;
_fetchMediaBlobUrl(rel, taskId, legacy).then((url) => {
_fetchMediaBlobUrl(rel, taskId, legacy, node.dataset.artifactId || "").then((url) => {
node.innerHTML = "";
if (cat === "image") {
const img = document.createElement("img");
@ -334,11 +338,12 @@ export function upgradeMediaArtifacts(root) {
});
}
export function downloadFile(rel, taskId = "", legacy = false) {
export function downloadFile(rel, taskId = "", legacy = false, artifactId = "") {
const base = taskId
? `/v1/tasks/${encodeURIComponent(taskId)}/files/download`
: "/v1/files/download";
fetch(base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : ""), {
fetch(base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : "")
+ (artifactId ? "&artifact_id=" + encodeURIComponent(artifactId) : ""), {
headers: { "Authorization": "Bearer " + state.token },
}).then(async (r) => {
if (!r.ok) { message("下载失败:" + r.status, "error"); return; }

View File

@ -68,6 +68,7 @@ export function _categorize(rel) {
let _fpCurrentRel = null;
let _fpCurrentTaskId = "";
let _fpCurrentLegacy = false;
let _fpCurrentArtifactId = "";
// Markdown / HTML 共用“预览 / 源文件”切换。HTML 在 opaque-origin sandbox iframe
// 中运行:可执行脚本、加载 HTTPS 资源,但不能读取 zcbot 页面或发起表单/顶层跳转。
@ -248,24 +249,27 @@ function _bindBodyWheel(bodyEl) {
}, { passive: false });
}
function _fileDownloadUrl(rel, taskId = "", legacy = false) {
function _fileDownloadUrl(rel, taskId = "", legacy = false, artifactId = "") {
const base = taskId
? `/v1/tasks/${encodeURIComponent(taskId)}/files/download`
: "/v1/files/download";
return base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : "");
return base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : "")
+ (artifactId ? "&artifact_id=" + encodeURIComponent(artifactId) : "");
}
function _pptPreviewUrl(rel, taskId = "", legacy = false) {
function _pptPreviewUrl(rel, taskId = "", legacy = false, artifactId = "") {
const base = taskId
? `/v1/tasks/${encodeURIComponent(taskId)}/files/preview_pdf`
: "/v1/files/preview_pdf";
return base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : "");
return base + "?path=" + encodeURIComponent(rel) + (legacy ? "&legacy=true" : "")
+ (artifactId ? "&artifact_id=" + encodeURIComponent(artifactId) : "");
}
export async function openFilePreview(rel, taskId = "", legacy = false) {
export async function openFilePreview(rel, taskId = "", legacy = false, artifactId = "") {
_fpCurrentRel = rel;
_fpCurrentTaskId = taskId;
_fpCurrentLegacy = legacy;
_fpCurrentArtifactId = artifactId;
const name = rel.split("/").pop() || rel;
$("fp-name").textContent = name;
$("fp-meta").textContent = "";
@ -284,9 +288,9 @@ export async function openFilePreview(rel, taskId = "", legacy = false) {
const cat = _categorize(rel);
// pptx/ppt:后端转 PDF 再复用现成 PDF iframe(非下载原文件),首次稍候 + 失败回退下载。
if (cat === "ppt") { await _showPptAsPdf(rel, $("fp-body"), $("fp-meta"), _showFallback, taskId, legacy); return; }
if (cat === "ppt") { await _showPptAsPdf(rel, $("fp-body"), $("fp-meta"), _showFallback, taskId, legacy, artifactId); return; }
try {
const r = await fetch(_fileDownloadUrl(rel, taskId, legacy), {
const r = await fetch(_fileDownloadUrl(rel, taskId, legacy, artifactId), {
headers: { "Authorization": "Bearer " + state.token },
});
if (!r.ok) throw new Error("HTTP " + r.status);
@ -580,13 +584,13 @@ async function _showPdf(blob) {
}
// pptx/ppt → 后端转 PDF → PDF.js。main / mini 共用各自 body / meta / fallback。
async function _showPptAsPdf(rel, body, metaEl, fallbackFn, taskId = "", legacy = false) {
async function _showPptAsPdf(rel, body, metaEl, fallbackFn, taskId = "", legacy = false, artifactId = "") {
body.className = "body center";
body.innerHTML = `<div class="ph"><div class="preview-spinner"></div>由 PPT 转换为 PDF · 首次稍候…</div>`;
if (metaEl) metaEl.textContent = "";
let r;
try {
r = await fetch(_pptPreviewUrl(rel, taskId, legacy), {
r = await fetch(_pptPreviewUrl(rel, taskId, legacy, artifactId), {
headers: { "Authorization": "Bearer " + state.token },
});
} catch (e) {
@ -690,7 +694,7 @@ function _showFallback(msg) {
dl.className = "primary";
dl.textContent = "下载原文件";
dl.style.marginTop = "12px";
dl.onclick = () => { if (_fpCurrentRel) downloadFile(_fpCurrentRel, _fpCurrentTaskId, _fpCurrentLegacy); };
dl.onclick = () => { if (_fpCurrentRel) downloadFile(_fpCurrentRel, _fpCurrentTaskId, _fpCurrentLegacy, _fpCurrentArtifactId); };
ph.appendChild(document.createElement("br"));
ph.appendChild(br);
ph.appendChild(dl);
@ -709,6 +713,7 @@ export function closeFilePreview() {
_fpCurrentRel = null;
_fpCurrentTaskId = "";
_fpCurrentLegacy = false;
_fpCurrentArtifactId = "";
}
let _mpCurrentRel = null;
@ -835,7 +840,7 @@ _bindBodyWheel($("fp-body"));
_bindBodyWheel($("mp-body"));
$("fp-close").onclick = closeFilePreview;
$("fp-download").onclick = () => { if (_fpCurrentRel) downloadFile(_fpCurrentRel, _fpCurrentTaskId, _fpCurrentLegacy); };
$("fp-download").onclick = () => { if (_fpCurrentRel) downloadFile(_fpCurrentRel, _fpCurrentTaskId, _fpCurrentLegacy, _fpCurrentArtifactId); };
$("fp-mode-preview").onclick = () => _renderTextMode("fp", "preview");
$("fp-mode-source").onclick = () => _renderTextMode("fp", "source");
$("file-preview-modal").addEventListener("click", (e) => {