diff --git a/DESIGN.md b/DESIGN.md index 362fbe6..543c3c4 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -460,11 +460,11 @@ scheduled_jobs(§8.5) channel_bindings(§8.7,判别列+JSONB) 第一阶段以 `docs/windows-node-mvp-intranet.md` 为实现契约:Windows Node 只作为受控执行节点,通过出站 HTTP/WS 主动连接 zcbot;首批能力固定为 `origin.plot@v1`。长期方案中的 mTLS、Service/DesktopRunner 双进程、完整租约与多节点调度暂不进入 MVP,但 URL path、Node ID、Bearer Header 和任务协议保留原位升级空间。 -云端控制面使用独立的 `compute_node_enrollments` 与 `compute_nodes`,不复用用户外部系统连接。管理员创建的一次性注册码具有 128 bit 随机熵,数据库只保存 SHA-256 摘要;节点注册在行锁事务中校验有效期、预期名称和允许能力,成功后原子消费。每个节点获得独立高熵 Token,数据库只保存 bcrypt 强哈希,明文仅在注册响应出现一次。 +云端控制面使用独立的 `software_node_enrollments` 与 `software_nodes`,不复用用户外部系统连接。管理员创建的一次性注册码具有 128 bit 随机熵,数据库只保存 SHA-256 摘要;节点注册在行锁事务中校验有效期、预期名称和允许能力,成功后原子消费。每个节点获得独立高熵 Token,数据库只保存 bcrypt 强哈希,明文仅在注册响应出现一次。 -Node 通过 `Authorization: Bearer` 与 `X-Node-Id` 建立 `/v1/compute/nodes/connect` WebSocket。进程内 Connection Manager 保证同一节点单活,新连接关闭旧连接;`hello`/`heartbeat` 更新版本、容量、软件健康与最后在线时间。管理员禁用节点时先持久化禁用态,再关闭现有连接;断线收尾不得覆盖禁用态。当前单活只覆盖单 Web 进程,生产启用多实例前必须增加 Redis/PG fencing 或将 Node API 固定路由到单一控制面实例。 +Node 通过 `Authorization: Bearer` 与 `X-Node-Id` 建立 `/v1/software-nodes/connect` WebSocket。进程内 Connection Manager 保证同一节点单活,新连接关闭旧连接;`hello`/`heartbeat` 更新版本、容量、软件健康与最后在线时间。管理员禁用节点时先持久化禁用态,再关闭现有连接;断线收尾不得覆盖禁用态。当前单活只覆盖单 Web 进程,生产启用多实例前必须增加 Redis/PG fencing 或将 Node API 固定路由到单一控制面实例。 -第二阶段已增加 `compute_jobs` 账本与 `origin.plot@v1` 的 offer/accept 骨架。用户只能在本人 task 下以幂等键提交固定 schema;云端规范化请求并记录 SHA-256,按当前进程真实在线、能力匹配、健康且有空闲 slot 的 Node 创建短期 offer。Node 再次校验 schema、图形类型和输出格式,使用 write-through、flush 与原子 rename 先落本机任务目录,再回 `job_accept`;重复 job 只有 digest 一致才接受。过期或发送失败的 offer 回到队列,lease、Node 和 digest 不匹配的响应被拒绝。Node 接收后云端进入 `dispatched` 而非 `running`,并将 slot 降为 0;只有固定 Worker 真正启动后才进入 `origin_running`。 +第二阶段已增加 `software_jobs`(专业软件任务)账本与 `origin.plot@v1` 的 offer/accept 骨架。用户只能在本人 task 下以幂等键提交固定 schema;云端规范化请求并记录 SHA-256,按当前进程真实在线、能力匹配、健康且有空闲 slot 的 Node 创建短期 offer。Node 再次校验 schema、图形类型和输出格式,使用 write-through、flush 与原子 rename 先落本机任务目录,再回 `job_accept`;重复 job 只有 digest 一致才接受。过期或发送失败的 offer 回到队列,lease、Node 和 digest 不匹配的响应被拒绝。Node 接收后云端进入 `dispatched` 而非 `running`,并将 slot 降为 0;只有固定 Worker 真正启动后才进入 `origin_running`。 第三阶段补齐输入下载与恢复状态协议:`input_id` 固定为用户已有 artifact UUID,提交时快照文件名、大小和 SHA-256,只允许 CSV/XLSX/JSON 且不超过 100 MiB。Node 以自身 Bearer 身份访问任务绑定的只读下载端点,流式写入本 job 的 `input/`,同时限制声明大小并校验 SHA-256,完成后原子 rename;不暴露工作区路径。Node 会原子读取/补报 `terminal.json`,断线后云端把活动任务标记 `disconnected` 并保留 Node/lease,重连按 job、lease、digest 恢复下载或幂等补报终态,不自动重派。 diff --git a/PROGRESS.md b/PROGRESS.md index b0bd5d5..b8385f6 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -341,7 +341,7 @@ core/llm_transport.py 438 ← wire 层健壮性:畸形/吐空检测+留 core/tool_registry.py 264 ← 声明式工具注册表((组名,gate,factory);secret/host 工具按实际能力 gate) core/context.py 95 ← LLM 调用前压缩旧 tool / load_skill 消息(带压力门槛),保 tool_call 协议字段 core/external_systems/*.py ← 外部系统目录/用户授权/凭据加密 + 通用 OpenAPI/MCP connector -core/compute_nodes.py ← Windows Node 注册码、身份认证与运行状态 +core/software_nodes.py ← Windows Node 注册码、身份认证与运行状态 core/sinks.py 101 core/paths.py 50 ← task_dir db form 归一 core/probe.py 243 @@ -363,7 +363,7 @@ tools/{base,output,fs,shell,run_python,skill_tool,skill_authoring,media_common,s main.py ~210 ← 入口:web / db / probe / user / sandbox check db/migrations/versions/ 0001-0030 web/app.py ~210 ← 工厂 + lifespan 编排(07-23 拆分;路由在 routers/,协程在 background 等) -web/routers/*.py ← 含 external_systems 用户连接与 compute_nodes 节点路由 +web/routers/*.py ← 含 external_systems 用户连接与 software_nodes 节点路由 web/{background,scheduler_runner,wechat_runner}.py ← lifespan 后台协程按域析出 web/{runs,common,schemas,model_gate,userfiles}.py ← BG worker/共享 helper/请求体/档位门控/路径安全 web/auth.py ~190 ← 邮箱密码 + platform_key → JWT diff --git a/RUN.md b/RUN.md index ff518db..fb5acfd 100644 --- a/RUN.md +++ b/RUN.md @@ -1054,14 +1054,14 @@ sudo xfs_quota -x -c "limit -p bhard=10g zcbot_" /opt ### Windows Node 内网 MVP(开发中) -先执行 `alembic upgrade head` 创建 `compute_node_enrollments` 和 `compute_nodes`。不要在未确认目标数据库时运行迁移;本机 `.env` 的 `ZCBOT_DB_URL` 可能是生产隧道。 +先执行 `alembic upgrade head` 创建 `software_node_enrollments` 和 `software_nodes`。不要在未确认目标数据库时运行迁移;本机 `.env` 的 `ZCBOT_DB_URL` 可能是生产隧道。 云端当前提供: -- 管理员 `POST /v1/admin/compute-node-enrollments` 创建一次性注册码; -- Node `POST /v1/compute/nodes/enroll` 注册并一次性取得 `node_id`、`node_token`; -- Node 携带 `Authorization: Bearer ` 和 `X-Node-Id` 连接 `WS /v1/compute/nodes/connect`; -- 管理员 `GET /v1/admin/compute-nodes` 查看节点,`PATCH /v1/admin/compute-nodes/{node_id}` 启停节点,`DELETE /v1/admin/compute-nodes/{node_id}` 永久删除节点身份。 +- 管理员 `POST /v1/admin/software-node-enrollments` 创建一次性注册码; +- Node `POST /v1/software-nodes/enroll` 注册并一次性取得 `node_id`、`node_token`; +- Node 携带 `Authorization: Bearer ` 和 `X-Node-Id` 连接 `WS /v1/software-nodes/connect`; +- 管理员 `GET /v1/admin/software-nodes` 查看节点,`PATCH /v1/admin/software-nodes/{node_id}` 启停节点,`DELETE /v1/admin/software-nodes/{node_id}` 永久删除节点身份。 Node API 只能绑定受控内网地址并由安全组限制来源 IP。当前 HTTP/WS 链路不加密;跨安全域、公网或不可信终端接入前,必须先升级 HTTPS/WSS。多 Web 实例部署时,Node API 暂时固定路由到单一实例,直至 Connection Manager 增加跨实例 fencing。 @@ -1088,7 +1088,7 @@ windows-node\install-origin-runtime.ps1 -BootstrapPython D:\programs\Python312\p 直接双击 EXE 启动托盘 UI:红点为未注册/身份失效,黄点为连接中,绿点为在线;双击托盘图标打开配置窗。原 CLI 注册入口继续保留,无 UI 模式使用 `Zcbot.WindowsNode.exe run --headless`。 -经 nginx 反代时,`/v1/compute/nodes/connect` 必须单独透传 WebSocket Upgrade/Connection 头并设置长连接超时,配置见 `deploy/nginx/zcbot.conf.example`。若注册成功后节点持续显示“连接中断,等待重连”,先用 WebSocket 握手检查该路径;返回普通 HTTP 404 通常表示请求落入了清空 `Connection` 头的默认 location。 +经 nginx 反代时,`/v1/software-nodes/connect` 必须单独透传 WebSocket Upgrade/Connection 头并设置长连接超时,配置见 `deploy/nginx/zcbot.conf.example`。若注册成功后节点持续显示“连接中断,等待重连”,先用 WebSocket 握手检查该路径;返回普通 HTTP 404 通常表示请求落入了清空 `Connection` 头的默认 location。 - **入口**:`main.py`(`web / db / probe / user`)→ `core/agent_builder.py::build_agent` diff --git a/core/compute_jobs.py b/core/software_jobs.py similarity index 77% rename from core/compute_jobs.py rename to core/software_jobs.py index 0f4a253..929bb49 100644 --- a/core/compute_jobs.py +++ b/core/software_jobs.py @@ -1,4 +1,4 @@ -"""Origin 受控计算任务的校验、幂等持久化和 offer 状态机。""" +"""专业软件任务的校验、幂等持久化和 offer 状态机。""" from __future__ import annotations @@ -11,9 +11,9 @@ from uuid import UUID, uuid4 from sqlalchemy import select from sqlalchemy.exc import IntegrityError -from core.compute_nodes import ComputeNodeError, SUPPORTED_CAPABILITIES +from core.software_nodes import SUPPORTED_CAPABILITIES from core.storage.engine import session_scope -from core.storage.models import Artifact, ComputeJob, ComputeNode, Task +from core.storage.models import Artifact, SoftwareJob, SoftwareNode, Task OFFER_SECONDS = 60 ALLOWED_PLOT_TYPES = frozenset( @@ -23,6 +23,10 @@ ALLOWED_OUTPUT_FORMATS = frozenset({"opju", "png", "svg", "pdf"}) ALLOWED_INPUT_SUFFIXES = frozenset({".csv", ".xlsx", ".json"}) MAX_INPUT_BYTES = 100 * 1024 * 1024 MAX_OUTPUT_ARTIFACT_BYTES = 256 * 1024 * 1024 +class SoftwareJobError(Exception): + pass + + MAX_OUTPUT_TOTAL_BYTES = 512 * 1024 * 1024 OUTPUT_ARTIFACTS = { "project": ("project.opju", "application/x-origin-project", "opju"), @@ -40,41 +44,41 @@ def _has_only(value: dict, fields: set[str]) -> bool: def _canonical_request(request: dict) -> tuple[dict, str]: if not isinstance(request, dict) or set(request) != {"schema_version", "input", "plot", "output"}: - raise ComputeNodeError("invalid origin plot request fields") + raise SoftwareJobError("invalid origin plot request fields") if request.get("schema_version") != 1: - raise ComputeNodeError("unsupported origin plot schema version") + raise SoftwareJobError("unsupported origin plot schema version") input_spec = request.get("input") plot = request.get("plot") output = request.get("output") if not all(isinstance(item, dict) for item in (input_spec, plot, output)): - raise ComputeNodeError("origin plot request sections must be objects") + raise SoftwareJobError("origin plot request sections must be objects") if not _has_only(input_spec, {"input_id", "sheet"}): - raise ComputeNodeError("unsupported origin input fields") + raise SoftwareJobError("unsupported origin input fields") try: UUID(str(input_spec.get("input_id") or "")) except ValueError as exc: - raise ComputeNodeError("input.input_id must be an artifact UUID") from exc + raise SoftwareJobError("input.input_id must be an artifact UUID") from exc if "sheet" in input_spec and ( not isinstance(input_spec["sheet"], str) or not 1 <= len(input_spec["sheet"]) <= 128 ): - raise ComputeNodeError("input.sheet must be a string") + raise SoftwareJobError("input.sheet must be a string") if not _has_only( plot, {"type", "x", "y", "template", "title", "x_axis", "y_axis", "legend", "error_bars"}, ): - raise ComputeNodeError("unsupported origin plot fields") + raise SoftwareJobError("unsupported origin plot fields") if plot.get("type") not in ALLOWED_PLOT_TYPES: - raise ComputeNodeError("unsupported origin plot type") + raise SoftwareJobError("unsupported origin plot type") if "title" in plot and ( not isinstance(plot["title"], str) or len(plot["title"]) > 500 ): - raise ComputeNodeError("plot.title must be a string") + raise SoftwareJobError("plot.title must be a string") if plot.get("template", "publication_double_column") != "publication_double_column": - raise ComputeNodeError("unsupported origin plot template") + raise SoftwareJobError("unsupported origin plot template") x_column = plot.get("x") y_columns = plot.get("y") if not isinstance(x_column, str) or not 1 <= len(x_column) <= 128: - raise ComputeNodeError("plot.x must be a column name") + raise SoftwareJobError("plot.x must be a column name") if isinstance(y_columns, str): y_columns = [y_columns] if ( @@ -83,7 +87,7 @@ def _canonical_request(request: dict) -> tuple[dict, str]: or len(y_columns) != len(set(y_columns)) or any(not isinstance(item, str) or not 1 <= len(item) <= 128 for item in y_columns) ): - raise ComputeNodeError("plot.y must contain 1 to 16 unique column names") + raise SoftwareJobError("plot.y must contain 1 to 16 unique column names") for axis_name in ("x_axis", "y_axis"): axis = plot.get(axis_name) if axis is not None and ( @@ -95,7 +99,7 @@ def _canonical_request(request: dict) -> tuple[dict, str]: for name in ("title", "unit") ) ): - raise ComputeNodeError(f"invalid {axis_name}") + raise SoftwareJobError(f"invalid {axis_name}") legend = plot.get("legend") if legend is not None and ( not isinstance(legend, dict) @@ -104,21 +108,21 @@ def _canonical_request(request: dict) -> tuple[dict, str]: or legend.get("enabled", True) is not True or legend.get("position", "top_right") != "top_right" ): - raise ComputeNodeError("invalid plot.legend") + raise SoftwareJobError("invalid plot.legend") if plot.get("error_bars") is not None: - raise ComputeNodeError("error bars are not supported in origin.plot@v1") + raise SoftwareJobError("error bars are not supported in origin.plot@v1") if not _has_only(output, {"formats", "dpi", "capture_screenshots", "record_video"}): - raise ComputeNodeError("unsupported origin output fields") + raise SoftwareJobError("unsupported origin output fields") if any( name in output and not isinstance(output[name], bool) for name in ("capture_screenshots", "record_video") ): - raise ComputeNodeError("origin output capture flags must be boolean") + raise SoftwareJobError("origin output capture flags must be boolean") if output.get("record_video", False): - raise ComputeNodeError("origin video recording is not supported") + raise SoftwareJobError("origin video recording is not supported") dpi = output.get("dpi", 300) if not isinstance(dpi, int) or isinstance(dpi, bool) or not 72 <= dpi <= 1200: - raise ComputeNodeError("output.dpi must be between 72 and 1200") + raise SoftwareJobError("output.dpi must be between 72 and 1200") formats = output.get("formats") if ( not isinstance(formats, list) @@ -126,15 +130,15 @@ def _canonical_request(request: dict) -> tuple[dict, str]: or len(formats) != len(set(formats)) or any(item not in ALLOWED_OUTPUT_FORMATS for item in formats) ): - raise ComputeNodeError("output.formats contains unsupported values") + raise SoftwareJobError("output.formats contains unsupported values") encoded = json.dumps(request, ensure_ascii=False, sort_keys=True, separators=(",", ":")) if len(encoded.encode("utf-8")) > 256 * 1024: - raise ComputeNodeError("origin plot request is too large") + raise SoftwareJobError("origin plot request is too large") normalized = json.loads(encoded) return normalized, sha256(encoded.encode("utf-8")).hexdigest() -def _job_dict(row: ComputeJob) -> dict: +def _job_dict(row: SoftwareJob) -> dict: return { "job_id": str(row.job_id), "task_id": str(row.task_id), @@ -163,16 +167,16 @@ def create_job( ) -> tuple[dict, bool]: key = idempotency_key.strip() if not key or len(key) > 200: - raise ComputeNodeError("idempotency_key must contain 1 to 200 characters") + raise SoftwareJobError("idempotency_key must contain 1 to 200 characters") if capability not in SUPPORTED_CAPABILITIES: - raise ComputeNodeError("unsupported capability") + raise SoftwareJobError("unsupported capability") normalized, digest = _canonical_request(request) with session_scope() as session: task = session.execute( select(Task.task_id).where(Task.task_id == task_id, Task.user_id == user_id) ).first() if task is None: - raise ComputeNodeError("task not found") + raise SoftwareJobError("task not found") artifact_id = UUID(normalized["input"]["input_id"]) artifact = session.execute( select(Artifact).where( @@ -182,10 +186,10 @@ def create_job( ) ).scalar_one_or_none() if artifact is None: - raise ComputeNodeError("input artifact not found") + raise SoftwareJobError("input artifact not found") suffix = "." + artifact.current_path.rsplit(".", 1)[-1].lower() if "." in artifact.current_path else "" if suffix not in ALLOWED_INPUT_SUFFIXES: - raise ComputeNodeError("input artifact type is not supported") + raise SoftwareJobError("input artifact type is not supported") if ( artifact.size_bytes is None or artifact.size_bytes < 0 @@ -193,7 +197,7 @@ def create_job( or not artifact.content_sha256 or len(artifact.content_sha256) != 64 ): - raise ComputeNodeError("input artifact metadata is incomplete or too large") + raise SoftwareJobError("input artifact metadata is incomplete or too large") input_manifest = { "artifact_id": str(artifact.artifact_id), "filename": artifact.current_path.replace("\\", "/").rsplit("/", 1)[-1], @@ -201,9 +205,9 @@ def create_job( "sha256": artifact.content_sha256, } existing = session.execute( - select(ComputeJob).where( - ComputeJob.user_id == user_id, - ComputeJob.idempotency_key == key, + select(SoftwareJob).where( + SoftwareJob.user_id == user_id, + SoftwareJob.idempotency_key == key, ) ).scalar_one_or_none() if existing is not None: @@ -212,9 +216,9 @@ def create_job( or existing.capability != capability or existing.request_digest != digest ): - raise ComputeNodeError("idempotency key was already used for a different request") + raise SoftwareJobError("idempotency key was already used for a different request") return _job_dict(existing), False - row = ComputeJob( + row = SoftwareJob( job_id=uuid4(), user_id=user_id, task_id=task_id, @@ -236,9 +240,9 @@ def create_job( return _job_dict(row), True except IntegrityError: existing = session.execute( - select(ComputeJob).where( - ComputeJob.user_id == user_id, - ComputeJob.idempotency_key == key, + select(SoftwareJob).where( + SoftwareJob.user_id == user_id, + SoftwareJob.idempotency_key == key, ) ).scalar_one() if ( @@ -246,7 +250,7 @@ def create_job( or existing.capability != capability or existing.request_digest != digest ): - raise ComputeNodeError( + raise SoftwareJobError( "idempotency key was already used for a different request" ) return _job_dict(existing), False @@ -255,7 +259,7 @@ def create_job( def get_job(user_id: UUID, job_id: UUID) -> dict | None: with session_scope() as session: row = session.execute( - select(ComputeJob).where(ComputeJob.job_id == job_id, ComputeJob.user_id == user_id) + select(SoftwareJob).where(SoftwareJob.job_id == job_id, SoftwareJob.user_id == user_id) ).scalar_one_or_none() return _job_dict(row) if row else None @@ -267,10 +271,10 @@ def offer_next_job(node_ids: set[UUID]) -> dict | None: now = datetime.now(timezone.utc) with session_scope() as session: expired = session.execute( - select(ComputeJob) + select(SoftwareJob) .where( - ComputeJob.status == "offered", - ComputeJob.lease_expires_at <= now, + SoftwareJob.status == "offered", + SoftwareJob.lease_expires_at <= now, ) .with_for_update(skip_locked=True) ).scalars() @@ -280,9 +284,9 @@ def offer_next_job(node_ids: set[UUID]) -> dict | None: item.lease_id = None item.lease_expires_at = None job = session.execute( - select(ComputeJob) - .where(ComputeJob.status == "queued") - .order_by(ComputeJob.created_at, ComputeJob.job_id) + select(SoftwareJob) + .where(SoftwareJob.status == "queued") + .order_by(SoftwareJob.created_at, SoftwareJob.job_id) .with_for_update(skip_locked=True) .limit(1) ).scalar_one_or_none() @@ -290,16 +294,16 @@ def offer_next_job(node_ids: set[UUID]) -> dict | None: return None busy_node_ids = set( session.execute( - select(ComputeJob.node_id).where( - ComputeJob.node_id.is_not(None), - ComputeJob.status.in_({"offered", "dispatched", "running"}), + select(SoftwareJob.node_id).where( + SoftwareJob.node_id.is_not(None), + SoftwareJob.status.in_({"offered", "dispatched", "running"}), ) ).scalars() ) nodes = session.execute( - select(ComputeNode) - .where(ComputeNode.node_id.in_(node_ids), ComputeNode.status == "online") - .order_by(ComputeNode.last_seen_at.desc()) + select(SoftwareNode) + .where(SoftwareNode.node_id.in_(node_ids), SoftwareNode.status == "online") + .order_by(SoftwareNode.last_seen_at.desc()) ).scalars() node = next( ( @@ -330,7 +334,7 @@ def offer_next_job(node_ids: set[UUID]) -> dict | None: "request": job.request, "input_transfer": { **job.input_manifest, - "download_path": f"/v1/compute/jobs/{job.job_id}/input", + "download_path": f"/v1/software-jobs/{job.job_id}/input", }, }, } @@ -340,10 +344,10 @@ def get_job_input(node_id: UUID, job_id: UUID) -> dict | None: """返回任务绑定的 artifact 定位信息;调用方仍需在 user_root 内安全解析。""" with session_scope() as session: job = session.execute( - select(ComputeJob).where( - ComputeJob.job_id == job_id, - ComputeJob.node_id == node_id, - ComputeJob.status.in_({"offered", "dispatched", "running", "disconnected"}), + select(SoftwareJob).where( + SoftwareJob.job_id == job_id, + SoftwareJob.node_id == node_id, + SoftwareJob.status.in_({"offered", "dispatched", "running", "disconnected"}), ) ).scalar_one_or_none() if job is None: @@ -369,9 +373,9 @@ def get_job_output_context(node_id: UUID, job_id: UUID, lease_id: UUID, digest: """返回 Node 输出上传上下文,不向 Node 暴露任何云端文件路径。""" with session_scope() as session: row = session.execute( - select(ComputeJob, Task.working_dir) - .join(Task, Task.task_id == ComputeJob.task_id) - .where(ComputeJob.job_id == job_id) + select(SoftwareJob, Task.working_dir) + .join(Task, Task.task_id == SoftwareJob.task_id) + .where(SoftwareJob.job_id == job_id) ).one_or_none() if row is None: return None @@ -395,7 +399,7 @@ def get_job_output_context(node_id: UUID, job_id: UUID, lease_id: UUID, digest: def validate_output_manifest(request: dict, manifest: object) -> list[dict]: if not isinstance(manifest, list): - raise ComputeNodeError("job artifact manifest must be a list") + raise SoftwareJobError("job artifact manifest must be a list") requested_formats = set(request.get("output", {}).get("formats") or []) expected_ids = {"plot_spec", "provenance"} expected_ids.update( @@ -404,7 +408,7 @@ def validate_output_manifest(request: dict, manifest: object) -> list[dict]: if output_format in requested_formats ) if len(manifest) != len(expected_ids): - raise ComputeNodeError("job artifact manifest is incomplete") + raise SoftwareJobError("job artifact manifest is incomplete") normalized: list[dict] = [] seen: set[str] = set() total = 0 @@ -412,24 +416,24 @@ def validate_output_manifest(request: dict, manifest: object) -> list[dict]: if not isinstance(raw, dict) or set(raw) != { "artifact_id", "filename", "media_type", "size_bytes", "sha256" }: - raise ComputeNodeError("job artifact manifest entry is invalid") + raise SoftwareJobError("job artifact manifest entry is invalid") local_id = raw.get("artifact_id") if local_id not in expected_ids or local_id in seen: - raise ComputeNodeError("job artifact manifest identity is invalid") + raise SoftwareJobError("job artifact manifest identity is invalid") filename, media_type, _ = OUTPUT_ARTIFACTS[local_id] size = raw.get("size_bytes") digest = raw.get("sha256") if raw.get("filename") != filename or raw.get("media_type") != media_type: - raise ComputeNodeError("job artifact manifest metadata does not match its identity") + raise SoftwareJobError("job artifact manifest metadata does not match its identity") if not isinstance(size, int) or isinstance(size, bool) or not 1 <= size <= MAX_OUTPUT_ARTIFACT_BYTES: - raise ComputeNodeError("job output artifact size is invalid") + raise SoftwareJobError("job output artifact size is invalid") if not isinstance(digest, str) or not re.fullmatch(r"[0-9a-f]{64}", digest): - raise ComputeNodeError("job output artifact digest is invalid") + raise SoftwareJobError("job output artifact digest is invalid") total += size seen.add(local_id) normalized.append(dict(raw)) if seen != expected_ids or total > MAX_OUTPUT_TOTAL_BYTES: - raise ComputeNodeError("job artifact manifest is incomplete or too large") + raise SoftwareJobError("job artifact manifest is incomplete or too large") return normalized @@ -442,7 +446,7 @@ def abandon_offer(node_id: UUID, payload: dict) -> None: return with session_scope() as session: job = session.execute( - select(ComputeJob).where(ComputeJob.job_id == job_id).with_for_update() + select(SoftwareJob).where(SoftwareJob.job_id == job_id).with_for_update() ).scalar_one_or_none() if ( job is not None @@ -461,14 +465,14 @@ def respond_to_offer(node_id: UUID, *, accepted: bool, payload: dict) -> None: job_id = UUID(str(payload.get("job_id", ""))) lease_id = UUID(str(payload.get("lease_id", ""))) except ValueError as exc: - raise ComputeNodeError("invalid job offer response identity") from exc + raise SoftwareJobError("invalid job offer response identity") from exc now = datetime.now(timezone.utc) with session_scope() as session: job = session.execute( - select(ComputeJob).where(ComputeJob.job_id == job_id).with_for_update() + select(SoftwareJob).where(SoftwareJob.job_id == job_id).with_for_update() ).scalar_one_or_none() if job is None or job.node_id != node_id or job.lease_id != lease_id: - raise ComputeNodeError("job offer is stale or does not belong to this node") + raise SoftwareJobError("job offer is stale or does not belong to this node") if ( accepted and job.status in {"dispatched", "running", "succeeded", "failed", "cancelled"} @@ -476,16 +480,16 @@ def respond_to_offer(node_id: UUID, *, accepted: bool, payload: dict) -> None: ): return if job.status != "offered": - raise ComputeNodeError("job offer is stale or does not belong to this node") + raise SoftwareJobError("job offer is stale or does not belong to this node") if job.lease_expires_at is None or job.lease_expires_at <= now: job.status = "queued" job.node_id = None job.lease_id = None job.lease_expires_at = None - raise ComputeNodeError("job offer has expired") + raise SoftwareJobError("job offer has expired") if accepted: if payload.get("request_digest") != job.request_digest: - raise ComputeNodeError("job request digest mismatch") + raise SoftwareJobError("job request digest mismatch") job.status = "dispatched" job.stage = "accepted" job.error = {} @@ -503,21 +507,21 @@ def update_job_state(node_id: UUID, payload: dict) -> None: progress = payload.get("progress") metrics = payload.get("metrics") or {} if not stage or len(stage) > 100: - raise ComputeNodeError("job stage is required") + raise SoftwareJobError("job stage is required") if not isinstance(progress, int) or isinstance(progress, bool) or not 0 <= progress <= 100: - raise ComputeNodeError("job progress must be between 0 and 100") + raise SoftwareJobError("job progress must be between 0 and 100") if not isinstance(metrics, dict) or len(json.dumps(metrics, ensure_ascii=False)) > 64 * 1024: - raise ComputeNodeError("job metrics are invalid") + raise SoftwareJobError("job metrics are invalid") now = datetime.now(timezone.utc) with session_scope() as session: job = session.execute( - select(ComputeJob).where(ComputeJob.job_id == job_id).with_for_update() + select(SoftwareJob).where(SoftwareJob.job_id == job_id).with_for_update() ).scalar_one_or_none() _assert_job_message(job, node_id, lease_id, digest) if job.status in {"succeeded", "failed", "cancelled"}: return if not _can_accept_state(job.status): - raise ComputeNodeError("job state cannot advance from its current status") + raise SoftwareJobError("job state cannot advance from its current status") job.status = ( "dispatched" if stage in {"accepted", "waiting_input", "ready_to_run"} @@ -534,17 +538,17 @@ def record_job_terminal(node_id: UUID, payload: dict) -> None: job_id, lease_id, digest = _message_identity(payload) terminal_status = payload.get("status") if terminal_status not in {"succeeded", "failed", "cancelled"}: - raise ComputeNodeError("invalid job terminal status") + raise SoftwareJobError("invalid job terminal status") error = payload.get("error") or {} manifest = payload.get("artifact_manifest") or [] if not isinstance(error, dict) or len(json.dumps(error, ensure_ascii=False)) > 64 * 1024: - raise ComputeNodeError("job terminal error is invalid") + raise SoftwareJobError("job terminal error is invalid") if not isinstance(manifest, list) or len(json.dumps(manifest, ensure_ascii=False)) > 256 * 1024: - raise ComputeNodeError("job artifact manifest is invalid") + raise SoftwareJobError("job artifact manifest is invalid") now = datetime.now(timezone.utc) with session_scope() as session: job = session.execute( - select(ComputeJob).where(ComputeJob.job_id == job_id).with_for_update() + select(SoftwareJob).where(SoftwareJob.job_id == job_id).with_for_update() ).scalar_one_or_none() _assert_job_message(job, node_id, lease_id, digest) if terminal_status == "succeeded": @@ -569,13 +573,13 @@ def record_job_terminal(node_id: UUID, payload: dict) -> None: or not item["path"].startswith(f"origin/{job.job_id}/") for item in manifest ): - raise ComputeNodeError("successful job artifacts have not been published") + raise SoftwareJobError("successful job artifacts have not been published") if job.status in {"succeeded", "failed", "cancelled"}: if job.status != terminal_status: - raise ComputeNodeError("job terminal status conflicts with existing terminal") + raise SoftwareJobError("job terminal status conflicts with existing terminal") return if job.status not in {"offered", "dispatched", "running", "disconnected"}: - raise ComputeNodeError("job terminal cannot advance from its current status") + raise SoftwareJobError("job terminal cannot advance from its current status") job.status = terminal_status job.stage = "terminal" job.progress = 100 if terminal_status == "succeeded" else job.progress @@ -596,10 +600,10 @@ def mark_node_jobs_disconnected(node_id: UUID) -> None: """连接丢失后保留 Node 归属和 lease,禁止任务被自动重派。""" with session_scope() as session: jobs = session.execute( - select(ComputeJob) + select(SoftwareJob) .where( - ComputeJob.node_id == node_id, - ComputeJob.status.in_({"dispatched", "running"}), + SoftwareJob.node_id == node_id, + SoftwareJob.status.in_({"dispatched", "running"}), ) .with_for_update() ).scalars() @@ -616,15 +620,15 @@ def _message_identity(payload: dict) -> tuple[UUID, UUID, str]: job_id = UUID(str(payload.get("job_id", ""))) lease_id = UUID(str(payload.get("lease_id", ""))) except ValueError as exc: - raise ComputeNodeError("invalid job message identity") from exc + raise SoftwareJobError("invalid job message identity") from exc digest = str(payload.get("request_digest") or "") if len(digest) != 64: - raise ComputeNodeError("invalid job request digest") + raise SoftwareJobError("invalid job request digest") return job_id, lease_id, digest def _assert_job_message( - job: ComputeJob | None, + job: SoftwareJob | None, node_id: UUID, lease_id: UUID, digest: str, @@ -635,4 +639,4 @@ def _assert_job_message( or job.lease_id != lease_id or job.request_digest != digest ): - raise ComputeNodeError("job message does not belong to this node or lease") + raise SoftwareJobError("job message does not belong to this node or lease") diff --git a/core/compute_nodes.py b/core/software_nodes.py similarity index 80% rename from core/compute_nodes.py rename to core/software_nodes.py index 6a25c1e..f5e49fd 100644 --- a/core/compute_nodes.py +++ b/core/software_nodes.py @@ -11,13 +11,13 @@ import bcrypt from sqlalchemy import select from core.storage.engine import session_scope -from core.storage.models import ComputeNode, ComputeNodeEnrollment +from core.storage.models import SoftwareNode, SoftwareNodeEnrollment SUPPORTED_CAPABILITIES = frozenset({"origin.plot@v1"}) MAX_ENROLLMENT_FAILURES = 5 -class ComputeNodeError(Exception): +class SoftwareNodeError(Exception): pass @@ -46,12 +46,12 @@ def create_enrollment( ) -> dict: allowed = list(dict.fromkeys(capabilities or ["origin.plot@v1"])) if not allowed or any(item not in SUPPORTED_CAPABILITIES for item in allowed): - raise ComputeNodeError("unsupported capability") + raise SoftwareNodeError("unsupported capability") if not 60 <= ttl_seconds <= 3600: - raise ComputeNodeError("ttl_seconds must be between 60 and 3600") + raise SoftwareNodeError("ttl_seconds must be between 60 and 3600") code = "ZCN-" + secrets.token_hex(16).upper() expires_at = datetime.now(timezone.utc) + timedelta(seconds=ttl_seconds) - row = ComputeNodeEnrollment( + row = SoftwareNodeEnrollment( enrollment_id=uuid4(), code_hash=_enrollment_digest(code), expected_name=(expected_name or "").strip() or None, @@ -80,22 +80,22 @@ def enroll_node( name = node_name.strip() requested = list(dict.fromkeys(capabilities)) if not name or not requested: - raise ComputeNodeError("node_name and capabilities are required") + raise SoftwareNodeError("node_name and capabilities are required") now = datetime.now(timezone.utc) error: str | None = None with session_scope() as session: enrollment = session.execute( - select(ComputeNodeEnrollment) + select(SoftwareNodeEnrollment) .where( - ComputeNodeEnrollment.code_hash == _enrollment_digest(enrollment_code), - ComputeNodeEnrollment.consumed_at.is_(None), - ComputeNodeEnrollment.expires_at > now, - ComputeNodeEnrollment.failed_attempts < MAX_ENROLLMENT_FAILURES, + SoftwareNodeEnrollment.code_hash == _enrollment_digest(enrollment_code), + SoftwareNodeEnrollment.consumed_at.is_(None), + SoftwareNodeEnrollment.expires_at > now, + SoftwareNodeEnrollment.failed_attempts < MAX_ENROLLMENT_FAILURES, ) .with_for_update() ).scalar_one_or_none() if enrollment is None: - raise ComputeNodeError("invalid or expired enrollment code") + raise SoftwareNodeError("invalid or expired enrollment code") enrollment.failed_attempts += 1 if enrollment.expected_name and enrollment.expected_name != name: error = "node name does not match enrollment" @@ -103,7 +103,7 @@ def enroll_node( error = "capability is not allowed by enrollment" else: existing = session.execute( - select(ComputeNode.node_id).where(ComputeNode.install_id == install_id) + select(SoftwareNode.node_id).where(SoftwareNode.install_id == install_id) ).first() if existing is not None: error = "install is already enrolled" @@ -111,7 +111,7 @@ def enroll_node( token = secrets.token_urlsafe(48) node_id = uuid4() session.add( - ComputeNode( + SoftwareNode( node_id=node_id, name=name, install_id=install_id, @@ -125,7 +125,7 @@ def enroll_node( ) enrollment.consumed_at = now if error is not None: - raise ComputeNodeError(error) + raise SoftwareNodeError(error) return { "node_id": str(node_id), "node_token": token, @@ -136,13 +136,13 @@ def enroll_node( def authenticate_node(node_id: UUID, token: str) -> dict: with session_scope() as session: - node = session.get(ComputeNode, node_id) + node = session.get(SoftwareNode, node_id) if ( node is None or node.status == "disabled" or not _verify_secret(token, node.token_hash) ): - raise ComputeNodeError("invalid node credentials") + raise SoftwareNodeError("invalid node credentials") return { "node_id": node.node_id, "install_id": node.install_id, @@ -152,9 +152,9 @@ def authenticate_node(node_id: UUID, token: str) -> dict: def update_node_runtime(node_id: UUID, *, status: str, runtime: dict) -> None: with session_scope() as session: - node = session.get(ComputeNode, node_id) + node = session.get(SoftwareNode, node_id) if node is None or node.status == "disabled": - raise ComputeNodeError("node is disabled or missing") + raise SoftwareNodeError("node is disabled or missing") node.status = status node.runtime = runtime node.last_seen_at = datetime.now(timezone.utc) @@ -163,14 +163,14 @@ def update_node_runtime(node_id: UUID, *, status: str, runtime: dict) -> None: def mark_node_offline(node_id: UUID) -> None: """仅把活动节点转离线;管理员禁用态不可被断线收尾覆盖。""" with session_scope() as session: - node = session.get(ComputeNode, node_id) + node = session.get(SoftwareNode, node_id) if node is not None and node.status != "disabled": node.status = "offline" def set_node_disabled(node_id: UUID, disabled: bool) -> bool: with session_scope() as session: - node = session.get(ComputeNode, node_id) + node = session.get(SoftwareNode, node_id) if node is None: return False node.status = "disabled" if disabled else "offline" @@ -180,7 +180,7 @@ def set_node_disabled(node_id: UUID, disabled: bool) -> bool: def delete_node(node_id: UUID) -> bool: """撤销并物理删除节点身份;当前节点表没有任务历史外键。""" with session_scope() as session: - node = session.get(ComputeNode, node_id) + node = session.get(SoftwareNode, node_id) if node is None: return False session.delete(node) @@ -190,7 +190,7 @@ def delete_node(node_id: UUID) -> bool: def list_nodes() -> list[dict]: with session_scope() as session: rows = ( - session.execute(select(ComputeNode).order_by(ComputeNode.created_at)) + session.execute(select(SoftwareNode).order_by(SoftwareNode.created_at)) .scalars() .all() ) diff --git a/core/storage/models.py b/core/storage/models.py index 090cbbe..487bc0e 100644 --- a/core/storage/models.py +++ b/core/storage/models.py @@ -424,10 +424,10 @@ class ChannelBinding(Base): ) -class ComputeNodeEnrollment(Base): +class SoftwareNodeEnrollment(Base): """Windows Node 一次性注册码;数据库只保存不可逆摘要。""" - __tablename__ = "compute_node_enrollments" + __tablename__ = "software_node_enrollments" enrollment_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), primary_key=True, default=uuid4) code_hash: Mapped[str] = mapped_column(Text, nullable=False, unique=True) expected_name: Mapped[Optional[str]] = mapped_column(Text, nullable=True) @@ -443,10 +443,10 @@ class ComputeNodeEnrollment(Base): ) -class ComputeNode(Base): +class SoftwareNode(Base): """平台托管的 Windows 执行节点身份与最后一次运行态。""" - __tablename__ = "compute_nodes" + __tablename__ = "software_nodes" node_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), primary_key=True, default=uuid4) name: Mapped[str] = mapped_column(Text, nullable=False) install_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), nullable=False, unique=True) @@ -465,14 +465,14 @@ class ComputeNode(Base): ) -class ComputeJob(Base): - """受控计算任务账本;请求只保存规范化参数和输入引用。""" +class SoftwareJob(Base): + """专业软件任务账本;请求只保存规范化参数和输入引用。""" - __tablename__ = "compute_jobs" + __tablename__ = "software_jobs" __table_args__ = ( - UniqueConstraint("user_id", "idempotency_key", name="uq_compute_jobs_user_idempotency"), - Index("ix_compute_jobs_status_created", "status", "created_at"), - Index("ix_compute_jobs_node_status", "node_id", "status"), + UniqueConstraint("user_id", "idempotency_key", name="uq_software_jobs_user_idempotency"), + Index("ix_software_jobs_status_created", "status", "created_at"), + Index("ix_software_jobs_node_status", "node_id", "status"), ) job_id: Mapped[UUID] = mapped_column(PG_UUID(as_uuid=True), primary_key=True, default=uuid4) @@ -488,7 +488,7 @@ class ComputeJob(Base): request_digest: Mapped[str] = mapped_column(Text, nullable=False) input_manifest: Mapped[dict[str, Any]] = mapped_column(JSONB, nullable=False, default=dict) node_id: Mapped[Optional[UUID]] = mapped_column( - PG_UUID(as_uuid=True), ForeignKey("compute_nodes.node_id", ondelete="SET NULL"), nullable=True + PG_UUID(as_uuid=True), ForeignKey("software_nodes.node_id", ondelete="SET NULL"), nullable=True ) lease_id: Mapped[Optional[UUID]] = mapped_column(PG_UUID(as_uuid=True), nullable=True) lease_expires_at: Mapped[Optional[datetime]] = mapped_column(DateTime(timezone=True), nullable=True) diff --git a/db/migrations/versions/20260813_1600_0032_compute_jobs.py b/db/migrations/versions/20260813_1600_0032_compute_jobs.py deleted file mode 100644 index 015fd65..0000000 --- a/db/migrations/versions/20260813_1600_0032_compute_jobs.py +++ /dev/null @@ -1,52 +0,0 @@ -"""Add the Windows compute job ledger. - -Revision ID: 0032 -Revises: 0031 -Create Date: 2026-08-13 -""" -from collections.abc import Sequence - -import sqlalchemy as sa -from alembic import op -from sqlalchemy.dialects import postgresql - -revision: str = "0032" -down_revision: str | None = "0031" -branch_labels: str | Sequence[str] | None = None -depends_on: str | Sequence[str] | None = None - - -def upgrade() -> None: - op.create_table( - "compute_jobs", - sa.Column("job_id", postgresql.UUID(as_uuid=True), primary_key=True), - sa.Column("user_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("users.user_id", ondelete="CASCADE"), nullable=False), - sa.Column("task_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("tasks.task_id", ondelete="CASCADE"), nullable=False), - sa.Column("idempotency_key", sa.Text(), nullable=False), - sa.Column("capability", sa.Text(), nullable=False), - sa.Column("request", postgresql.JSONB(), nullable=False), - sa.Column("request_digest", sa.Text(), nullable=False), - sa.Column("input_manifest", postgresql.JSONB(), nullable=False), - sa.Column("node_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("compute_nodes.node_id", ondelete="SET NULL"), nullable=True), - sa.Column("lease_id", postgresql.UUID(as_uuid=True), nullable=True), - sa.Column("lease_expires_at", sa.DateTime(timezone=True), nullable=True), - sa.Column("status", sa.Text(), server_default="queued", nullable=False), - sa.Column("stage", sa.Text(), server_default="", nullable=False), - sa.Column("progress", sa.Integer(), server_default="0", nullable=False), - sa.Column("metrics", postgresql.JSONB(), nullable=False), - sa.Column("error", postgresql.JSONB(), nullable=False), - sa.Column("artifact_manifest", postgresql.JSONB(), nullable=False), - sa.Column("started_at", sa.DateTime(timezone=True), nullable=True), - sa.Column("terminal_at", sa.DateTime(timezone=True), 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.UniqueConstraint("user_id", "idempotency_key", name="uq_compute_jobs_user_idempotency"), - ) - op.create_index("ix_compute_jobs_status_created", "compute_jobs", ["status", "created_at"]) - op.create_index("ix_compute_jobs_node_status", "compute_jobs", ["node_id", "status"]) - - -def downgrade() -> None: - op.drop_index("ix_compute_jobs_node_status", table_name="compute_jobs") - op.drop_index("ix_compute_jobs_status_created", table_name="compute_jobs") - op.drop_table("compute_jobs") diff --git a/db/migrations/versions/20260813_1600_0032_software_jobs.py b/db/migrations/versions/20260813_1600_0032_software_jobs.py new file mode 100644 index 0000000..4d74647 --- /dev/null +++ b/db/migrations/versions/20260813_1600_0032_software_jobs.py @@ -0,0 +1,102 @@ +"""Rename software nodes and add the professional software job ledger. + +Revision ID: 0032 +Revises: 0031 +Create Date: 2026-08-13 +""" +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +revision: str = "0032" +down_revision: str | None = "0031" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.rename_table("compute_node_enrollments", "software_node_enrollments") + op.rename_table("compute_nodes", "software_nodes") + op.execute( + "ALTER INDEX ix_compute_nodes_status RENAME TO ix_software_nodes_status" + ) + op.execute( + "ALTER TABLE software_node_enrollments RENAME CONSTRAINT " + "compute_node_enrollments_pkey TO software_node_enrollments_pkey" + ) + op.execute( + "ALTER TABLE software_node_enrollments RENAME CONSTRAINT " + "compute_node_enrollments_code_hash_key TO software_node_enrollments_code_hash_key" + ) + op.execute( + "ALTER TABLE software_node_enrollments RENAME CONSTRAINT " + "compute_node_enrollments_created_by_fkey TO software_node_enrollments_created_by_fkey" + ) + op.execute( + "ALTER TABLE software_nodes RENAME CONSTRAINT " + "compute_nodes_pkey TO software_nodes_pkey" + ) + op.execute( + "ALTER TABLE software_nodes RENAME CONSTRAINT " + "compute_nodes_install_id_key TO software_nodes_install_id_key" + ) + op.create_table( + "software_jobs", + sa.Column("job_id", postgresql.UUID(as_uuid=True), primary_key=True), + sa.Column("user_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("users.user_id", ondelete="CASCADE"), nullable=False), + sa.Column("task_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("tasks.task_id", ondelete="CASCADE"), nullable=False), + sa.Column("idempotency_key", sa.Text(), nullable=False), + sa.Column("capability", sa.Text(), nullable=False), + sa.Column("request", postgresql.JSONB(), nullable=False), + sa.Column("request_digest", sa.Text(), nullable=False), + sa.Column("input_manifest", postgresql.JSONB(), nullable=False), + sa.Column("node_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("software_nodes.node_id", ondelete="SET NULL"), nullable=True), + sa.Column("lease_id", postgresql.UUID(as_uuid=True), nullable=True), + sa.Column("lease_expires_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("status", sa.Text(), server_default="queued", nullable=False), + sa.Column("stage", sa.Text(), server_default="", nullable=False), + sa.Column("progress", sa.Integer(), server_default="0", nullable=False), + sa.Column("metrics", postgresql.JSONB(), nullable=False), + sa.Column("error", postgresql.JSONB(), nullable=False), + sa.Column("artifact_manifest", postgresql.JSONB(), nullable=False), + sa.Column("started_at", sa.DateTime(timezone=True), nullable=True), + sa.Column("terminal_at", sa.DateTime(timezone=True), 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.UniqueConstraint("user_id", "idempotency_key", name="uq_software_jobs_user_idempotency"), + ) + op.create_index("ix_software_jobs_status_created", "software_jobs", ["status", "created_at"]) + op.create_index("ix_software_jobs_node_status", "software_jobs", ["node_id", "status"]) + + +def downgrade() -> None: + op.drop_index("ix_software_jobs_node_status", table_name="software_jobs") + op.drop_index("ix_software_jobs_status_created", table_name="software_jobs") + op.drop_table("software_jobs") + op.execute( + "ALTER TABLE software_nodes RENAME CONSTRAINT " + "software_nodes_install_id_key TO compute_nodes_install_id_key" + ) + op.execute( + "ALTER TABLE software_nodes RENAME CONSTRAINT " + "software_nodes_pkey TO compute_nodes_pkey" + ) + op.execute( + "ALTER TABLE software_node_enrollments RENAME CONSTRAINT " + "software_node_enrollments_created_by_fkey TO compute_node_enrollments_created_by_fkey" + ) + op.execute( + "ALTER TABLE software_node_enrollments RENAME CONSTRAINT " + "software_node_enrollments_code_hash_key TO compute_node_enrollments_code_hash_key" + ) + op.execute( + "ALTER TABLE software_node_enrollments RENAME CONSTRAINT " + "software_node_enrollments_pkey TO compute_node_enrollments_pkey" + ) + op.execute( + "ALTER INDEX ix_software_nodes_status RENAME TO ix_compute_nodes_status" + ) + op.rename_table("software_nodes", "compute_nodes") + op.rename_table("software_node_enrollments", "compute_node_enrollments") diff --git a/deploy/nginx/zcbot.conf.example b/deploy/nginx/zcbot.conf.example index 07bc84e..0306508 100644 --- a/deploy/nginx/zcbot.conf.example +++ b/deploy/nginx/zcbot.conf.example @@ -60,7 +60,7 @@ server { # ★ Windows Node 长连接:注册走普通 POST,注册后的节点控制通道走 WebSocket。 # 必须独立于下面会清空 Connection 头的默认 location;15s 应用心跳保持链路活跃。 - location = /v1/compute/nodes/connect { + location = /v1/software-nodes/connect { proxy_pass http://zcbot_backend; proxy_http_version 1.1; proxy_set_header Upgrade $http_upgrade; diff --git a/docs/windows-node-design.md b/docs/windows-node-design.md index 5ab7184..19b88ea 100644 --- a/docs/windows-node-design.md +++ b/docs/windows-node-design.md @@ -34,7 +34,7 @@ Windows Node 不是完整的本地 zcbot: ### 2.1 目标 - 云端 zcbot 可以发现节点能力、容量、软件版本和在线状态。 -- 用户可以提交、查询、取消长时间 Windows 计算任务。 +- 用户可以提交、查询、取消长时间运行的专业软件任务。 - 网络中断、云端重启或节点重启后,任务可以确定性对账。 - 节点可以上报阶段、进度、结构化指标、日志摘要和事件截图。 - 中间产物和最终产物支持校验、断点上传与按需导入工作目录。 @@ -79,7 +79,7 @@ flowchart LR |---|---| | `NodeRegistry` | 节点注册、证书指纹、启停、能力和管理员标签 | | `NodeConnectionManager` | WSS 连接、心跳、消息 ACK、同节点单活连接 | -| `ComputeJobService` | 用户授权、幂等提交、节点选择、租约、取消和终态 | +| `SoftwareJobService` | 用户授权、幂等提交、节点选择、租约、取消和终态 | | `ComputeTransferService` | 输入下载凭证、分块上传、SHA-256、容量与保留期 | | `ComputeBroker` | 把任务事件推送到 Web UI;不承担持久化事实源 | | `ComputeTools` | agent 可调用的能力发现、提交、查询、取消、产物导入工具 | @@ -157,7 +157,7 @@ sequenceDiagram participant Node as Windows Node Admin->>Cloud: 创建一次性 enrollment token Node->>Node: 生成设备密钥对 - Node->>Cloud: POST /v1/compute/nodes/enroll + Node->>Cloud: POST /v1/software-nodes/enroll Cloud->>Cloud: 消耗 token,创建 node_id Cloud-->>Node: 客户端证书、CA、云端地址 Node->>Node: 私钥写入 Windows Certificate Store @@ -178,7 +178,7 @@ sequenceDiagram 节点连接: ```text -WSS /v1/compute/nodes/connect +WSS /v1/software-nodes/connect ``` 统一消息 envelope: @@ -252,13 +252,13 @@ Node 在 `hello` 和心跳中声明由本机可信配置生成的能力: 建议新增三张表,不复用外部系统连接表: ```text -compute_nodes( +software_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( +software_jobs( job_id pk, user_id fk, task_id fk, tool_call_id, capability, schema_version, request jsonb, idempotency_key, request_digest, @@ -268,7 +268,7 @@ compute_jobs( created_at, started_at, terminal_at, updated_at ) -compute_job_events( +software_job_events( event_id pk, job_id fk, sequence, kind, level, payload jsonb, created_at ) @@ -352,7 +352,7 @@ Node 不访问整个用户 workspace。云端只为显式引用的文件创建 大产物不通过 WSS 消息传输。Node 使用 HTTPS 分块上传: ```text -POST /v1/compute/jobs/{job_id}/artifacts/upload-session +POST /v1/software-jobs/{job_id}/artifacts/upload-session PUT /v1/compute/transfers/{transfer_id}/parts/{part_number} POST /v1/compute/transfers/{transfer_id}/complete ``` @@ -364,7 +364,7 @@ POST /v1/compute/transfers/{transfer_id}/complete Node 上传完成后先进入: ```text -/.zcbot_cache//compute_jobs// +/.zcbot_cache//software_jobs// ``` 该目录默认隐藏且有 TTL。用户或 agent 明确导入后复制到: @@ -379,10 +379,10 @@ Node 上传完成后先进入: ```text compute_capability_list -compute_job_submit -compute_job_status -compute_job_cancel -compute_job_artifact_import +software_job_submit +software_job_status +software_job_cancel +software_job_artifact_import ``` - capability list 只返回用户有权使用且有健康节点承载的能力; @@ -418,7 +418,7 @@ process.cancel_adapter - 优先捕获目标窗口; - 敏感信息上传前遮罩; - 视频默认关闭,显式启用时建议 1280×720、5–10 FPS、H.264 分段; -- 采集失败不得使计算任务失败; +- 采集失败不得使专业软件任务失败; - 模型只按需读取关键帧,不持续消费完整视频。 ## 12. Origin 首批适配器 diff --git a/docs/windows-node-mvp-intranet.md b/docs/windows-node-mvp-intranet.md index 02a9137..51e7b7e 100644 --- a/docs/windows-node-mvp-intranet.md +++ b/docs/windows-node-mvp-intranet.md @@ -60,7 +60,7 @@ sequenceDiagram participant N as Windows Node A->>Z: 创建一次性注册码 A->>N: 输入内网地址和注册码 - N->>Z: HTTP POST /v1/compute/nodes/enroll + N->>Z: HTTP POST /v1/software-nodes/enroll Z->>Z: 校验并原子消费注册码 Z-->>N: node_id + node_token + 配置 N->>N: DPAPI 加密保存 node_token @@ -71,7 +71,7 @@ sequenceDiagram 注册请求: ```http -POST http://zcbot.internal:8765/v1/compute/nodes/enroll +POST http://zcbot.internal:8765/v1/software-nodes/enroll Content-Type: application/json ``` @@ -116,7 +116,7 @@ Content-Type: application/json ### 4.1 WS 连接 ```http -GET ws://zcbot.internal:8765/v1/compute/nodes/connect +GET ws://zcbot.internal:8765/v1/software-nodes/connect Authorization: Bearer X-Node-Id: Upgrade: websocket @@ -170,13 +170,13 @@ RDP:不向公网开放,使用 VPN、堡垒机或云安全登录 云端首期只增加: ```text -compute_nodes( +software_nodes( node_id pk, name, install_id, token_hash, status, capabilities jsonb, last_seen_at, created_at, updated_at ) -compute_jobs( +software_jobs( job_id pk, user_id fk, task_id fk, idempotency_key, capability, request jsonb, node_id fk, status, progress, diff --git a/tests/test_compute_nodes.py b/tests/test_software_nodes.py similarity index 89% rename from tests/test_compute_nodes.py rename to tests/test_software_nodes.py index 0a41d4e..9b95387 100644 --- a/tests/test_compute_nodes.py +++ b/tests/test_software_nodes.py @@ -11,13 +11,13 @@ from alembic.operations import Operations from sqlalchemy import create_mock_engine from sqlalchemy.dialects import postgresql -from core.compute_nodes import ( +from core.software_nodes import ( _enrollment_digest, _hash_secret, _verify_secret, delete_node, ) -from core.compute_jobs import ( +from core.software_jobs import ( _canonical_request, abandon_offer, mark_node_jobs_disconnected, @@ -26,10 +26,10 @@ from core.compute_jobs import ( update_job_state, validate_output_manifest, ) -from web.routers.compute_nodes import NodeConnectionManager, _bearer +from web.routers.software_nodes import NodeConnectionManager, _bearer -class ComputeNodeSecurityTests(unittest.TestCase): +class SoftwareNodeSecurityTests(unittest.TestCase): def test_secret_hash_is_salted_and_verifiable(self) -> None: first = _hash_secret("node-secret") second = _hash_secret("node-secret") @@ -50,9 +50,9 @@ class ComputeNodeSecurityTests(unittest.TestCase): def test_websocket_auth_rejection_uses_explicit_application_close_code(self) -> None: source = ( - Path(__file__).resolve().parents[1] / "web" / "routers" / "compute_nodes.py" + Path(__file__).resolve().parents[1] / "web" / "routers" / "software_nodes.py" ).read_text(encoding="utf-8") - rejection = source.split("except (ValueError, ComputeNodeError):", 1)[1].split( + rejection = source.split("except (ValueError, SoftwareNodeError):", 1)[1].split( "await node_connections.activate", 1 )[0] self.assertLess( @@ -61,7 +61,7 @@ class ComputeNodeSecurityTests(unittest.TestCase): self.assertIn('code=4003, reason="invalid node credentials"', rejection) -class ComputeNodeConnectionTests(unittest.IsolatedAsyncioTestCase): +class SoftwareNodeConnectionTests(unittest.IsolatedAsyncioTestCase): async def test_new_connection_replaces_old_without_removing_new(self) -> None: manager = NodeConnectionManager() node_id = uuid4() @@ -87,7 +87,7 @@ class ComputeNodeConnectionTests(unittest.IsolatedAsyncioTestCase): self.assertFalse(await manager.remove(node_id, websocket)) -class ComputeNodeMigrationTests(unittest.TestCase): +class SoftwareNodeMigrationTests(unittest.TestCase): def test_0030_upgrade_compiles_as_postgresql_ddl(self) -> None: statements: list[str] = [] @@ -116,18 +116,20 @@ class ComputeNodeMigrationTests(unittest.TestCase): engine = create_mock_engine("postgresql+psycopg://", capture) operations = Operations(MigrationContext.configure(engine.connect())) migration = importlib.import_module( - "db.migrations.versions.20260813_1600_0032_compute_jobs" + "db.migrations.versions.20260813_1600_0032_software_jobs" ) with patch.object(migration, "op", operations): migration.upgrade() rendered = "\n".join(statements) - self.assertIn("compute_jobs", rendered) - self.assertIn("uq_compute_jobs_user_idempotency", rendered) - self.assertIn("ix_compute_jobs_status_created", rendered) + self.assertIn("ALTER TABLE compute_node_enrollments RENAME TO software_node_enrollments", rendered) + self.assertIn("ALTER TABLE compute_nodes RENAME TO software_nodes", rendered) + self.assertIn("software_jobs", rendered) + self.assertIn("uq_software_jobs_user_idempotency", rendered) + self.assertIn("ix_software_jobs_status_created", rendered) -class ComputeJobProtocolTests(unittest.TestCase): +class SoftwareJobProtocolTests(unittest.TestCase): def test_origin_request_is_canonical_and_rejects_extra_fields(self) -> None: request = { "schema_version": 1, @@ -182,7 +184,7 @@ class ComputeJobProtocolTests(unittest.TestCase): [{**manifest[0], "filename": "anything.opju"}, *manifest[1:]], ) - @patch("core.compute_jobs.session_scope") + @patch("core.software_jobs.session_scope") def test_stale_offer_cannot_be_accepted_by_another_node(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value job = type("Job", (), {})() @@ -197,7 +199,7 @@ class ComputeJobProtocolTests(unittest.TestCase): payload={"job_id": str(uuid4()), "lease_id": str(job.lease_id)}, ) - @patch("core.compute_jobs.session_scope") + @patch("core.software_jobs.session_scope") def test_failed_delivery_only_abandons_matching_offer(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value node_id = uuid4() @@ -216,7 +218,7 @@ class ComputeJobProtocolTests(unittest.TestCase): def test_dispatcher_excludes_nodes_with_active_jobs(self) -> None: source = ( - Path(__file__).resolve().parents[1] / "core" / "compute_jobs.py" + Path(__file__).resolve().parents[1] / "core" / "software_jobs.py" ).read_text(encoding="utf-8") self.assertIn('{"offered", "dispatched", "running"}', source) self.assertIn("item.node_id not in busy_node_ids", source) @@ -224,12 +226,12 @@ class ComputeJobProtocolTests(unittest.TestCase): def test_input_download_rechecks_file_digest(self) -> None: source = ( Path(__file__).resolve().parents[1] - / "web" / "routers" / "compute_nodes.py" + / "web" / "routers" / "software_nodes.py" ).read_text(encoding="utf-8") self.assertIn("digest = sha256()", source) self.assertIn('digest.hexdigest() != item["sha256"]', source) - @patch("core.compute_jobs.session_scope") + @patch("core.software_jobs.session_scope") def test_job_state_restores_disconnected_job(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value node_id = uuid4() @@ -253,7 +255,7 @@ class ComputeJobProtocolTests(unittest.TestCase): self.assertEqual(job.status, "dispatched") self.assertEqual(job.stage, "waiting_input") - @patch("core.compute_jobs.session_scope") + @patch("core.software_jobs.session_scope") def test_ready_to_run_is_not_reported_as_running(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value node_id = uuid4() @@ -273,7 +275,7 @@ class ComputeJobProtocolTests(unittest.TestCase): }) self.assertEqual(job.status, "dispatched") - @patch("core.compute_jobs.session_scope") + @patch("core.software_jobs.session_scope") def test_terminal_replay_is_idempotent(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value node_id = uuid4() @@ -295,7 +297,7 @@ class ComputeJobProtocolTests(unittest.TestCase): }) self.assertEqual(job.status, "failed") - @patch("core.compute_jobs.session_scope") + @patch("core.software_jobs.session_scope") def test_disconnect_does_not_requeue_active_jobs(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value first = type("Job", (), {"status": "running"})() @@ -306,8 +308,8 @@ class ComputeJobProtocolTests(unittest.TestCase): self.assertEqual(second.status, "disconnected") -class ComputeNodeDeleteTests(unittest.TestCase): - @patch("core.compute_nodes.session_scope") +class SoftwareNodeDeleteTests(unittest.TestCase): + @patch("core.software_nodes.session_scope") def test_delete_node_removes_existing_identity(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value node = object() @@ -316,7 +318,7 @@ class ComputeNodeDeleteTests(unittest.TestCase): self.assertTrue(delete_node(uuid4())) session.delete.assert_called_once_with(node) - @patch("core.compute_nodes.session_scope") + @patch("core.software_nodes.session_scope") def test_delete_node_reports_missing_identity(self, session_scope) -> None: session = session_scope.return_value.__enter__.return_value session.get.return_value = None diff --git a/tests/test_compute_output_publish.py b/tests/test_software_output_publish.py similarity index 78% rename from tests/test_compute_output_publish.py rename to tests/test_software_output_publish.py index 43b0d34..f192005 100644 --- a/tests/test_compute_output_publish.py +++ b/tests/test_software_output_publish.py @@ -7,16 +7,16 @@ from pathlib import Path from unittest.mock import patch from uuid import uuid4 -from web.routers.compute_nodes import _publish_compute_outputs +from web.routers.software_nodes import _publish_software_job_outputs -class ComputeOutputPublishTests(unittest.TestCase): +class SoftwareOutputPublishTests(unittest.TestCase): def test_complete_set_moves_atomically_and_can_be_replayed(self) -> None: with tempfile.TemporaryDirectory() as directory: root = Path(directory) job_id = uuid4() working_dir = root / "research" - staging = root / ".zcbot_compute_staging" / str(job_id) + staging = root / ".zcbot_software_job_staging" / str(job_id) staging.mkdir(parents=True) working_dir.mkdir() content = b"origin-result" @@ -44,14 +44,14 @@ class ComputeOutputPublishTests(unittest.TestCase): } for ref in kwargs["refs"]) with ( - patch("web.routers.compute_nodes.load_user_root", return_value=root), + patch("web.routers.software_nodes.load_user_root", return_value=root), patch( - "web.routers.compute_nodes.register_published_artifacts", + "web.routers.software_nodes.register_published_artifacts", side_effect=register, ), ): - first = _publish_compute_outputs(job_id, context, manifest) - second = _publish_compute_outputs(job_id, context, manifest) + first = _publish_software_job_outputs(job_id, context, manifest) + second = _publish_software_job_outputs(job_id, context, manifest) published = working_dir / "origin" / str(job_id) / "figure.png" self.assertEqual(published.read_bytes(), content) diff --git a/tests/test_static_vendor.py b/tests/test_static_vendor.py index 77dc763..22cfea1 100644 --- a/tests/test_static_vendor.py +++ b/tests/test_static_vendor.py @@ -32,14 +32,14 @@ class StaticVendorTests(unittest.TestCase): self.assertIn('id="node-enrollment-modal" class="modal"', html) self.assertIn("生成 Windows Node 注册码", html) - self.assertIn('"/v1/admin/compute-node-enrollments"', admin_js) + self.assertIn('"/v1/admin/software-node-enrollments"', admin_js) self.assertIn('capabilities: ["origin.plot@v1"]', admin_js) self.assertIn('origin.health === "ready"', admin_js) self.assertIn("ttl_seconds: 600", admin_js) self.assertIn("navigator.clipboard.writeText(value)", admin_js) - self.assertIn('apiGet("/v1/admin/compute-nodes")', admin_js) - self.assertIn('apiSend("PATCH", `/v1/admin/compute-nodes/${node.node_id}`', admin_js) - self.assertIn('apiSend("DELETE", `/v1/admin/compute-nodes/${node.node_id}`', admin_js) + self.assertIn('apiGet("/v1/admin/software-nodes")', admin_js) + self.assertIn('apiSend("PATCH", `/v1/admin/software-nodes/${node.node_id}`', admin_js) + self.assertIn('apiSend("DELETE", `/v1/admin/software-nodes/${node.node_id}`', admin_js) self.assertIn("最近心跳", admin_js) self.assertIn("重新启用", admin_js) self.assertIn("永久删除", admin_js) diff --git a/tests/test_web_routes_nodb.py b/tests/test_web_routes_nodb.py index 21a9820..f344b15 100644 --- a/tests/test_web_routes_nodb.py +++ b/tests/test_web_routes_nodb.py @@ -110,8 +110,8 @@ class AuthGateTests(unittest.TestCase): ("POST", "/v1/tasks"), ("POST", "/v1/asr/transcribe"), ("GET", "/v1/admin/overview"), - ("GET", "/v1/admin/compute-nodes"), - ("DELETE", "/v1/admin/compute-nodes/00000000-0000-0000-0000-000000000000"), + ("GET", "/v1/admin/software-nodes"), + ("DELETE", "/v1/admin/software-nodes/00000000-0000-0000-0000-000000000000"), ("GET", "/v1/admin/tool-wire-health"), ("GET", "/v1/admin/external-system-definitions"), ("GET", "/v1/admin/external-system-users"), diff --git a/tests/test_windows_node_source.py b/tests/test_windows_node_source.py index e019a12..e5a4c89 100644 --- a/tests/test_windows_node_source.py +++ b/tests/test_windows_node_source.py @@ -20,8 +20,8 @@ class WindowsNodeSourceTests(unittest.TestCase): def test_node_protocol_and_secret_storage_markers_are_present(self) -> None: source = "\n".join(path.read_text(encoding="utf-8") for path in PROJECT.glob("*.cs")) for marker in ( - "v1/compute/nodes/enroll", - "v1/compute/nodes/connect", + "v1/software-nodes/enroll", + "v1/software-nodes/connect", 'SetRequestHeader("Authorization"', 'SetRequestHeader("X-Node-Id"', "DataProtectionScope.LocalMachine", diff --git a/web/app.py b/web/app.py index f788b37..a2107db 100644 --- a/web/app.py +++ b/web/app.py @@ -49,7 +49,7 @@ from .background import ( from .broker import broker from .routers.asr import register_asr_routes from .routers.authroutes import register_auth_routes -from .routers.compute_nodes import register_compute_node_routes +from .routers.software_nodes import register_software_node_routes from .routers.external_systems import register_external_system_routes from .routers.files import register_file_routes from .routers.kb import register_kb_routes @@ -206,7 +206,7 @@ def create_app() -> FastAPI: register_asr_routes(app, require_user=require_user, auth_cfg=auth_cfg) register_task_routes(app, require_user=require_user) register_message_routes(app, require_user=require_user) - register_compute_node_routes( + register_software_node_routes( app, require_user=require_user, require_admin=require_admin ) diff --git a/web/routers/compute_nodes.py b/web/routers/software_nodes.py similarity index 84% rename from web/routers/compute_nodes.py rename to web/routers/software_nodes.py index 74e94f1..24af62a 100644 --- a/web/routers/compute_nodes.py +++ b/web/routers/software_nodes.py @@ -11,8 +11,8 @@ from uuid import UUID from fastapi import Depends, Header, HTTPException, Request, WebSocket, WebSocketDisconnect, status from fastapi.responses import FileResponse -from core.compute_nodes import ( - ComputeNodeError, +from core.software_nodes import ( + SoftwareNodeError, authenticate_node, create_enrollment, delete_node, @@ -22,7 +22,7 @@ from core.compute_nodes import ( set_node_disabled, update_node_runtime, ) -from core.compute_jobs import ( +from core.software_jobs import ( MAX_OUTPUT_ARTIFACT_BYTES, MAX_OUTPUT_TOTAL_BYTES, OUTPUT_ARTIFACTS, @@ -37,13 +37,14 @@ from core.compute_jobs import ( respond_to_offer, update_job_state, validate_output_manifest, + SoftwareJobError, ) from core.artifact_lifecycle import register_published_artifacts from web.schemas import ( - ComputeEnrollmentCreateRequest, - ComputeJobCreateRequest, - ComputeNodeDisableRequest, - ComputeNodeEnrollRequest, + SoftwareEnrollmentCreateRequest, + SoftwareJobCreateRequest, + SoftwareNodeDisableRequest, + SoftwareNodeEnrollRequest, ) from web.userfiles import load_user_root, safe_join @@ -108,7 +109,7 @@ node_connections = NodeConnectionManager() def _bearer(authorization: str | None) -> str: scheme, _, token = (authorization or "").partition(" ") if scheme.lower() != "bearer" or not token: - raise ComputeNodeError("missing node bearer token") + raise SoftwareNodeError("missing node bearer token") return token @@ -123,11 +124,11 @@ def _authenticate_output_request( node_id = UUID(x_node_id) lease_id = UUID(x_lease_id) authenticate_node(node_id, _bearer(authorization)) - except (ValueError, ComputeNodeError) as exc: + except (ValueError, SoftwareNodeError) as exc: raise HTTPException(401, "invalid node credentials or job identity") from exc context = get_job_output_context(node_id, job_id, lease_id, x_request_digest) if context is None: - raise HTTPException(404, "compute job output target not found") + raise HTTPException(404, "software job output target not found") return node_id, lease_id, context @@ -145,13 +146,13 @@ def _reject_symlink_path(root: Path, target: Path) -> None: for part in target.relative_to(root).parts: current = current / part if current.is_symlink(): - raise HTTPException(409, "compute output path contains a symbolic link") + raise HTTPException(409, "software job output path contains a symbolic link") -def _publish_compute_outputs(job_id: UUID, context: dict, manifest: list[dict]) -> list[dict]: +def _publish_software_job_outputs(job_id: UUID, context: dict, manifest: list[dict]) -> list[dict]: root = load_user_root(context["user_id"]) working_dir = safe_join(root, context["working_dir"]) - staging = safe_join(root, f".zcbot_compute_staging/{job_id}") + staging = safe_join(root, f".zcbot_software_job_staging/{job_id}") relative_output = Path("origin") / str(job_id) destination = safe_join(working_dir, relative_output.as_posix()) source = staging if staging.is_dir() else destination @@ -164,11 +165,11 @@ def _publish_compute_outputs(job_id: UUID, context: dict, manifest: list[dict]) or path.stat().st_size != item["size_bytes"] or _hash_file(path) != item["sha256"] ): - raise ComputeNodeError(f"uploaded artifact is missing or invalid: {item['artifact_id']}") + raise SoftwareJobError(f"uploaded artifact is missing or invalid: {item['artifact_id']}") if source == staging: destination.parent.mkdir(parents=True, exist_ok=True) if destination.exists(): - raise ComputeNodeError("compute output destination already exists unexpectedly") + raise SoftwareJobError("software job output destination already exists unexpectedly") os.replace(staging, destination) try: staging.parent.rmdir() @@ -198,20 +199,20 @@ def _publish_compute_outputs(job_id: UUID, context: dict, manifest: list[dict]) ] -def register_compute_node_routes(app, *, require_user, require_admin) -> None: +def register_software_node_routes(app, *, require_user, require_admin) -> None: @app.post( - "/v1/compute/nodes/enroll", - tags=["compute-nodes"], + "/v1/software-nodes/enroll", + tags=["software-nodes"], status_code=status.HTTP_201_CREATED, ) - def node_enroll(body: ComputeNodeEnrollRequest): + def node_enroll(body: SoftwareNodeEnrollRequest): try: return enroll_node(**body.model_dump()) - except ComputeNodeError as exc: + except SoftwareNodeError as exc: raise HTTPException(400, str(exc)) from exc - @app.get("/v1/compute/jobs/{job_id}/input", tags=["compute-nodes"]) - def download_compute_job_input( + @app.get("/v1/software-jobs/{job_id}/input", tags=["software-nodes"]) + def download_software_job_input( job_id: UUID, authorization: str | None = Header(default=None), x_node_id: str = Header(default=""), @@ -219,23 +220,23 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: try: node_id = UUID(x_node_id) authenticate_node(node_id, _bearer(authorization)) - except (ValueError, ComputeNodeError) as exc: + except (ValueError, SoftwareNodeError) as exc: raise HTTPException(401, "invalid node credentials") from exc item = get_job_input(node_id, job_id) if item is None: - raise HTTPException(404, "compute job input not found") + raise HTTPException(404, "software job input not found") target = safe_join(load_user_root(item["user_id"]), item["current_path"]) if not target.is_file(): - raise HTTPException(404, "compute job input file not found") + raise HTTPException(404, "software job input file not found") stat = target.stat() if stat.st_size != item["size_bytes"]: - raise HTTPException(409, "compute job input changed after submission") + raise HTTPException(409, "software job input changed after submission") digest = sha256() with target.open("rb") as handle: for chunk in iter(lambda: handle.read(1024 * 1024), b""): digest.update(chunk) if digest.hexdigest() != item["sha256"]: - raise HTTPException(409, "compute job input changed after submission") + raise HTTPException(409, "software job input changed after submission") return FileResponse( path=str(target), filename=item["filename"], @@ -247,11 +248,11 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: ) @app.put( - "/v1/compute/jobs/{job_id}/outputs/{artifact_id}", - tags=["compute-nodes"], + "/v1/software-jobs/{job_id}/outputs/{artifact_id}", + tags=["software-nodes"], status_code=status.HTTP_204_NO_CONTENT, ) - async def upload_compute_job_output( + async def upload_software_job_output( job_id: UUID, artifact_id: str, request: Request, @@ -286,7 +287,7 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: if published.stat().st_size == x_content_length and _hash_file(published) == x_content_sha256: return None raise HTTPException(409, "published output conflicts with uploaded artifact") - staging = safe_join(root, f".zcbot_compute_staging/{job_id}") + staging = safe_join(root, f".zcbot_software_job_staging/{job_id}") _reject_symlink_path(root, staging) staging.mkdir(parents=True, exist_ok=True) destination = staging / filename @@ -298,7 +299,7 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: item.stat().st_size for item in staging.iterdir() if item.is_file() ) if staged_total + x_content_length > MAX_OUTPUT_TOTAL_BYTES: - raise HTTPException(413, "compute job outputs exceed the total size limit") + raise HTTPException(413, "software job outputs exceed the total size limit") temporary = destination.with_name(destination.name + ".tmp-" + os.urandom(8).hex()) digest = sha256() total = 0 @@ -319,8 +320,8 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: temporary.unlink(missing_ok=True) return None - @app.post("/v1/compute/jobs/{job_id}/outputs/complete", tags=["compute-nodes"]) - async def complete_compute_job_outputs( + @app.post("/v1/software-jobs/{job_id}/outputs/complete", tags=["software-nodes"]) + async def complete_software_job_outputs( job_id: UUID, request: Request, authorization: str | None = Header(default=None), @@ -337,7 +338,9 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: raise HTTPException(400, "output completion body must be an object") try: manifest = validate_output_manifest(context["request"], body.get("artifact_manifest")) - published = await asyncio.to_thread(_publish_compute_outputs, job_id, context, manifest) + published = await asyncio.to_thread( + _publish_software_job_outputs, job_id, context, manifest + ) terminal = { "job_id": str(job_id), "lease_id": str(lease_id), @@ -347,17 +350,17 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: "artifact_manifest": published, } await asyncio.to_thread(record_job_terminal, node_id, terminal) - except (ComputeNodeError, KeyError, TypeError) as exc: + except (SoftwareJobError, KeyError, TypeError) as exc: raise HTTPException(409, str(exc)) from exc return {"status": "succeeded", "artifact_manifest": published} - @app.websocket("/v1/compute/nodes/connect") + @app.websocket("/v1/software-nodes/connect") async def node_connect(websocket: WebSocket): try: node_id = UUID(websocket.headers.get("x-node-id", "")) token = _bearer(websocket.headers.get("authorization")) identity = await asyncio.to_thread(authenticate_node, node_id, token) - except (ValueError, ComputeNodeError): + except (ValueError, SoftwareNodeError): # 握手前 close 会被 ASGI 统一表现为 HTTP 403,客户端无法区分 # “凭据无效”和“代理/路由没有正确转发 WebSocket”。先升级再用 # 应用关闭码给已持有 Node ID/Token 的节点返回明确诊断。 @@ -443,17 +446,17 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: await asyncio.to_thread( abandon_offer, offer["node_id"], offer["payload"] ) - except (ComputeNodeError, WebSocketDisconnect, RuntimeError, ValueError): + except (SoftwareNodeError, WebSocketDisconnect, RuntimeError, ValueError): pass finally: if await node_connections.remove(node_id, websocket): await asyncio.to_thread(mark_node_offline, node_id) await asyncio.to_thread(mark_node_jobs_disconnected, node_id) - @app.post("/v1/tasks/{task_id}/compute-jobs", tags=["compute-jobs"]) - async def submit_compute_job( + @app.post("/v1/tasks/{task_id}/software-jobs", tags=["software-jobs"]) + async def submit_software_job( task_id: UUID, - body: ComputeJobCreateRequest, + body: SoftwareJobCreateRequest, user_id: UUID = Depends(require_user), # noqa: B008 ): try: @@ -472,42 +475,42 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: abandon_offer, offer["node_id"], offer["payload"] ) return {**job, "created": created} - except ComputeNodeError as exc: + except SoftwareJobError as exc: detail = str(exc) raise HTTPException(404 if detail == "task not found" else 400, detail) from exc - @app.get("/v1/compute-jobs/{job_id}", tags=["compute-jobs"]) - def read_compute_job( + @app.get("/v1/software-jobs/{job_id}", tags=["software-jobs"]) + def read_software_job( job_id: UUID, user_id: UUID = Depends(require_user), # noqa: B008 ): job = get_job(user_id, job_id) if job is None: - raise HTTPException(404, "compute job not found") + raise HTTPException(404, "software job not found") return job - @app.post("/v1/admin/compute-node-enrollments", tags=["admin"]) - def admin_create_compute_enrollment( - body: ComputeEnrollmentCreateRequest, + @app.post("/v1/admin/software-node-enrollments", tags=["admin"]) + def admin_create_software_enrollment( + body: SoftwareEnrollmentCreateRequest, user_id: UUID = Depends(require_admin), # noqa: B008 ): try: return create_enrollment(user_id, **body.model_dump()) - except ComputeNodeError as exc: + except SoftwareNodeError as exc: raise HTTPException(400, str(exc)) from exc - @app.get("/v1/admin/compute-nodes", tags=["admin"]) - def admin_compute_nodes(user_id: UUID = Depends(require_admin)): # noqa: B008 + @app.get("/v1/admin/software-nodes", tags=["admin"]) + def admin_software_nodes(user_id: UUID = Depends(require_admin)): # noqa: B008 return {"results": list_nodes()} - @app.patch("/v1/admin/compute-nodes/{node_id}", tags=["admin"]) - async def admin_disable_compute_node( + @app.patch("/v1/admin/software-nodes/{node_id}", tags=["admin"]) + async def admin_disable_software_node( node_id: UUID, - body: ComputeNodeDisableRequest, + body: SoftwareNodeDisableRequest, user_id: UUID = Depends(require_admin), # noqa: B008 ): if not await asyncio.to_thread(set_node_disabled, node_id, body.disabled): - raise HTTPException(404, "compute node not found") + raise HTTPException(404, "software node not found") if body.disabled: await node_connections.close(node_id) return { @@ -515,13 +518,13 @@ def register_compute_node_routes(app, *, require_user, require_admin) -> None: "status": "disabled" if body.disabled else "offline", } - @app.delete("/v1/admin/compute-nodes/{node_id}", tags=["admin"]) - async def admin_delete_compute_node( + @app.delete("/v1/admin/software-nodes/{node_id}", tags=["admin"]) + async def admin_delete_software_node( node_id: UUID, user_id: UUID = Depends(require_admin), # noqa: B008 ): # 先撤掉在线连接,避免删除后的旧 socket 继续上报运行态。 await node_connections.close(node_id) if not await asyncio.to_thread(delete_node, node_id): - raise HTTPException(404, "compute node not found") + raise HTTPException(404, "software node not found") return {"node_id": str(node_id), "status": "deleted"} diff --git a/web/schemas.py b/web/schemas.py index 5603bdf..6d009f4 100644 --- a/web/schemas.py +++ b/web/schemas.py @@ -116,13 +116,13 @@ class ExternalSystemCredentialsRequest(BaseModel): credentials: dict[str, str] = Field(default_factory=dict) -class ComputeEnrollmentCreateRequest(BaseModel): +class SoftwareEnrollmentCreateRequest(BaseModel): expected_name: str = "" capabilities: list[str] = Field(default_factory=lambda: ["origin.plot@v1"]) ttl_seconds: int = 600 -class ComputeNodeEnrollRequest(BaseModel): +class SoftwareNodeEnrollRequest(BaseModel): enrollment_code: str node_name: str install_id: UUID @@ -131,11 +131,11 @@ class ComputeNodeEnrollRequest(BaseModel): capabilities: list[str] -class ComputeNodeDisableRequest(BaseModel): +class SoftwareNodeDisableRequest(BaseModel): disabled: bool = True -class ComputeJobCreateRequest(BaseModel): +class SoftwareJobCreateRequest(BaseModel): idempotency_key: str capability: str = "origin.plot@v1" request: dict = Field(default_factory=dict) diff --git a/web/static/js/admin.js b/web/static/js/admin.js index d80b3ef..7ac2421 100644 --- a/web/static/js/admin.js +++ b/web/static/js/admin.js @@ -55,7 +55,7 @@ let externalDefinitions = []; let externalUsers = []; let externalDefinitionsLoaded = false; let externalEditingId = ""; -let computeNodes = []; +let softwareNodes = []; // ───── 格式化 ───── function fmtCNY(n) { @@ -167,7 +167,7 @@ function nodeStatusHTML(status) { } function renderWindowsNodes() { - const rows = computeNodes.map(node => { + const rows = softwareNodes.map(node => { const runtime = node.runtime || {}; const origin = runtime.origin || {}; const originState = origin.health === "ready" ? "Origin 可用" : "Origin 不可用"; @@ -197,22 +197,22 @@ function renderWindowsNodes() { }).join("") || `尚无已注册的 Windows Node`; $("s-windows-node").innerHTML = `
` - + `

Windows Node(${computeNodes.length})

查看节点状态;禁用会立即断开节点并拒绝后续连接。
` + + `

Windows Node(${softwareNodes.length})

查看节点状态;禁用会立即断开节点并拒绝后续连接。
` + `` + `
` + `${rows}
节点状态运行环境版本能力最近心跳操作
`; $("node-enrollment-open").onclick = openNodeEnrollmentModal; $("s-windows-node").querySelectorAll("[data-node-toggle]").forEach(button => { - button.onclick = () => toggleComputeNode(button); + button.onclick = () => toggleSoftwareNode(button); }); $("s-windows-node").querySelectorAll("[data-node-delete]").forEach(button => { - button.onclick = () => deleteComputeNode(button); + button.onclick = () => deleteSoftwareNode(button); }); } -async function toggleComputeNode(button) { +async function toggleSoftwareNode(button) { const row = button.closest("tr[data-node-id]"); - const node = computeNodes.find(item => item.node_id === row?.dataset.nodeId); + const node = softwareNodes.find(item => item.node_id === row?.dataset.nodeId); if (!node) return; const disabling = button.dataset.nodeToggle === "disable"; const confirmed = await dialogConfirm({ @@ -226,20 +226,20 @@ async function toggleComputeNode(button) { if (!confirmed) return; button.disabled = true; try { - await apiSend("PATCH", `/v1/admin/compute-nodes/${node.node_id}`, { + await apiSend("PATCH", `/v1/admin/software-nodes/${node.node_id}`, { disabled: disabling, }); message(disabling ? "节点已禁用" : "节点已重新启用,请在节点电脑上立即重连", "success", 5000); - await loadComputeNodes(); + await loadSoftwareNodes(); } catch (err) { if (err.code !== "auth") message("更新节点失败:" + (err.message || String(err)), "error", 5000); button.disabled = false; } } -async function deleteComputeNode(button) { +async function deleteSoftwareNode(button) { const row = button.closest("tr[data-node-id]"); - const node = computeNodes.find(item => item.node_id === row?.dataset.nodeId); + const node = softwareNodes.find(item => item.node_id === row?.dataset.nodeId); if (!node) return; const confirmed = await dialogConfirm({ title: "删除 Windows Node", @@ -250,9 +250,9 @@ async function deleteComputeNode(button) { if (!confirmed) return; button.disabled = true; try { - await apiSend("DELETE", `/v1/admin/compute-nodes/${node.node_id}`, {}); + await apiSend("DELETE", `/v1/admin/software-nodes/${node.node_id}`, {}); message("节点已删除,本机需重新注册后才能使用", "success", 5000); - await loadComputeNodes(); + await loadSoftwareNodes(); } catch (err) { if (err.code !== "auth") message("删除节点失败:" + (err.message || String(err)), "error", 5000); button.disabled = false; @@ -310,7 +310,7 @@ async function createNodeEnrollment(e) { submit.disabled = true; submit.textContent = "生成中…"; try { - const result = await apiSend("POST", "/v1/admin/compute-node-enrollments", { + const result = await apiSend("POST", "/v1/admin/software-node-enrollments", { expected_name: $("node-expected-name").value.trim(), capabilities: ["origin.plot@v1"], ttl_seconds: 600, @@ -1035,10 +1035,10 @@ async function loadExternalDefinitions(force = false) { } catch (e) { /* overview 统一处理鉴权 */ } } -async function loadComputeNodes() { +async function loadSoftwareNodes() { try { - const result = await apiGet("/v1/admin/compute-nodes"); - computeNodes = result.results || []; + const result = await apiGet("/v1/admin/software-nodes"); + softwareNodes = result.results || []; renderWindowsNodes(); } catch (e) { /* overview 统一处理鉴权 */ } } @@ -1052,7 +1052,7 @@ async function refresh() { loadModels(); loadUserUsage(userPage); loadStorage(storagePage); - loadComputeNodes(); + loadSoftwareNodes(); loadExternalDefinitions(); loadToolFailures(); } catch (e) { diff --git a/windows-node/Zcbot.WindowsNode/EnrollOptions.cs b/windows-node/Zcbot.WindowsNode/EnrollOptions.cs index 1bdfdb0..7cb08dc 100644 --- a/windows-node/Zcbot.WindowsNode/EnrollOptions.cs +++ b/windows-node/Zcbot.WindowsNode/EnrollOptions.cs @@ -52,7 +52,7 @@ internal static class NodeUri internal static Uri WebSocketEndpoint(Uri serverUrl) { - var builder = new UriBuilder(new Uri(serverUrl, "v1/compute/nodes/connect")) + var builder = new UriBuilder(new Uri(serverUrl, "v1/software-nodes/connect")) { Scheme = serverUrl.Scheme == Uri.UriSchemeHttps ? "wss" : "ws" }; diff --git a/windows-node/Zcbot.WindowsNode/EnrollmentClient.cs b/windows-node/Zcbot.WindowsNode/EnrollmentClient.cs index 9488b7c..a740c40 100644 --- a/windows-node/Zcbot.WindowsNode/EnrollmentClient.cs +++ b/windows-node/Zcbot.WindowsNode/EnrollmentClient.cs @@ -28,7 +28,7 @@ internal static class EnrollmentClient using var client = new HttpClient { BaseAddress = options.ServerUrl, Timeout = TimeSpan.FromSeconds(30) }; using var response = await client.PostAsJsonAsync( - "v1/compute/nodes/enroll", request, cancellationToken); + "v1/software-nodes/enroll", request, cancellationToken); if (!response.IsSuccessStatusCode) { var detail = await response.Content.ReadAsStringAsync(cancellationToken); diff --git a/windows-node/Zcbot.WindowsNode/JobInboxStore.cs b/windows-node/Zcbot.WindowsNode/JobInboxStore.cs index 98d6d08..87252ed 100644 --- a/windows-node/Zcbot.WindowsNode/JobInboxStore.cs +++ b/windows-node/Zcbot.WindowsNode/JobInboxStore.cs @@ -277,7 +277,7 @@ internal sealed class JobInboxStore(string jobsDirectory) && transfer.TryGetProperty("sha256", out var sha) && sha.GetString() is { Length: 64 } && transfer.TryGetProperty("download_path", out var downloadPath) - && downloadPath.GetString()?.StartsWith("/v1/compute/jobs/", StringComparison.Ordinal) == true + && downloadPath.GetString()?.StartsWith("/v1/software-jobs/", StringComparison.Ordinal) == true && downloadPath.GetString()?.EndsWith("/input", StringComparison.Ordinal) == true; internal string InputPath(RecoverableJob job) diff --git a/windows-node/Zcbot.WindowsNode/JobInputDownloader.cs b/windows-node/Zcbot.WindowsNode/JobInputDownloader.cs index 0290309..d17b0d6 100644 --- a/windows-node/Zcbot.WindowsNode/JobInputDownloader.cs +++ b/windows-node/Zcbot.WindowsNode/JobInputDownloader.cs @@ -13,7 +13,7 @@ internal sealed class JobInputDownloader(NodeConfig config, JobInboxStore inbox) throw new InvalidDataException("Job input transfer is missing."); } var downloadPath = transfer.GetProperty("download_path").GetString()!; - if (!downloadPath.StartsWith("/v1/compute/jobs/", StringComparison.Ordinal) + if (!downloadPath.StartsWith("/v1/software-jobs/", StringComparison.Ordinal) || !downloadPath.EndsWith("/input", StringComparison.Ordinal) || !Uri.TryCreate(downloadPath, UriKind.Relative, out var relativeUri)) { diff --git a/windows-node/Zcbot.WindowsNode/JobOutputUploader.cs b/windows-node/Zcbot.WindowsNode/JobOutputUploader.cs index 338b148..8fd114c 100644 --- a/windows-node/Zcbot.WindowsNode/JobOutputUploader.cs +++ b/windows-node/Zcbot.WindowsNode/JobOutputUploader.cs @@ -58,7 +58,7 @@ internal sealed class JobOutputUploader(NodeConfig config) content.Headers.Add("X-Content-SHA256", expectedDigest); content.Headers.Add("X-Content-Length", expectedSize.ToString()); using var response = await client.PutAsync( - $"/v1/compute/jobs/{job.JobId:D}/outputs/{Uri.EscapeDataString(localId)}", + $"/v1/software-jobs/{job.JobId:D}/outputs/{Uri.EscapeDataString(localId)}", content); response.EnsureSuccessStatusCode(); } @@ -68,7 +68,7 @@ internal sealed class JobOutputUploader(NodeConfig config) Encoding.UTF8, "application/json"); using var completeResponse = await client.PostAsync( - $"/v1/compute/jobs/{job.JobId:D}/outputs/complete", completeContent); + $"/v1/software-jobs/{job.JobId:D}/outputs/complete", completeContent); completeResponse.EnsureSuccessStatusCode(); var responseBody = await completeResponse.Content.ReadAsByteArrayAsync(); AtomicWrite(completionPath, responseBody);