Files
cowork-local/core/co4e_run_manager.py
T
f9f6bc01fd
CI / test (push) Canceled after 0s
Feature/delta team/epic r04 (#7)
## Summary

epic r04 - begin refactor

## Change Type

- [x] Cowork feature
- [ ] Bug fix
- [ ] Core AI contribution
- [ ] Test / hardening
- [ ] Performance
- [ ] Documentation

## Related Work

Cowork Task:

Core Repo: http://34.143.229.138/gitea-admin/fsg-ai-core-assets

Core AI Issue:

Core Task:

Related PR:

## Scope

What is intentionally included?

What is intentionally NOT included?

## Validation

- [ ] Unit tests
- [ ] Integration tests
- [ ] Manual verification
- [ ] Regression check

Commands / evidence:

## Security Impact

Permission / credential / network / customer data impact:

## Compatibility

- [ ] No breaking change
- [ ] Breaking change documented

## Reviewer Notes

Anything Cowork reviewers should pay attention to.

---------

Co-authored-by: Anh Tran Nguyen Minh <anhtnm1@fpt.com>
Co-authored-by: Huong Le Thi Thien <huongltt35@fpt.com>
Co-authored-by: Nam Pham Dinh Thanh <nampdt@fpt.com>
Co-authored-by: Vu Dam Tuan <vudt15@fpt.com>
Co-authored-by: Hiep Ha Van <hiephv3@fpt.com>
Co-authored-by: Lam Hoang Van <lamhv7@fpt.com>
Reviewed-on: #7
Co-authored-by: Duy Le Huu <duylh19@fpt.com>
2026-08-31 05:15:13 +00:00

373 lines
17 KiB
Python

"""Co4E run manager — tracks concurrent flow runs and their live status.
The app already runs several Cowork/Code turns at once (each on its own
``AgentWorker`` QThread). This brings the same to Co4E: any number of flows can
run in parallel, each on its own worker, with live per-run status
(running / done / error / stopped) and progress (done / total steps).
The manager is the single source of truth — the Co4E "Running flows" list reads
it and refreshes on every ``changed`` signal (and when the tab is re-shown), so
switching sub-tabs never loses, stales, or drops status. ``event`` re-emits each
run's node-level events tagged with the run id, so the canvas/chat can mirror
the run that is currently open.
"""
from __future__ import annotations
from ..infrastructure.persistence.json.atomic_json_file import AtomicJsonFile
import json
from pathlib import Path
from typing import Dict, List, Optional
from PySide6.QtCore import QObject, Signal
from .co4e import STEP_DONE, STEP_ERROR, STEP_PLANNED, Workflow
_TERMINAL_NODE = {STEP_DONE, STEP_ERROR, STEP_PLANNED}
_HISTORY_CAP = 500 # keep the most-recent N runs on disk
def _now_str() -> str:
"""Mốc thời gian hiện tại dạng 'YYYY-MM-DD HH:MM' cho lịch sử run."""
from datetime import datetime
return datetime.now().strftime("%Y-%m-%d %H:%M")
def _current_user() -> str:
"""Best-effort creator name for a run (signed-in MS365 identity → OS user)."""
import os
return (os.environ.get("USERNAME") or os.environ.get("USER") or "you")
class RunHandle:
"""Live state for one flow run. Mutated by the manager as events arrive."""
def __init__(self, run_id: str, wf_id: str, name: str, total: int,
plan_mode: bool, manual: bool, created_by: str = "", created_at: str = "",
project_id: str = ""):
"""Một lượt chạy workflow đang sống trong bộ nhớ.
``total`` âm bị kẹp về 0 — số bước không thể âm, và để lọt xuống thì thanh
tiến độ vẽ ngược.
"""
self.id = run_id
self.wf_id = wf_id
self.name = name
self.project_id = project_id # workspace this run belongs to (Flow Status is per-project)
self.total = max(0, total)
self.done = 0
self.status = "running" # running | done | error | stopped
self.plan_mode = plan_mode
self.manual = manual
self.created_by = created_by
self.created_at = created_at
self.error = ""
self.node_status: Dict[str, str] = {}
self.worker = None
self.wf = None # the Workflow this run ran — lets Runs reopen it
# even after its tab was closed / it was never saved
self.out_dir = "" # workspace folder this run wrote its files into
@property
def running(self) -> bool:
"""Lượt chạy này còn đang chạy hay không."""
return self.status == "running"
def progress_text(self) -> str:
"""Chuỗi tiến độ 'xong/tổng'; chưa biết tổng thì hiện trạng thái."""
return f"{self.done}/{self.total}" if self.total else self.status
# ---- persistence ------------------------------------------------------
def to_record(self) -> dict:
"""Serialize for the on-disk run history. The workflow snapshot is kept
so a past run can be reopened even if its saved flow was later edited or
deleted."""
from .co4e import workflow_to_dict
return {
"id": self.id, "wf_id": self.wf_id, "name": self.name,
"total": self.total, "done": self.done, "status": self.status,
"plan_mode": self.plan_mode, "manual": self.manual,
"created_by": self.created_by, "created_at": self.created_at,
"error": self.error, "node_status": dict(self.node_status),
"wf": workflow_to_dict(self.wf) if self.wf is not None else None,
"out_dir": self.out_dir, "project_id": self.project_id,
}
@classmethod
def from_record(cls, rec: dict) -> "RunHandle":
"""Dựng lại một ``RunHandle`` từ bản ghi đọc trong lịch sử trên đĩa."""
from .co4e import workflow_from_dict
rec = dict(rec or {})
h = cls(str(rec.get("id", "")), str(rec.get("wf_id", "")),
rec.get("name", ""), int(rec.get("total", 0) or 0),
bool(rec.get("plan_mode")), bool(rec.get("manual")),
created_by=rec.get("created_by", ""), created_at=rec.get("created_at", ""))
h.done = int(rec.get("done", 0) or 0)
h.status = rec.get("status", "done")
# a run persisted as "running" means the app closed mid-run — its worker
# is gone, so it's no longer live: settle it as "stopped".
if h.status == "running":
h.status = "stopped"
h.error = rec.get("error", "")
h.node_status = dict(rec.get("node_status") or {})
h.out_dir = rec.get("out_dir", "")
h.project_id = rec.get("project_id", "")
wfd = rec.get("wf")
h.wf = workflow_from_dict(wfd) if wfd else None
return h
class Co4ERunManager(QObject):
"""Quản lý vòng đời nhiều lượt chạy luồng Co4E cùng lúc.
Flow Status lọc theo project, nên hầu hết truy vấn ở đây chỉ tính run thuộc
workspace ĐANG chọn — xem ``_belongs``.
"""
changed = Signal() # any run's status/progress changed → refresh views
event = Signal(str, dict) # (run_id, ev) — node-level events, for mirroring
def __init__(self, ctx):
"""Dựng bộ quản lý run và khôi phục lịch sử cũ ngay, để tab Flow Status có nội
dung ngay khi mở chứ không trống cho tới lần chạy đầu tiên.
"""
super().__init__()
self.ctx = ctx
self._runs: Dict[str, RunHandle] = {}
self._seq = 0
self._output_root: Optional[Path] = None # active workspace's co4e output base
self._project_id: str = "" # active workspace — Flow Status is filtered to it
self._load_history() # restore past runs so the Flow Status tab
# keeps its full history across restarts
# every status/progress change is persisted, so history is never lost
self.changed.connect(self._save_history)
# ---- persistence ------------------------------------------------------
def _history_path(self) -> Path:
"""Đường dẫn file lịch sử run."""
from .co4e import CO4E_DIR
return CO4E_DIR / "run_history.json"
def _load_history(self) -> None:
"""Khôi phục lịch sử run từ đĩa lúc khởi động; file hỏng thì bỏ qua lặng lẽ."""
path = self._history_path()
try:
data = json.loads(path.read_text(encoding="utf-8"))
except (OSError, ValueError):
return
max_seq = 0
for rec in data.get("runs", []):
try:
handle = RunHandle.from_record(rec)
except Exception:
continue
if not handle.id:
continue
self._runs[handle.id] = handle
if handle.id.startswith("run") and handle.id[3:].isdigit():
max_seq = max(max_seq, int(handle.id[3:]))
self._seq = max_seq # avoid minting ids that collide with history
def _save_history(self) -> None:
"""Ghi ``_HISTORY_CAP`` run gần nhất xuống đĩa."""
path = self._history_path()
runs = list(self._runs.values())[-_HISTORY_CAP:]
payload = {"runs": [h.to_record() for h in runs]}
try:
path.parent.mkdir(parents=True, exist_ok=True)
# AtomicJsonFile thay cho tmp+replace tự viết: bản cũ thiếu fsync
# (dữ liệu có thể còn trong bộ đệm khi mất điện) và dùng thẳng
# Path.replace, vốn thỉnh thoảng bị Defender từ chối trên Windows.
AtomicJsonFile(path).write(payload)
except OSError:
pass
# ---- lifecycle --------------------------------------------------------
def _next_id(self) -> str:
"""Sinh id run kế tiếp dạng 'runN'."""
self._seq += 1
return f"run{self._seq}"
def start(self, wf: Workflow, *, skill_map: Optional[Dict[str, str]] = None,
plan_mode: bool = False, only_nodes: Optional[set] = None,
seed_outputs: Optional[Dict[str, str]] = None,
manual: bool = False, label: Optional[str] = None) -> str:
"""Launch a flow (or a subset via ``only_nodes``) on its own worker and
return the new run id. Runs concurrently with every other active run."""
from . import co4e_runner
from .worker import AgentWorker
import copy as _copy
run_id = self._next_id()
total = len(only_nodes) if only_nodes else len(wf.nodes)
handle = RunHandle(run_id, wf.id, label or wf.name, total, plan_mode, manual,
created_by=_current_user(), created_at=_now_str(),
project_id=self._project_id)
# Keep a DEEP COPY so Runs can reopen the flow exactly as it ran — even if
# the live canvas is later edited (nodes moved, edges removed) while the
# run is still tracked. A shared reference caused reopened runs to show a
# disconnected/blank graph.
handle.wf = _copy.deepcopy(wf)
nodes = list(wf.nodes)
edges = list(wf.edges)
out_dir = self._out_dir(wf)
handle.out_dir = str(out_dir) # where this run writes its files (workspace)
ctx = self.ctx
sk = dict(skill_map or {})
only = set(only_nodes) if only_nodes else None
seed = dict(seed_outputs or {})
run_label = handle.name
def job(worker: AgentWorker):
"""Chạy nền: thực thi luồng, chuyển tiếp sự kiện tiến độ và cờ huỷ."""
return co4e_runner.run_workflow(
ctx, nodes, edges, out_dir, worker.emit_event, worker.is_cancelled,
plan_mode=plan_mode, skill_map=sk, only_nodes=only, seed_outputs=seed,
usage_label=run_label)
worker = AgentWorker(job)
handle.worker = worker
worker.event.connect(lambda ev, rid=run_id: self._on_event(rid, ev))
worker.finished_ok.connect(lambda _r, rid=run_id: self._on_finished(rid))
worker.failed.connect(lambda e, rid=run_id: self._on_failed(rid, e))
self._runs[run_id] = handle
worker.start()
self.changed.emit()
return run_id
# ---- worker callbacks -------------------------------------------------
def _on_event(self, run_id: str, ev: dict) -> None:
"""Nhận sự kiện từ luồng đang chạy và cập nhật trạng thái/tiến độ của run."""
handle = self._runs.get(run_id)
if handle is not None and isinstance(ev, dict):
t = ev.get("type")
if t == "node_status":
handle.node_status[ev.get("node_id")] = ev.get("status")
handle.done = sum(1 for s in handle.node_status.values() if s in _TERMINAL_NODE)
self.changed.emit()
elif t == "run_done":
if handle.status == "running":
handle.status = "done" if ev.get("ok", True) else "error"
self.changed.emit()
self.event.emit(run_id, ev)
def _on_finished(self, run_id: str) -> None:
"""Job kết thúc mà không phát ``run_done``: chốt trạng thái về 'done'.
Lẽ ra không xảy ra, nhưng thiếu bước này thì run kẹt ở 'running' mãi.
"""
handle = self._runs.get(run_id)
if handle is not None and handle.status == "running":
# job returned without a run_done event (shouldn't happen) — settle it
handle.status = "done"
self.changed.emit()
def _on_failed(self, run_id: str, err: str) -> None:
"""Job ném lỗi: ghi lỗi vào bản ghi run và báo ra ngoài."""
handle = self._runs.get(run_id)
if handle is not None:
handle.status = "error"
handle.error = str(err)
self.event.emit(run_id, {"type": "run_error", "error": str(err)})
self.changed.emit()
# ---- control ----------------------------------------------------------
def stop(self, run_id: str) -> None:
"""Yêu cầu dừng một run đang chạy."""
handle = self._runs.get(run_id)
if handle is not None and handle.worker is not None and handle.running:
handle.worker.request_stop()
handle.status = "stopped"
self.changed.emit()
def stop_all(self) -> None:
# Only the CURRENT workspace's runs (Flow Status is per-project).
"""Dừng mọi run của workspace đang chọn."""
for run_id in [r for r, h in self._runs.items() if self._belongs(h)]:
self.stop(run_id)
def rename(self, run_id: str, new_name: str) -> None:
"""Rename a run in the Flow Status history (and its kept workflow snapshot),
then persist + refresh views. No-op on a blank name / unknown run."""
handle = self._runs.get(run_id)
new_name = (new_name or "").strip()
if handle is None or not new_name or new_name == handle.name:
return
handle.name = new_name
if handle.wf is not None:
handle.wf.name = new_name
self.changed.emit()
def remove(self, run_id: str) -> None:
"""Xoá một run khỏi lịch sử; đang chạy thì dừng trước."""
handle = self._runs.get(run_id)
if handle is not None and handle.running:
self.stop(run_id)
self._runs.pop(run_id, None)
self.changed.emit()
def clear_finished(self) -> None:
# Only clear finished runs of the CURRENT workspace.
"""Xoá mọi run đã kết thúc của workspace đang chọn, giữ nguyên run đang chạy."""
for run_id in [r for r, h in self._runs.items() if not h.running and self._belongs(h)]:
self._runs.pop(run_id, None)
self.changed.emit()
# ---- queries ----------------------------------------------------------
def _belongs(self, h: RunHandle) -> bool:
"""Whether a run belongs to the currently-selected workspace."""
return getattr(h, "project_id", "") == self._project_id
def runs(self) -> List[RunHandle]:
"""Runs of the CURRENT workspace only — Flow Status is per-project."""
return [h for h in self._runs.values() if self._belongs(h)]
def all_runs(self) -> List[RunHandle]:
"""Every tracked run across all workspaces (background tracking)."""
return list(self._runs.values())
def get(self, run_id: str) -> Optional[RunHandle]:
"""Bản ghi của một run theo id; ``None`` nếu không có."""
return self._runs.get(run_id)
def active_count(self) -> int:
"""Số run đang chạy của workspace đang chọn."""
return sum(1 for h in self._runs.values() if h.running and self._belongs(h))
def set_current_project(self, project_id: str) -> None:
"""Filter Flow Status (and new runs) to this workspace. Runs started while
this is set are tagged with it; the Runs view shows only matching runs."""
pid = project_id or ""
if pid != self._project_id:
self._project_id = pid
self.changed.emit() # re-render Flow Status for the new workspace
def set_output_root(self, root: Optional[Path]) -> None:
"""Point flow outputs at the SELECTED workspace's co4e folder (set by the
Co4E tab when a project is chosen). ``None`` → fall back to the global
Cowork output dir."""
self._output_root = Path(root) if root else None
def _out_dir(self, wf: Workflow) -> Path:
# Flow deliverables are written into the SELECTED workspace (the active
# project's folder) so they land where the user works with files (Folder
# tab), not in the config/install folder. One subfolder per flow keeps
# runs tidy. Falls back to the global Cowork output dir when no workspace
# is selected.
"""Thư mục ghi kết quả của một luồng, tạo sẵn nếu chưa có.
Ưu tiên thư mục của workspace đang chọn để file rơi đúng chỗ người dùng làm
việc (màn Thư mục), không rơi vào thư mục cài đặt.
"""
from .co4e import slugify
base = self._output_root
if base is None:
try:
base = self.ctx.config.cowork_output_dir() / "co4e"
except Exception: # noqa: BLE001 - fall back to the config dir if unavailable
from .co4e import CO4E_DIR
base = CO4E_DIR / "runs" / "co4e"
d = Path(base) / slugify(wf.name or "flow")
d.mkdir(parents=True, exist_ok=True)
return d