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/
1 change: 1 addition & 0 deletions document/en/architecture/frontend.md
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ SSE streams bypass the JSON channel of `api.ts`; they are consumed directly from
| `pageConfigStore` | Page configuration (branding, navigation, copy — drives white-labeling) |
| `editionStore` | Consumer of the `/v1/meta/edition` probe: edition and license feature-bit map |
| `modelCapabilitiesStore` | Main-model capability probing (thinking / vision etc.) |
| `sidebarOrderStore` | Manual drag-and-drop order of the sidebar chat list (persisted locally + in `users_shadow.metadata`) |

## Hooks

Expand Down
1 change: 1 addition & 0 deletions document/zh-CN/architecture/frontend.md
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ SSE 流式不走 `api.ts` 的 JSON 通道,由 `hooks/useStreaming.ts` 直接
| `pageConfigStore` | 页面配置(品牌、导航、文案——驱动 white-label) |
| `editionStore` | `/v1/meta/edition` 探针的消费端:版本与 license 能力位布尔表 |
| `modelCapabilitiesStore` | 主模型能力探测(思考 / 视觉等) |
| `sidebarOrderStore` | 侧边栏对话列表的手动拖拽顺序(本地 + `users_shadow.metadata` 双层持久化) |

## Hooks

Expand Down
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
68 changes: 68 additions & 0 deletions src/backend/api/routes/v1/chats.py
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,18 @@ class UpdateChatRequest(BaseModel):
metadata: Optional[dict] = Field(None, description="Additional metadata")


class UpdateSidebarOrderRequest(BaseModel):
"""Request model for persisting the sidebar's manual (drag-and-drop) order."""

order: List[str] = Field(default_factory=list, description="Chat ids in manual order")


# Sidebar manual order lives in users_shadow.metadata (no schema migration needed):
# it is a pure UI preference of "which chat sits where", not session state.
SIDEBAR_ORDER_KEY = "sidebar_chat_order"
SIDEBAR_ORDER_MAX = 500


def _session_to_dict(s) -> dict:
"""Convert a ChatSession ORM object to the edition-neutral API response."""
return {
Expand Down Expand Up @@ -287,6 +299,53 @@ async def list_pending_confirms(
return success_response(data={"items": items})


def _dedup_id_list(raw: Optional[list]) -> List[str]:
"""Clean + de-duplicate an id list, preserving first-seen order."""
seen: set = set()
out: List[str] = []
for cid in _clean_id_list(raw):
if cid in seen:
continue
seen.add(cid)
out.append(cid)
return out


@router.get("/sidebar-order", summary="获取侧边栏手动排序")
async def get_sidebar_order(
user: UserContext = Depends(get_current_user),
db: Session = Depends(get_db),
):
"""侧边栏对话列表的手动拖拽顺序(chat_id 序列)。

没拖过的账号返回空数组——前端据此退回「置顶 + 最近更新」默认排序。

注意:本路由必须声明在 ``GET /{chat_id}`` **之前**,否则会被 path 参数
路由吞掉(FastAPI 按声明顺序匹配)。
"""
user_settings = UserService(db).get_user_settings(str(user.user_id))
order = _dedup_id_list(user_settings.get(SIDEBAR_ORDER_KEY))[:SIDEBAR_ORDER_MAX]
return success_response(data={"order": order})


@router.put("/sidebar-order", summary="保存侧边栏手动排序")
async def update_sidebar_order(
request: UpdateSidebarOrderRequest,
user: UserContext = Depends(get_current_user),
db: Session = Depends(get_db),
):
"""整表覆盖写入手动顺序;空数组 = 恢复默认排序。

不校验 chat_id 是否存在:顺序表是纯 UI 偏好,已删除的会话留在表里也只是
查不到对应项而被忽略,反倒省掉一次全表校验。超出上限的尾部直接截断。
"""
order = _dedup_id_list(request.order)[:SIDEBAR_ORDER_MAX]
UserService(db).update_user_metadata(
user_id=str(user.user_id), patch={SIDEBAR_ORDER_KEY: order}
)
return success_response(data={"order": order})


@router.get("/{chat_id}", summary="获取会话详情")
async def get_chat(
chat_id: str, user: UserContext = Depends(get_current_user), db: Session = Depends(get_db)
Expand Down Expand Up @@ -788,6 +847,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 +863,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 +875,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 +913,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 +994,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 +1196,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