Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -166,3 +166,4 @@ node_modules/
# 建站模板必须完整入库(覆盖上面的 lib/ 与 AGENTS.md 规则),但依赖树除外
!docker/site-template/**
docker/site-template/react-vite/node_modules/
.jobtest/
101 changes: 101 additions & 0 deletions src/backend/alembic/versions/ce_0006_job_runtime.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,101 @@
"""CE: add job runtime tables (jobs / job_items / job_calls)

Revision ID: ce_0006
Revises: ce_0005
Create Date: 2026-08-17

作业编排运行时的持久化底座。CE 的 agent_factory 一直注册着 `run_job` 工具、
工作流模式的提示词也在,但这三张表从来没进过 CE 迁移链——作业一提交就会在
建台账那一步炸掉。这条迁移把 CE 补齐到与主仓一致。

台账(job_items)放 DB 而不是沙箱文件:沙箱池化复用会让新 job 读到旧 job 残留的
账本,以 (job_id, item_key) 作复合主键从结构上避免。
"""

import sqlalchemy as sa
from alembic import op
from sqlalchemy.dialects import postgresql

revision = "ce_0006"
down_revision = "ce_0005"
branch_labels = None
depends_on = None

JSONB = postgresql.JSONB(astext_type=sa.Text()).with_variant(sa.JSON(), "sqlite")


def upgrade() -> None:
op.create_table(
"jobs",
sa.Column("job_id", sa.String(64), primary_key=True),
sa.Column("user_id", sa.String(64), nullable=False),
sa.Column("chat_id", sa.String(64)),
sa.Column("name", sa.String(255), server_default=""),
sa.Column("status", sa.String(20), nullable=False, server_default="pending"),
sa.Column("script_path", sa.Text(), server_default=""),
sa.Column("script_text", sa.Text(), server_default=""),
sa.Column("sandbox_session_id", sa.String(128)),
sa.Column("budget", JSONB),
sa.Column("usage", JSONB),
sa.Column("metadata", JSONB),
sa.Column("error_message", sa.Text()),
sa.Column("created_at", sa.TIMESTAMP(timezone=True), server_default=sa.func.now()),
sa.Column("started_at", sa.TIMESTAMP(timezone=True)),
sa.Column("completed_at", sa.TIMESTAMP(timezone=True)),
sa.Column("updated_at", sa.TIMESTAMP(timezone=True), server_default=sa.func.now()),
sa.ForeignKeyConstraint(["user_id"], ["users_shadow.user_id"], ondelete="CASCADE"),
sa.CheckConstraint(
"status IN ('pending','running','paused','completed','failed','cancelled','interrupted')",
name="jobs_status_check",
),
)
op.create_index("idx_jobs_user_id", "jobs", ["user_id"])
op.create_index("idx_jobs_chat_status", "jobs", ["chat_id", "status"])
op.create_index("idx_jobs_status", "jobs", ["status"])

op.create_table(
"job_items",
sa.Column("job_id", sa.String(64), primary_key=True, nullable=False),
sa.Column("item_key", sa.String(128), primary_key=True, nullable=False),
sa.Column("status", sa.String(20), nullable=False, server_default="pending"),
sa.Column("payload", JSONB),
sa.Column("result", JSONB),
sa.Column("review", JSONB),
sa.Column("attempts", sa.Integer(), server_default="0"),
sa.Column("error", sa.Text()),
sa.Column("updated_at", sa.TIMESTAMP(timezone=True), server_default=sa.func.now()),
sa.ForeignKeyConstraint(["job_id"], ["jobs.job_id"], ondelete="CASCADE"),
sa.CheckConstraint(
"status IN ('pending','running','done','not_found','failed','needs_review')",
name="job_items_status_check",
),
)
op.create_index("idx_job_items_job_status", "job_items", ["job_id", "status"])

op.create_table(
"job_calls",
sa.Column("call_id", sa.String(64), primary_key=True),
sa.Column("job_id", sa.String(64), nullable=False),
sa.Column("item_key", sa.String(128)),
sa.Column("seq", sa.Integer(), server_default="0"),
sa.Column("prompt_hash", sa.String(64)),
sa.Column("model", sa.String(128)),
sa.Column("tokens", JSONB),
sa.Column("duration_ms", sa.BigInteger(), server_default="0"),
sa.Column("status", sa.String(20), server_default="running"),
sa.Column("error", sa.Text()),
sa.Column("created_at", sa.TIMESTAMP(timezone=True), server_default=sa.func.now()),
sa.ForeignKeyConstraint(["job_id"], ["jobs.job_id"], ondelete="CASCADE"),
)
op.create_index("idx_job_calls_job_id", "job_calls", ["job_id"])


def downgrade() -> None:
op.drop_index("idx_job_calls_job_id", table_name="job_calls")
op.drop_table("job_calls")
op.drop_index("idx_job_items_job_status", table_name="job_items")
op.drop_table("job_items")
op.drop_index("idx_jobs_status", table_name="jobs")
op.drop_index("idx_jobs_chat_status", table_name="jobs")
op.drop_index("idx_jobs_user_id", table_name="jobs")
op.drop_table("jobs")
44 changes: 44 additions & 0 deletions src/backend/api/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,8 @@ async def lifespan(app: FastAPI):
await _startup_local_sidecars()
await _startup_recover_chat_runs()
await _startup_resume_loops()
await _startup_recover_jobs()
await _startup_orphan_job_reaper()
await _startup_stale_run_reaper()
await _startup_warm_sandbox_pool()
await _startup_idle_session_reaper()
Expand All @@ -78,6 +80,7 @@ async def lifespan(app: FastAPI):
yield
# ── shutdown ──
await _shutdown_stale_run_reaper()
await _shutdown_orphan_job_reaper()
await _shutdown_kb_wiki_worker()
await _shutdown_channel_manager()
await _shutdown_datasource_sidecar_recovery()
Expand Down Expand Up @@ -472,6 +475,36 @@ async def _startup_resume_loops():
logger.warning("[startup] autonomous loop resume failed: %s", exc)


async def _startup_recover_jobs():
"""作业编排:进程重启后活跃 job 全是孤儿(进程内没有 driver task),归位 interrupted。

与 chat_run 不同的是,job 的工作项台账在 DB —— 归位后 ``run_job action=resume``
可直接断点续跑,已完成的项不会重做。
"""
try:
from orchestration import job_runtime

count = await job_runtime.resume_running_jobs()
if count:
logger.info("[startup] orphan jobs recovered: %d", count)
except Exception as exc:
logger.warning("[startup] job orphan recovery failed: %s", exc)


async def _startup_orphan_job_reaper():
"""周期性对账失联的批量作业(驱动没了就没人管,护栏跟着一起没)。"""
try:
import asyncio

from orchestration import job_runtime

task = asyncio.create_task(job_runtime.run_job_reaper_loop())
app.state.orphan_job_reaper_task = task
logger.info("[startup] job orphan reaper started")
except Exception as exc:
logger.warning("[startup] job orphan reaper start failed: %s", exc)


async def _startup_stale_run_reaper():
"""Periodically reap zombie chat_runs stuck in running (fallback behind the watchdog)."""
try:
Expand All @@ -497,6 +530,17 @@ async def _shutdown_stale_run_reaper():
await task


async def _shutdown_orphan_job_reaper():
import asyncio
import contextlib

task = getattr(app.state, "orphan_job_reaper_task", None)
if task is not None:
task.cancel()
with contextlib.suppress(asyncio.CancelledError):
await task


async def _shutdown_channel_manager():
try:
from core.channels.manager import get_manager
Expand Down
5 changes: 5 additions & 0 deletions src/backend/api/routes/v1/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,11 @@
("myspace_folders", "router"),
("batch", "router"),
("internal_batch", "router"),
# 作业编排:沙箱里的作业脚本经 internal_jobs 回调后端(建台账、派子智能体、报终态),
# jobs 是给人看的只读进度视图(输入框上方那条状态条)。两条都必须在——CE 的
# agent_factory 已经注册了 run_job 工具,少了回调路由,作业会在第一发回调 404 当场死掉。
("internal_jobs", "router"),
("jobs", "router"),
("internal_sites", "router"),
("projects", "router"),
("api_keys", "router"),
Expand Down
9 changes: 9 additions & 0 deletions src/backend/api/routes/v1/chats.py
Original file line number Diff line number Diff line change
Expand Up @@ -788,6 +788,7 @@ def _build_ctx(
"plugin_name": request.plugin_name,
"plan_chat": request.plan_chat,
"batch_chat": request.batch_chat,
"workflow_chat": request.workflow_chat,
"disable_batch_plan": request.disable_batch_plan,
**project_ctx,
}
Expand All @@ -803,6 +804,7 @@ def _ensure_chat_session(
agent_name: Optional[str] = None,
plan_chat: bool = False,
batch_chat: bool = False,
workflow_chat: bool = False,
project_id: Optional[str] = None,
):
extra_data: Dict[str, Any] = {"chat_id": chat_id}
Expand All @@ -814,6 +816,8 @@ def _ensure_chat_session(
extra_data["plan_chat"] = True
if batch_chat:
extra_data["batch_chat"] = True
if workflow_chat:
extra_data["workflow_chat"] = True
# Prefer the edition-aware access resolver before creating a session.
pair = chat_service.get_session_with_access(chat_id, user_id)
if pair is not None:
Expand Down Expand Up @@ -850,6 +854,9 @@ def _ensure_chat_session(
if batch_chat and not existing_meta.get("batch_chat"):
merged["batch_chat"] = True
dirty = True
if workflow_chat and not existing_meta.get("workflow_chat"):
merged["workflow_chat"] = True
dirty = True
if dirty:
chat_service.update_session(chat_id, user_id, {"extra_data": merged})
return session
Expand Down Expand Up @@ -928,6 +935,7 @@ async def chat_send(
agent_name=_agent_name,
plan_chat=request.plan_chat,
batch_chat=request.batch_chat,
workflow_chat=request.workflow_chat,
project_id=request.project_id,
)
# Link orphan artifacts (uploaded before session existed) to this chat
Expand Down Expand Up @@ -1129,6 +1137,7 @@ def _read_messages():
agent_name=_agent_name_stream,
plan_chat=request.plan_chat,
batch_chat=request.batch_chat,
workflow_chat=request.workflow_chat,
project_id=request.project_id,
)

Expand Down
Loading
Loading