diff --git a/.tdd/spec-real-agent-heartbeat-truth-v1.md b/.tdd/spec-real-agent-heartbeat-truth-v1.md new file mode 100644 index 0000000..4562fcd --- /dev/null +++ b/.tdd/spec-real-agent-heartbeat-truth-v1.md @@ -0,0 +1,93 @@ +# Real Agent Heartbeat Truth v1 + +## Problem + +The Department Campus currently treats any fresh lifecycle event as live work, +including terminal `done` and `failed` events, while six idle residents animate +without a running task. Bridge lifecycle timestamps do not prove that an AI +process is still executing. + +## Locked criteria + +### AC-1 — live state requires a verified heartbeat + +- Campus `state=active` and live agent sprites are driven only by fresh, + verified heartbeat records in `working` state. +- The heartbeat identifies an exact canonical project, its responsible + canonical agent, a safe run id, a safe session id, and a UTC heartbeat time. +- A valid heartbeat is no older than 45 seconds and is not in the future. + +### AC-2 — lifecycle history is not live presence + +- Bridge `pending`, `running`, `done`, and `failed` rows do not create live + Campus agents without a matching verified heartbeat. +- Direct terminal projection events (`done` and `failed`) never produce + `state=active`, active-agent counts, movement, or live routes. +- Existing completed-task and Git history surfaces remain unchanged. + +### AC-3 — exact fail-closed identity + +- The heartbeat project must exist in `CAMPUS_PROJECTS` and its `agent_id` must + equal that project's registered owner. +- Unknown projects/agents, mismatches, unsafe ids, duplicate live identities, + malformed JSON, missing fields, stale timestamps, and future timestamps are + omitted without partial or inferred activity. +- Raw task text, prompts, paths, tool output, credentials, and heartbeat file + contents never enter the public projection. + +### AC-4 — honest motion + +- All seven persistent residents are static while idle. +- Only a verified live ephemeral agent may hide its matching resident and use + route/sprite movement. +- Existing keyboard, mobile, and `prefers-reduced-motion` behavior remains. + +### AC-5 — JARVIS producer keeps truth fresh only while a provider runs + +- `jarvis-pixel-agent-event` stores the canonical project, canonical agent id, + run id, session id, state, and heartbeat timestamp atomically. +- A `heartbeat` command refreshes an existing working record without changing + its task/status copy and without posting a synthetic Pixel tool event. +- `jarvis-agent-pipeline` refreshes heartbeat at a bounded interval while the + Claude/Codex provider process runs, stops the loop afterward, and writes idle + on completion/failure. + +### AC-6 — concurrent producer updates preserve every agent + +- Heartbeat storage updates are protected across independent producer + processes, so one agent's read-modify-write cannot discard another agent's + newly written or refreshed record. +- The protection remains dependency-free, bounded, and fail-closed; an + abandoned lock must not block the producer forever. + +### AC-7 — terminal state follows the complete provider process tree + +- On `HUP`, `INT`, `TERM`, or pipeline exit, cleanup stops and waits for the + complete provider process group, not only its shell wrapper. +- The terminal Pixel `done` event is written only after that process group is + no longer running, so Campus cannot report idle while Claude/Codex work + continues in an orphaned descendant. + +## Error and boundary criteria + +- ERR-1: unreadable or malformed heartbeat storage returns no live events. +- ERR-2: a valid heartbeat mixed with a duplicate or identity conflict for the + same canonical agent fails that identity closed. +- EC-1: exactly 45 seconds old is accepted; older is stale. +- EC-2: terminal Bridge history may remain available elsewhere but never + changes Campus live counts or motion. +- EC-3: two independent producer processes updating different canonical agents + preserve both records. +- EC-4: interrupted cleanup is bounded and still removes provider temp output. + +## Constraints + +- Public/read-only behavior and owner-field privacy remain unchanged. +- No new dependency, credential, network service, dispatch control, or public + mutation endpoint. +- Maximum three visible live tasks remains. + +## Out of scope + +- Codex desktop tasks that do not emit the heartbeat contract. +- Merge, deploy, restart, credentials, permissions, pairing, or Remote changes. diff --git a/builder/dashboard-assets/script.js b/builder/dashboard-assets/script.js index 8002e6a..2700e01 100644 --- a/builder/dashboard-assets/script.js +++ b/builder/dashboard-assets/script.js @@ -103,6 +103,7 @@ waiting: 'department', failed: 'department', }; + const liveStatuses = ['active', 'testing']; const movingStatuses = ['active', 'testing']; const stateMessages = { loading: 'Загрузка кампуса…', @@ -473,7 +474,9 @@ function renderDepartmentCampus(payload) { const state = payload && typeof payload.state === 'string' ? payload.state : 'unavailable'; - const events = state === 'active' && Array.isArray(payload.events) ? payload.events : []; + const events = state === 'active' && Array.isArray(payload.events) + ? payload.events.filter((event) => liveStatuses.includes(event?.status)) + : []; if ( state === 'empty' || state === 'stale' @@ -494,7 +497,7 @@ } const matchedEvents = []; events.forEach((event) => { - if (!Object.prototype.hasOwnProperty.call(statusLabels, event?.status)) return; + if (!liveStatuses.includes(event?.status)) return; const folder = projectFolderForEvent(event); if (!folder) return; matchedEvents.push(event); @@ -533,7 +536,10 @@ animateCampusJourney(button, shouldAnimate, managerOriginRect); }); const activeAgentCount = new Set( - matchedEvents.map((event) => String(event.agent_id || '')).filter(Boolean), + matchedEvents + .filter((event) => liveStatuses.includes(event.status)) + .map((event) => String(event.agent_id || '')) + .filter(Boolean), ).size; setCampusState( 'active', diff --git a/builder/dashboard-server-m4.py b/builder/dashboard-server-m4.py index d29548a..2453e87 100644 --- a/builder/dashboard-server-m4.py +++ b/builder/dashboard-server-m4.py @@ -40,6 +40,12 @@ JARVIS_PIPELINE_SCRIPT = SCRIPTS_DIR / "jarvis-agent-pipeline" JARVIS_PIPELINE_REPORT_DIR = HOME / "Library" / "Logs" / "jarvis-agent-pipeline" JARVIS_PIPELINE_LOG_FILE = HOME / "Library" / "Logs" / "dashboard-jarvis-pipeline-run.log" +JARVIS_PIXEL_AGENT_DETAILS_FILE = Path( + os.environ.get( + "JARVIS_PIXEL_AGENT_DETAILS_FILE", + str(HOME / ".pixel-agents" / "jarvis-agent-details.json"), + ) +).expanduser() JARVIS_REPO = HOME / "jarvis" JARVIS_PYTHON = JARVIS_REPO / ".venv" / "bin" / "python" JARVIS_VAULT_ROOT = Path( @@ -1193,7 +1199,7 @@ def _department_snapshot_time(value: object) -> datetime | None: def _department_campus_state(state: str, *, now: datetime) -> dict: - payload = department_campus_projection([], now=now) + payload = department_campus_projection([], heartbeats=[], now=now) payload["state"] = state return payload @@ -1278,9 +1284,34 @@ def _campus_bridge_event(task: dict) -> dict | None: return event +def _department_campus_heartbeats(path: Path) -> list[dict]: + """Read only the six heartbeat proof fields from local producer storage.""" + try: + document = json.loads(path.read_text(encoding="utf-8")) + except (OSError, UnicodeError, json.JSONDecodeError): + return [] + agents = document.get("agents") if isinstance(document, dict) else None + if not isinstance(agents, dict): + return [] + normalized: list[dict] = [] + for record in agents.values(): + if not isinstance(record, dict): + continue + normalized.append({ + "project": record.get("project"), + "agent_id": record.get("agentId", record.get("agent_id")), + "run_id": record.get("runId", record.get("run_id")), + "session_id": record.get("sessionId", record.get("session_id")), + "state": record.get("state"), + "heartbeat_at": record.get("heartbeatAt", record.get("heartbeat_at")), + }) + return normalized + + def _department_campus_payload( data: object, *, + heartbeat_path: Path | None = None, now: datetime | None = None, owner_view: bool = False, ) -> dict: @@ -1289,8 +1320,16 @@ def _department_campus_payload( if current.tzinfo is None: current = current.replace(tzinfo=timezone.utc) current = current.astimezone(timezone.utc) + heartbeats = _department_campus_heartbeats( + heartbeat_path or JARVIS_PIXEL_AGENT_DETAILS_FILE + ) if not isinstance(data, dict) or not isinstance(data.get("tasks"), list): - return department_campus_projection(None, now=current, owner_view=owner_view) + return department_campus_projection( + None, + heartbeats=heartbeats, + now=current, + owner_view=owner_view, + ) candidates: list[tuple[int, datetime, list]] = [] malformed_verified_snapshot = False @@ -1323,7 +1362,12 @@ def _department_campus_payload( if not candidates: if malformed_verified_snapshot: - return department_campus_projection(None, now=current, owner_view=owner_view) + return department_campus_projection( + None, + heartbeats=heartbeats, + now=current, + owner_view=owner_view, + ) events = [ event for task in data["tasks"] @@ -1331,8 +1375,12 @@ def _department_campus_payload( for event in [_campus_bridge_event(task)] if event is not None ] + # A heartbeat can stand alone only when Bridge has no lifecycle rows. + # Non-empty but unverified rows are not allowed to borrow that proof. + fallback_heartbeats = heartbeats if events or not data["tasks"] else [] return department_campus_projection( events, + heartbeats=fallback_heartbeats, now=current, max_tasks=3, owner_view=owner_view, @@ -1344,6 +1392,7 @@ def _department_campus_payload( return _department_campus_state("stale", now=current) return department_campus_projection( events, + heartbeats=heartbeats, now=current, max_tasks=3, owner_view=owner_view, diff --git a/builder/dashboard_builder/department_campus.py b/builder/dashboard_builder/department_campus.py index 027f632..deb50b1 100644 --- a/builder/dashboard_builder/department_campus.py +++ b/builder/dashboard_builder/department_campus.py @@ -170,9 +170,9 @@ def campus_project_for_event(event: object) -> MappingProxyType | None: "sprite_x": "-96px", "sprite_step_x": "-128px", "sprite_y": "-192px", - "wandering": True, - "walk_duration": "11s", - "walk_delay": "-4s", + "wandering": False, + "walk_duration": "0s", + "walk_delay": "0s", }, "BUILDER": { "name": "Разработчик", @@ -186,9 +186,9 @@ def campus_project_for_event(event: object) -> MappingProxyType | None: "sprite_x": "-288px", "sprite_step_x": "-320px", "sprite_y": "-192px", - "wandering": True, - "walk_duration": "13s", - "walk_delay": "-7s", + "wandering": False, + "walk_duration": "0s", + "walk_delay": "0s", }, "DESIGNER": { "name": "Дизайнер", @@ -197,9 +197,9 @@ def campus_project_for_event(event: object) -> MappingProxyType | None: "sprite_x": "-192px", "sprite_step_x": "-224px", "sprite_y": "-64px", - "wandering": True, - "walk_duration": "14s", - "walk_delay": "-5s", + "wandering": False, + "walk_duration": "0s", + "walk_delay": "0s", }, "INFRASTRUCTURE": { "name": "Инженер инфраструктуры", @@ -208,9 +208,9 @@ def campus_project_for_event(event: object) -> MappingProxyType | None: "sprite_x": "-288px", "sprite_step_x": "-320px", "sprite_y": "-64px", - "wandering": True, - "walk_duration": "15s", - "walk_delay": "-9s", + "wandering": False, + "walk_duration": "0s", + "walk_delay": "0s", }, "VAULT": { "name": "Хранитель знаний", @@ -219,9 +219,9 @@ def campus_project_for_event(event: object) -> MappingProxyType | None: "sprite_x": "0px", "sprite_step_x": "-32px", "sprite_y": "-64px", - "wandering": True, - "walk_duration": "12s", - "walk_delay": "-2s", + "wandering": False, + "walk_duration": "0s", + "walk_delay": "0s", }, "ANALYST": { "name": "Аналитик", @@ -230,9 +230,9 @@ def campus_project_for_event(event: object) -> MappingProxyType | None: "sprite_x": "-96px", "sprite_step_x": "-128px", "sprite_y": "-64px", - "wandering": True, - "walk_duration": "10s", - "walk_delay": "-6s", + "wandering": False, + "walk_duration": "0s", + "walk_delay": "0s", }, } @@ -274,6 +274,7 @@ def campus_project_for_event(event: object) -> MappingProxyType | None: "issue_url", ) _FRESH_SECONDS = 30 * 60 +_HEARTBEAT_FRESH_SECONDS = 45 _SAFE_ID = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.:-]{0,95}$") _UNSAFE_TEXT = re.compile( r"(?:" @@ -508,6 +509,7 @@ def _validated_event( or type(event.get("ephemeral")) is not bool or event.get("ephemeral") is not True or updated is None + or campus_project_for_event(event) is None ): return None, "invalid", None age = (_utc_now(now) - updated).total_seconds() @@ -537,9 +539,112 @@ def _validated_event( return public, "valid", updated +def _validated_heartbeat( + heartbeat: object, + *, + now: datetime, +) -> tuple[dict[str, Any] | None, str]: + """Validate only the public identity proof needed for live presence.""" + if not isinstance(heartbeat, dict): + return None, "invalid" + project_name = heartbeat.get("project") + agent_id = heartbeat.get("agent_id") + run_id = _safe_id(heartbeat.get("run_id")) + session_id = _safe_id(heartbeat.get("session_id")) + heartbeat_at = _parse_time(heartbeat.get("heartbeat_at")) + project = next( + ( + record + for record in CAMPUS_PROJECTS + if record["project"] == project_name + ), + None, + ) + if ( + project is None + or agent_id != project["agent_id"] + or run_id is None + or session_id is None + or heartbeat.get("state") != "working" + or heartbeat_at is None + ): + return None, "invalid" + age = (_utc_now(now) - heartbeat_at).total_seconds() + if age < 0: + return None, "invalid" + if age > _HEARTBEAT_FRESH_SECONDS: + return None, "stale" + return { + "project": project["project"], + "department_id": project["department_id"], + "agent_id": project["agent_id"], + "run_id": run_id, + "heartbeat_at": heartbeat_at, + }, "valid" + + +def _heartbeat_event( + heartbeat: dict[str, Any], + events: list[object], + *, + now: datetime, + owner_view: bool, +) -> tuple[dict[str, Any] | None, datetime]: + """Join a heartbeat to safe lifecycle copy or synthesize a minimal event.""" + matching: list[tuple[int, dict[str, Any], datetime]] = [] + for index, raw_event in enumerate(events): + if not isinstance(raw_event, dict): + continue + if ( + raw_event.get("project") != heartbeat["project"] + or raw_event.get("agent_id") != heartbeat["agent_id"] + or raw_event.get("task_id") != heartbeat["run_id"] + ): + continue + if raw_event.get("status") in {"done", "failed"}: + return None, heartbeat["heartbeat_at"] + event, _, updated = _validated_event( + raw_event, + now=now, + owner_view=owner_view, + ) + if event is None or updated is None: + return None, heartbeat["heartbeat_at"] + if event is not None and updated is not None and event["status"] in {"active", "testing"}: + matching.append((index, event, updated)) + if matching: + # The newest safe lifecycle copy enriches the heartbeat. Stable sorting + # retains the first source row when timestamps tie. + _, event, _ = sorted( + matching, + key=lambda item: item[2], + reverse=True, + )[0] + return event, heartbeat["heartbeat_at"] + + zone = DEPARTMENT_ZONES[heartbeat["department_id"]] + timestamp = heartbeat["heartbeat_at"].isoformat().replace("+00:00", "Z") + return { + "event_id": heartbeat["run_id"], + "task_id": heartbeat["run_id"], + "department_id": heartbeat["department_id"], + "department_label": zone["label"], + "project": heartbeat["project"], + "agent_id": heartbeat["agent_id"], + "role": zone["roles"][0], + "status": "active", + "updated_at": timestamp, + "next_step": "", + "evidence_count": 0, + "ephemeral": True, + "zone_id": zone["zone_id"], + }, heartbeat["heartbeat_at"] + + def department_campus_projection( events: object, *, + heartbeats: object = None, now: datetime, max_tasks: int = 3, owner_view: bool = False, @@ -547,46 +652,68 @@ def department_campus_projection( """Return a strict projection, optionally including validated owner fields.""" if not isinstance(events, list): return _empty_projection("unavailable", now) - if not events: - return _empty_projection("empty", now) - validated: list[tuple[int, dict[str, Any], datetime]] = [] + # Lifecycle rows remain history. They are inspected for staleness and + # terminal conflicts, but only a fresh exact heartbeat can create presence. saw_stale = False - for index, raw_event in enumerate(events): - event, validation_state, updated = _validated_event( + for raw_event in events: + _, validation_state, _ = _validated_event( raw_event, now=now, owner_view=owner_view, ) saw_stale = saw_stale or validation_state == "stale" - if event is not None and updated is not None: - validated.append((index, event, updated)) - if not validated: + + if not isinstance(heartbeats, list): return _empty_projection("stale" if saw_stale else "empty", now) - # Newest update wins for both event id and task+agent identity. Python's - # stable sort preserves the first source item when timestamps tie. - deduped: list[tuple[int, dict[str, Any], datetime]] = [] - seen_event_ids: set[str] = set() - seen_identities: set[tuple[str, str]] = set() - for item in sorted(validated, key=lambda candidate: candidate[2], reverse=True): - _, event, _ = item - identity = (event["task_id"], event["agent_id"]) - if event["event_id"] in seen_event_ids or identity in seen_identities: + # Any repeated canonical agent claim is ambiguous, even when one record is + # malformed. Fail that identity closed instead of choosing a winner. + identity_claims: dict[str, int] = {} + for raw_heartbeat in heartbeats: + if not isinstance(raw_heartbeat, dict): + continue + claimed_agent = raw_heartbeat.get("agent_id") + claimed_at = _parse_time(raw_heartbeat.get("heartbeat_at")) + claim_age = ( + (_utc_now(now) - claimed_at).total_seconds() + if claimed_at is not None + else None + ) + if ( + claimed_agent in _CAMPUS_AGENT_DEPARTMENTS + and raw_heartbeat.get("state") == "working" + and claim_age is not None + and 0 <= claim_age <= _HEARTBEAT_FRESH_SECONDS + ): + identity_claims[claimed_agent] = identity_claims.get(claimed_agent, 0) + 1 + + live: list[tuple[int, dict[str, Any], datetime]] = [] + for index, raw_heartbeat in enumerate(heartbeats): + heartbeat, validation_state = _validated_heartbeat(raw_heartbeat, now=now) + saw_stale = saw_stale or validation_state == "stale" + if heartbeat is None or identity_claims.get(heartbeat["agent_id"]) != 1: continue - seen_event_ids.add(event["event_id"]) - seen_identities.add(identity) - deduped.append(item) - deduped.sort(key=lambda item: item[0]) + event, updated = _heartbeat_event( + heartbeat, + events, + now=now, + owner_view=owner_view, + ) + if event is not None: + live.append((index, event, updated)) + + if not live: + return _empty_projection("stale" if saw_stale else "empty", now) requested = max_tasks if type(max_tasks) is int else 3 lane_limit = min(3, max(0, requested)) all_tasks: list[str] = [] - for _, event, _ in deduped: + for _, event, _ in live: if event["task_id"] not in all_tasks: all_tasks.append(event["task_id"]) visible_tasks = set(all_tasks[:lane_limit]) - visible_events = [event for _, event, _ in deduped if event["task_id"] in visible_tasks] + visible_events = [event for _, event, _ in live if event["task_id"] in visible_tasks] return { "state": "active" if visible_events else "empty", "generated_at": _utc_now(now).isoformat().replace("+00:00", "Z"), diff --git a/builder/jarvis-agent-pipeline b/builder/jarvis-agent-pipeline index b23052b..d8f0064 100755 --- a/builder/jarvis-agent-pipeline +++ b/builder/jarvis-agent-pipeline @@ -1,5 +1,6 @@ #!/usr/bin/env zsh set -euo pipefail +unsetopt BG_NICE export PATH="/opt/homebrew/bin:/usr/local/bin:$HOME/.local/bin:$PATH" @@ -17,6 +18,8 @@ if [[ -z "$TASK" ]]; then echo "usage: jarvis-agent-pipeline ''" >&2 exit 2 fi +PIPELINE_SCRIPT="${0:A}" +PROVIDER_WORKER_ROLE="${JARVIS_AGENT_PROVIDER_WORKER_ROLE:-}" PING_ROLE="${JARVIS_AGENT_PING_ROLE:-}" case "$PING_ROLE" in @@ -86,6 +89,14 @@ fi CODEX_MODEL="${JARVIS_CODEX_MODEL:-}" CODEX_REASONING_EFFORT="${JARVIS_CODEX_REASONING_EFFORT:-low}" CODEX_APPROVAL_POLICY="${JARVIS_CODEX_APPROVAL_POLICY:-never}" +JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS="${JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS:-15}" +if [[ "$JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS" != <-> ]]; then + JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS="15" +elif (( JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS < 5 )); then + JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS="5" +elif (( JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS > 30 )); then + JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS="30" +fi SUPERVISOR_PERMISSION="${JARVIS_SUPERVISOR_PERMISSION:-dontAsk}" BUILDER_PERMISSION="${JARVIS_BUILDER_PERMISSION:-auto}" @@ -220,6 +231,46 @@ print(json.dumps(event, ensure_ascii=False)) PY } +pixel_project_name() { + if [[ -n "$PING_ROLE" ]]; then + # A role-only availability ping does not identify the approved project + # selected by the runner, so it must not create live Campus presence. + echo "" + return 0 + fi + local project_key="${PROJECT_NAME:u}" + project_key="${project_key//_/ }" + project_key="${project_key//-/ }" + case "$project_key" in + "MAIN MANAGER") echo "MAIN MANAGER" ;; + "AI STUDIO") echo "AI STUDIO" ;; + "MY DICTIONARY") echo "MY DICTIONARY" ;; + "ACCOUNTABLE OS") echo "ACCOUNTABLE OS" ;; + "HEALTH OS") echo "HEALTH OS" ;; + "CONTEXT NEWS") echo "CONTEXT NEWS" ;; + "PIXELVERSE DASHBOARD"|"AGENT DASHBOARD") echo "PIXELVERSE DASHBOARD" ;; + "JARVIS") echo "JARVIS" ;; + "GITHUB HYGIENE") echo "GITHUB HYGIENE" ;; + "SKILLS LIBRARY") echo "SKILLS LIBRARY" ;; + "UNFINISHED STUFF") echo "UNFINISHED STUFF" ;; + "FINANCIAL OS") echo "FINANCIAL OS" ;; + *) echo "" ;; + esac +} + +pixel_agent_id_for_project() { + case "$1" in + "MAIN MANAGER") echo "COORDINATOR" ;; + "AI STUDIO") echo "RESEARCHER" ;; + "MY DICTIONARY"|"ACCOUNTABLE OS"|"HEALTH OS"|"CONTEXT NEWS") echo "BUILDER" ;; + "PIXELVERSE DASHBOARD") echo "DESIGNER" ;; + "JARVIS"|"GITHUB HYGIENE") echo "INFRASTRUCTURE" ;; + "SKILLS LIBRARY"|"UNFINISHED STUFF") echo "VAULT" ;; + "FINANCIAL OS") echo "ANALYST" ;; + *) echo "" ;; + esac +} + pixel_event() { if [[ ! -x "$PIXEL_EVENT" ]]; then return 0 @@ -230,8 +281,19 @@ pixel_event() { local event_phase="${4:-}" local event_handoff="${5:-}" local event_evidence="${6:-}" - - JARVIS_PIXEL_PROJECT_NAME="$PROJECT_NAME" \ + if [[ "$event_command" == "done" \ + && "${PROVIDER_TERMINAL_ALLOWED:-1}" != "1" ]]; then + return 0 + fi + local canonical_project_name="" + local canonical_agent_id="" + canonical_project_name="$(pixel_project_name)" + canonical_agent_id="$(pixel_agent_id_for_project "$canonical_project_name" "$event_role")" + + JARVIS_PIXEL_PROJECT_NAME="$canonical_project_name" \ + JARVIS_PIXEL_AGENT_ID="$canonical_agent_id" \ + JARVIS_PIXEL_RUN_ID="$RUN_ID" \ + JARVIS_PIXEL_SESSION_ID="$RUN_ID-$event_role" \ JARVIS_PIXEL_PROJECT_DIR="$PROJECT_DIR" \ JARVIS_PIXEL_TASK="$TASK_REDACTED" \ JARVIS_PIXEL_PHASE="$event_phase" \ @@ -248,9 +310,112 @@ pixel_cleanup() { "$PIXEL_CLEANUP" >/dev/null 2>&1 || true } -if [[ -z "$PING_ROLE" ]]; then - trap pixel_cleanup EXIT - pixel_cleanup +ACTIVE_PROVIDER_PID="" +ACTIVE_PROVIDER_PROCESS_GROUP="" +ACTIVE_HEARTBEAT_PID="" +ACTIVE_PROVIDER_OUTPUT="" +ACTIVE_PROVIDER_CONTEXT="" +ACTIVE_PROVIDER_ROLE="" +PROVIDER_SPAWN_CRITICAL="0" +PENDING_SIGNAL_STATUS="" +PROVIDER_TERMINAL_ALLOWED="1" +ROLE_PROVIDER_RESULT="" +SELF_HEAL_RESULT="" + +terminate_active_provider_group() { + if [[ "$ACTIVE_PROVIDER_PROCESS_GROUP" != <-> ]]; then + return 0 + fi + + local provider_pid="$ACTIVE_PROVIDER_PID" + local term_deadline=$(( SECONDS + 3 )) + local kill_deadline="" + if ! kill -TERM -- "-$ACTIVE_PROVIDER_PROCESS_GROUP" 2>/dev/null; then + kill -TERM "$provider_pid" 2>/dev/null || true + fi + while (( SECONDS < term_deadline )); do + if ! kill -0 -- "-$ACTIVE_PROVIDER_PROCESS_GROUP" 2>/dev/null \ + && ! kill -0 "$provider_pid" 2>/dev/null; then + break + fi + sleep 0.05 + done + if kill -0 -- "-$ACTIVE_PROVIDER_PROCESS_GROUP" 2>/dev/null \ + || kill -0 "$provider_pid" 2>/dev/null; then + kill -KILL -- "-$ACTIVE_PROVIDER_PROCESS_GROUP" 2>/dev/null || true + kill -KILL "$provider_pid" 2>/dev/null || true + fi + kill_deadline=$(( SECONDS + 3 )) + while (( SECONDS < kill_deadline )); do + if ! kill -0 -- "-$ACTIVE_PROVIDER_PROCESS_GROUP" 2>/dev/null \ + && ! kill -0 "$provider_pid" 2>/dev/null; then + break + fi + sleep 0.05 + done + if kill -0 -- "-$ACTIVE_PROVIDER_PROCESS_GROUP" 2>/dev/null \ + || kill -0 "$provider_pid" 2>/dev/null; then + PROVIDER_TERMINAL_ALLOWED="0" + return 1 + fi + if [[ "$provider_pid" == <-> ]]; then + wait "$ACTIVE_PROVIDER_PID" 2>/dev/null || true + fi + ACTIVE_PROVIDER_PID="" + ACTIVE_PROVIDER_PROCESS_GROUP="" +} + +cleanup_active_role() { + local cleanup_proven="1" + if [[ "$ACTIVE_HEARTBEAT_PID" == <-> ]]; then + kill "$ACTIVE_HEARTBEAT_PID" 2>/dev/null || true + wait "$ACTIVE_HEARTBEAT_PID" 2>/dev/null || true + fi + ACTIVE_HEARTBEAT_PID="" + terminate_active_provider_group || cleanup_proven="0" + if [[ -n "$ACTIVE_PROVIDER_OUTPUT" ]]; then + rm -f -- "$ACTIVE_PROVIDER_OUTPUT" 2>/dev/null || true + fi + if [[ -n "$ACTIVE_PROVIDER_CONTEXT" ]]; then + rm -rf -- "$ACTIVE_PROVIDER_CONTEXT" 2>/dev/null || true + fi + ACTIVE_PROVIDER_OUTPUT="" + ACTIVE_PROVIDER_CONTEXT="" + ACTIVE_PROVIDER_ROLE="" + [[ "$cleanup_proven" == "1" ]] +} + +pipeline_exit_cleanup() { + local interrupted_role="$ACTIVE_PROVIDER_ROLE" + local cleanup_proven="1" + cleanup_active_role || cleanup_proven="0" + if [[ -n "$interrupted_role" && "$cleanup_proven" == "1" \ + && "$PROVIDER_TERMINAL_ALLOWED" == "1" ]]; then + pixel_event done "$interrupted_role" "Provider interrupted" "Provider stopped" "Retry the exact task" + fi + if [[ -z "$PING_ROLE" ]]; then + pixel_cleanup + fi +} + +handle_pipeline_signal() { + local signal_status="$1" + if [[ "$PROVIDER_SPAWN_CRITICAL" == "1" ]]; then + PENDING_SIGNAL_STATUS="$signal_status" + return 0 + fi + exit "$signal_status" +} + +if [[ -z "$PROVIDER_WORKER_ROLE" ]]; then + trap pipeline_exit_cleanup EXIT + trap 'handle_pipeline_signal 129' HUP + trap 'handle_pipeline_signal 130' INT + trap 'handle_pipeline_signal 143' TERM + + if [[ -z "$PING_ROLE" ]]; then + pixel_cleanup + fi fi run_claude_role() { @@ -448,6 +613,71 @@ run_agent_role() { esac } +run_agent_role_with_heartbeat() { + local role="$1" + local permission="$2" + local system_prompt="$3" + local user_prompt="$4" + local provider_code=0 + local provider_output="" + local provider_context="" + local heartbeat_elapsed="$JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS" + + ROLE_PROVIDER_RESULT="" + PROVIDER_TERMINAL_ALLOWED="1" + PENDING_SIGNAL_STATUS="" + provider_output="$(mktemp "$REPORT_DIR/.provider-output.XXXXXX")" + ACTIVE_PROVIDER_OUTPUT="$provider_output" + provider_context="$(mktemp -d "$REPORT_DIR/.provider-worker.XXXXXX")" + ACTIVE_PROVIDER_CONTEXT="$provider_context" + ACTIVE_PROVIDER_ROLE="$role" + print -rn -- "$permission" > "$provider_context/permission" + print -rn -- "$system_prompt" > "$provider_context/system-prompt" + print -rn -- "$user_prompt" > "$provider_context/user-prompt" + chmod 700 "$provider_context" + chmod 600 "$provider_context/permission" "$provider_context/system-prompt" "$provider_context/user-prompt" + PROVIDER_SPAWN_CRITICAL="1" + JARVIS_AGENT_PROVIDER_WORKER_ROLE="$role" \ + JARVIS_AGENT_PROVIDER_WORKER_CONTEXT="$provider_context" \ + python3 -c 'import os, sys; os.setsid(); os.execvp(sys.argv[1], sys.argv[1:])' \ + zsh "$PIPELINE_SCRIPT" "$TASK" > "$provider_output" 2>&1 & + ACTIVE_PROVIDER_PID=$! + ACTIVE_PROVIDER_PROCESS_GROUP="$ACTIVE_PROVIDER_PID" + PROVIDER_SPAWN_CRITICAL="0" + if [[ "$PENDING_SIGNAL_STATUS" == <-> ]]; then + exit "$PENDING_SIGNAL_STATUS" + fi + + while kill -0 "$ACTIVE_PROVIDER_PID" 2>/dev/null; do + if (( heartbeat_elapsed >= JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS )); then + pixel_event heartbeat "$role" + heartbeat_elapsed=0 + fi + sleep 1 & + ACTIVE_HEARTBEAT_PID=$! + wait "$ACTIVE_HEARTBEAT_PID" 2>/dev/null || true + ACTIVE_HEARTBEAT_PID="" + (( heartbeat_elapsed += 1 )) + done + + set +e + wait "$ACTIVE_PROVIDER_PID" + provider_code=$? + set -e + if ! terminate_active_provider_group; then + PROVIDER_TERMINAL_ALLOWED="0" + exit 125 + fi + ACTIVE_PROVIDER_ROLE="" + + ROLE_PROVIDER_RESULT="$(<"$provider_output")" + rm -f "$provider_output" 2>/dev/null || true + rm -rf "$provider_context" 2>/dev/null || true + ACTIVE_PROVIDER_OUTPUT="" + ACTIVE_PROVIDER_CONTEXT="" + return "$provider_code" +} + role_output_failed() { local output="$1" print -r -- "$output" | grep -q '^ROLE_FAILED=' @@ -527,6 +757,7 @@ attempt_role_self_heal() { local can_retry="1" local failure_tail="" + SELF_HEAL_RESULT="" diagnosis="$(classify_role_failure "$failed_output")" failure_tail="$(print -r -- "$failed_output" | tail -60)" record_self_learning "$role" "$diagnosis" "detected" "$failure_tail" @@ -544,7 +775,7 @@ attempt_role_self_heal() { if [[ "$SELF_HEAL_ENABLED" != "1" || "$can_retry" != "1" || "$SELF_HEAL_MAX_RETRIES" -lt 1 ]]; then record_self_learning "$role" "$diagnosis" "not_retried" "$failure_tail" - print -r -- "$failed_output" + SELF_HEAL_RESULT="$failed_output" return 1 fi @@ -564,18 +795,19 @@ $failure_tail Original role prompt: $user_prompt" - retry_output="$(run_agent_role "$role" "$permission" "$retry_system" "$retry_prompt" || true)" + run_agent_role_with_heartbeat "$role" "$permission" "$retry_system" "$retry_prompt" || true + retry_output="$ROLE_PROVIDER_RESULT" append_section "$role Self-Heal Retry" "$retry_output" if role_output_failed "$retry_output"; then record_self_learning "$role" "$diagnosis" "retry_failed" "$(print -r -- "$retry_output" | tail -60)" - print -r -- "$retry_output" + SELF_HEAL_RESULT="$retry_output" return 1 fi record_self_learning "$role" "$diagnosis" "recovered" "$retry_output" append_section "$role" "$retry_output" - print -r -- "$retry_output" + SELF_HEAL_RESULT="$retry_output" } supervisor_needs_input() { @@ -1003,6 +1235,28 @@ finally: PY } +if [[ -n "$PROVIDER_WORKER_ROLE" ]]; then + PROVIDER_WORKER_CONTEXT="${JARVIS_AGENT_PROVIDER_WORKER_CONTEXT:-}" + if [[ ! -d "$PROVIDER_WORKER_CONTEXT" \ + || "${PROVIDER_WORKER_CONTEXT:A:h}" != "${REPORT_DIR:A}" \ + || "${PROVIDER_WORKER_CONTEXT:t}" != .provider-worker.* ]]; then + echo "ROLE_FAILED=$PROVIDER_WORKER_ROLE" + echo "EXIT_CODE=2" + echo + echo "Provider worker context is missing" + exit 2 + fi + PROVIDER_WORKER_PERMISSION="$(<"$PROVIDER_WORKER_CONTEXT/permission")" + PROVIDER_WORKER_SYSTEM="$(<"$PROVIDER_WORKER_CONTEXT/system-prompt")" + PROVIDER_WORKER_PROMPT="$(<"$PROVIDER_WORKER_CONTEXT/user-prompt")" + run_agent_role \ + "$PROVIDER_WORKER_ROLE" \ + "$PROVIDER_WORKER_PERMISSION" \ + "$PROVIDER_WORKER_SYSTEM" \ + "$PROVIDER_WORKER_PROMPT" + exit $? +fi + if [[ -n "$PING_ROLE" ]]; then ROUTE_LABEL="$PING_ROLE availability ping" else @@ -1182,12 +1436,14 @@ Budget policy: $BUDGET_POLICY 5. Критерии готовности" pixel_event start supervisor "Request: $TASK_REDACTED" "Defining acceptance criteria and routing plan" "Builder receives scope, risks, and checks" -SUPERVISOR_OUTPUT="$(run_agent_role supervisor "$SUPERVISOR_PERMISSION" "$SUPERVISOR_SYSTEM" "$SUPERVISOR_PROMPT" || true)" +run_agent_role_with_heartbeat supervisor "$SUPERVISOR_PERMISSION" "$SUPERVISOR_SYSTEM" "$SUPERVISOR_PROMPT" || true +SUPERVISOR_OUTPUT="$ROLE_PROVIDER_RESULT" append_section "Supervisor" "$SUPERVISOR_OUTPUT" pixel_event done supervisor "Acceptance criteria ready for Builder" "Handoff ready" "Builder receives Supervisor section from report" "Report: $REPORT" if role_output_failed "$SUPERVISOR_OUTPUT"; then - HEALED_OUTPUT="$(attempt_role_self_heal supervisor "$SUPERVISOR_PERMISSION" "$SUPERVISOR_SYSTEM" "$SUPERVISOR_PROMPT" "$SUPERVISOR_OUTPUT" || true)" + attempt_role_self_heal supervisor "$SUPERVISOR_PERMISSION" "$SUPERVISOR_SYSTEM" "$SUPERVISOR_PROMPT" "$SUPERVISOR_OUTPUT" || true + HEALED_OUTPUT="$SELF_HEAL_RESULT" if ! role_output_failed "$HEALED_OUTPUT"; then SUPERVISOR_OUTPUT="$HEALED_OUTPUT" pixel_event done supervisor "Supervisor recovered and produced criteria" "Self-healing complete" "Builder receives recovered Supervisor output" "Report: $REPORT" @@ -1228,12 +1484,14 @@ $SUPERVISOR_OUTPUT 4. Handoff для Tester" pixel_event start builder "Build/check: $TASK_REDACTED" "Reading context and producing implementation handoff" "Tester receives files, commands, and result" -BUILDER_OUTPUT="$(run_agent_role builder "$BUILDER_PERMISSION" "$BUILDER_SYSTEM" "$BUILDER_PROMPT" || true)" +run_agent_role_with_heartbeat builder "$BUILDER_PERMISSION" "$BUILDER_SYSTEM" "$BUILDER_PROMPT" || true +BUILDER_OUTPUT="$ROLE_PROVIDER_RESULT" append_section "Builder" "$BUILDER_OUTPUT" pixel_event done builder "Builder result ready for Tester" "Implementation/check complete" "Tester receives Builder section from report" "Report: $REPORT" if role_output_failed "$BUILDER_OUTPUT"; then - HEALED_OUTPUT="$(attempt_role_self_heal builder "$BUILDER_PERMISSION" "$BUILDER_SYSTEM" "$BUILDER_PROMPT" "$BUILDER_OUTPUT" || true)" + attempt_role_self_heal builder "$BUILDER_PERMISSION" "$BUILDER_SYSTEM" "$BUILDER_PROMPT" "$BUILDER_OUTPUT" || true + HEALED_OUTPUT="$SELF_HEAL_RESULT" if ! role_output_failed "$HEALED_OUTPUT"; then BUILDER_OUTPUT="$HEALED_OUTPUT" pixel_event done builder "Builder recovered and produced result" "Self-healing complete" "Tester receives recovered Builder output" "Report: $REPORT" @@ -1268,12 +1526,14 @@ $BUILDER_OUTPUT 4. Следующий шаг" pixel_event start tester "Verify: $TASK_REDACTED" "Checking request vs result and evidence" "Supervisor receives pass/fail verdict" -TESTER_OUTPUT="$(run_agent_role tester "$TESTER_PERMISSION" "$TESTER_SYSTEM" "$TESTER_PROMPT" || true)" +run_agent_role_with_heartbeat tester "$TESTER_PERMISSION" "$TESTER_SYSTEM" "$TESTER_PROMPT" || true +TESTER_OUTPUT="$ROLE_PROVIDER_RESULT" append_section "Tester" "$TESTER_OUTPUT" pixel_event done tester "Tester verdict ready" "Verification complete" "Supervisor receives Tester section from report" "Report: $REPORT" if role_output_failed "$TESTER_OUTPUT"; then - HEALED_OUTPUT="$(attempt_role_self_heal tester "$TESTER_PERMISSION" "$TESTER_SYSTEM" "$TESTER_PROMPT" "$TESTER_OUTPUT" || true)" + attempt_role_self_heal tester "$TESTER_PERMISSION" "$TESTER_SYSTEM" "$TESTER_PROMPT" "$TESTER_OUTPUT" || true + HEALED_OUTPUT="$SELF_HEAL_RESULT" if ! role_output_failed "$HEALED_OUTPUT"; then TESTER_OUTPUT="$HEALED_OUTPUT" pixel_event done tester "Tester recovered and produced verdict" "Self-healing complete" "Supervisor receives recovered Tester output" "Report: $REPORT" diff --git a/builder/jarvis-pixel-agent-event b/builder/jarvis-pixel-agent-event index 3f7bd2d..2a06505 100755 --- a/builder/jarvis-pixel-agent-event +++ b/builder/jarvis-pixel-agent-event @@ -11,10 +11,33 @@ const cwd = process.env.JARVIS_PROJECT_DIR || '/Users/pirajoke/pixel-agents'; const sessionPrefix = process.env.JARVIS_PIXEL_SESSION_PREFIX || 'jarvis-visual'; -const sessionId = `${sessionPrefix}-${role}`; +const explicitProjectName = (process.env.JARVIS_PIXEL_PROJECT_NAME || '').trim(); +const explicitAgentId = (process.env.JARVIS_PIXEL_AGENT_ID || '').trim(); +const explicitRunId = (process.env.JARVIS_PIXEL_RUN_ID || '').trim(); +const explicitSessionId = (process.env.JARVIS_PIXEL_SESSION_ID || '').trim(); +const hasExactHeartbeatIdentity = Boolean( + explicitProjectName && explicitAgentId && explicitRunId && explicitSessionId, +); +const sessionId = explicitSessionId || `${sessionPrefix}-${role}`; +const runId = explicitRunId || sessionId; const serverPath = path.join(os.homedir(), '.pixel-agents', 'server.json'); const detailsPath = path.join(os.homedir(), '.pixel-agents', 'jarvis-agent-details.json'); -const envProjectName = process.env.JARVIS_PIXEL_PROJECT_NAME || path.basename(cwd); +const detailsLockPath = `${detailsPath}.lock`; +const lockWaitTimeoutMs = 3000; +const lockStaleAfterMs = 10000; +const envProjectName = explicitProjectName || path.basename(cwd); +const defaultAgentIds = { + supervisor: 'COORDINATOR', + coordinator: 'COORDINATOR', + researcher: 'RESEARCHER', + builder: 'BUILDER', + tester: 'BUILDER', + designer: 'DESIGNER', + infrastructure: 'INFRASTRUCTURE', + vault: 'VAULT', + analyst: 'ANALYST', +}; +const agentId = explicitAgentId || defaultAgentIds[role] || role.toUpperCase(); const envProjectDir = process.env.JARVIS_PIXEL_PROJECT_DIR || process.env.JARVIS_PROJECT_DIR || cwd; const envTask = process.env.JARVIS_PIXEL_TASK || ''; const envPhase = process.env.JARVIS_PIXEL_PHASE || ''; @@ -24,7 +47,7 @@ const envAcceptance = process.env.JARVIS_PIXEL_ACCEPTANCE || ''; const envEvidence = process.env.JARVIS_PIXEL_EVIDENCE || ''; function usage() { - console.error('usage: jarvis-pixel-agent-event [message]'); + console.error('usage: jarvis-pixel-agent-event [message]'); process.exit(2); } @@ -34,9 +57,8 @@ function readConfig() { try { return JSON.parse(fs.readFileSync(serverPath, 'utf8')); } catch (error) { - console.error(`Pixel Agents server config not found: ${serverPath}`); - console.error(error instanceof Error ? error.message : String(error)); - process.exit(1); + const reason = error instanceof Error ? error.message : String(error); + throw new Error(`Pixel Agents server config not found: ${serverPath}: ${reason}`); } } @@ -48,6 +70,140 @@ function readJson(filePath, fallback) { } } +function writeJsonAtomically(filePath, value) { + fs.mkdirSync(path.dirname(filePath), { recursive: true }); + const tmpPath = `${filePath}.${process.pid}.${Date.now()}.tmp`; + try { + fs.writeFileSync(tmpPath, `${JSON.stringify(value, null, 2)}\n`, 'utf8'); + fs.renameSync(tmpPath, filePath); + } finally { + try { + fs.unlinkSync(tmpPath); + } catch (error) { + if (!(error instanceof Error) || !('code' in error) || error.code !== 'ENOENT') { + throw error; + } + } + } +} + +function sleepSync(milliseconds) { + const sleeper = new Int32Array(new SharedArrayBuffer(4)); + Atomics.wait(sleeper, 0, 0, milliseconds); +} + +function ticketOwnerIsAlive(owner) { + if (!owner || !Number.isInteger(owner.pid) || owner.pid < 1) return false; + try { + process.kill(owner.pid, 0); + return true; + } catch (error) { + return Boolean( + error instanceof Error && 'code' in error && error.code === 'EPERM', + ); + } +} + +function ticketOwner(ticketPath) { + try { + return JSON.parse(fs.readFileSync(ticketPath, 'utf8')); + } catch { + return null; + } +} + +function removeAbandonedTicket(ticketPath, now) { + let ticketStat = null; + try { + ticketStat = fs.statSync(ticketPath); + } catch (error) { + if (error instanceof Error && 'code' in error && error.code === 'ENOENT') { + return; + } + throw error; + } + + const owner = ticketOwner(ticketPath); + if (ticketOwnerIsAlive(owner) || now - ticketStat.mtimeMs < lockStaleAfterMs) { + return; + } + + // Ticket names are unique and never reused. Re-check the inode and owner + // token immediately before unlinking only this abandoned contender. + try { + const currentStat = fs.statSync(ticketPath); + const currentOwner = ticketOwner(ticketPath); + if ( + currentStat.ino !== ticketStat.ino + || currentOwner?.token !== owner?.token + || ticketOwnerIsAlive(currentOwner) + ) return; + fs.unlinkSync(ticketPath); + } catch (error) { + if (!(error instanceof Error) || !('code' in error) || error.code !== 'ENOENT') { + throw error; + } + } +} + +function withDetailsLock(callback) { + fs.mkdirSync(detailsLockPath, { recursive: true, mode: 0o700 }); + fs.chmodSync(detailsLockPath, 0o700); + const deadline = Date.now() + lockWaitTimeoutMs; + const lockToken = `${process.pid}-${Date.now()}-${Math.random()}`; + const orderKey = process.hrtime.bigint().toString().padStart(20, '0'); + const ticketName = `${orderKey}-${process.pid}-${Math.random().toString(16).slice(2)}.ticket`; + const ticketPath = path.join(detailsLockPath, ticketName); + const ticketFd = fs.openSync(ticketPath, 'wx', 0o600); + try { + fs.writeFileSync( + ticketFd, + JSON.stringify({ pid: process.pid, token: lockToken }), + 'utf8', + ); + } catch (error) { + try { + fs.unlinkSync(ticketPath); + } catch { + // A failed ticket write remains fail-closed and later expires as stale. + } + throw error; + } finally { + fs.closeSync(ticketFd); + } + + try { + while (true) { + const now = Date.now(); + const tickets = fs.readdirSync(detailsLockPath) + .filter((name) => name.endsWith('.ticket')) + .sort(); + for (const queuedTicket of tickets) { + if (queuedTicket !== ticketName) { + removeAbandonedTicket(path.join(detailsLockPath, queuedTicket), now); + } + } + const remaining = fs.readdirSync(detailsLockPath) + .filter((name) => name.endsWith('.ticket')) + .sort(); + if (remaining[0] === ticketName) break; + if (now >= deadline) { + throw new Error('heartbeat details lock timeout'); + } + sleepSync(10); + } + return callback(); + } finally { + try { + fs.unlinkSync(ticketPath); + } catch (error) { + if (!(error instanceof Error) || !('code' in error) || error.code !== 'ENOENT') { + throw error; + } + } + } +} + function stateForCommand(eventCommand) { switch (eventCommand) { case 'start': @@ -137,11 +293,31 @@ function defaultHandoff(eventCommand, eventRole) { return ''; } -function writeDetails(eventCommand, label) { +function writeDetailsUnlocked(eventCommand, label) { const now = new Date().toISOString(); const details = readJson(detailsPath, { updatedAt: null, agents: {} }); const agents = details.agents && typeof details.agents === 'object' ? details.agents : {}; const previous = agents[role] || {}; + if (eventCommand === 'heartbeat') { + if ( + !hasExactHeartbeatIdentity + || previous.state !== 'working' + || previous.project !== envProjectName + || previous.agentId !== agentId + || previous.runId !== runId + || previous.sessionId !== sessionId + ) { + return false; + } + agents[role] = { + ...previous, + updatedAt: now, + heartbeatAt: now, + }; + const next = {updatedAt: now, agents}; + writeJsonAtomically(detailsPath, next); + return true; + } const history = Array.isArray(previous.history) ? previous.history : []; const task = cleanTask(label) || previous.task || ''; const phase = envPhase || defaultPhase(eventCommand, role); @@ -152,6 +328,9 @@ function writeDetails(eventCommand, label) { const acceptanceCriteria = splitList(envAcceptance); const evidence = splitList(envEvidence); const shouldResetRunContext = eventCommand === 'spawn' || eventCommand === 'start'; + const heartbeatIdentity = hasExactHeartbeatIdentity + ? { agentId, runId, sessionId, heartbeatAt: now } + : { agentId: '', runId: '', sessionId, heartbeatAt: null }; const event = { at: now, type: eventCommand, @@ -166,7 +345,7 @@ function writeDetails(eventCommand, label) { agents[role] = { role, - sessionId, + ...heartbeatIdentity, updatedAt: now, state: stateForCommand(eventCommand), current: currentForCommand(eventCommand, label), @@ -194,10 +373,12 @@ function writeDetails(eventCommand, label) { agents, }; - fs.mkdirSync(path.dirname(detailsPath), { recursive: true }); - const tmpPath = `${detailsPath}.tmp`; - fs.writeFileSync(tmpPath, `${JSON.stringify(next, null, 2)}\n`, 'utf8'); - fs.renameSync(tmpPath, detailsPath); + writeJsonAtomically(detailsPath, next); + return true; +} + +function writeDetails(eventCommand, label) { + return withDetailsLock(() => writeDetailsUnlocked(eventCommand, label)); } async function post(config, payload) { @@ -254,39 +435,57 @@ async function wait(config) { } async function main() { - const config = readConfig(); const label = message || `${role}: ${command}`; switch (command) { - case 'spawn': - await toolStart(config, `${role}: online`); - await toolDone(config); - await wait(config); - break; - case 'start': - await toolStart(config, `${role}: ${label}`); - break; - case 'done': - await toolDone(config); - await post(config, base('Stop')); - break; - case 'wait': - await wait(config); - break; - case 'permission': - await post(config, base('PermissionRequest')); - break; - case 'end': - await post(config, { - ...base('SessionEnd'), - reason: 'exit', - }); - break; + case 'heartbeat': + if (!writeDetails(command, label)) { + throw new Error('heartbeat rejected: exact working identity is required'); + } + console.log(`${command} ${role} ${sessionId}`); + return; default: - usage(); + break; } - writeDetails(command, label); + let detailsCommand = command; + try { + const config = readConfig(); + switch (command) { + case 'spawn': + await toolStart(config, `${role}: online`); + await toolDone(config); + await wait(config); + break; + case 'start': + await toolStart(config, `${role}: ${label}`); + break; + case 'done': + await toolDone(config); + await post(config, base('Stop')); + break; + case 'wait': + await wait(config); + break; + case 'permission': + await post(config, base('PermissionRequest')); + break; + case 'end': + await post(config, { + ...base('SessionEnd'), + reason: 'exit', + }); + break; + default: + usage(); + } + } catch (error) { + if (command === 'start') detailsCommand = 'done'; + throw error; + } finally { + // A failed Pixel hook must never leave the durable heartbeat looking live. + writeDetails(detailsCommand, label); + } console.log(`${command} ${role} ${sessionId}`); } diff --git a/builder/tests/test_department_campus.py b/builder/tests/test_department_campus.py index 9affd21..e124269 100644 --- a/builder/tests/test_department_campus.py +++ b/builder/tests/test_department_campus.py @@ -49,6 +49,16 @@ "ephemeral", "zone_id", ) +_AUTO_HEARTBEATS = object() +_CANONICAL_BY_DEPARTMENT = { + "hq": ("MAIN MANAGER", "COORDINATOR"), + "sales": ("AI STUDIO", "RESEARCHER"), + "development": ("MY DICTIONARY", "BUILDER"), + "design": ("PIXELVERSE DASHBOARD", "DESIGNER"), + "infrastructure": ("JARVIS", "INFRASTRUCTURE"), + "internal": ("SKILLS LIBRARY", "VAULT"), + "finance": ("FINANCIAL OS", "ANALYST"), +} class DepartmentCampusContractTests(unittest.TestCase): @@ -79,13 +89,14 @@ def _event(self, department_id="development", **overrides): zone = self._zones()[department_id] roles = zone.get("roles", ()) self.assertTrue(roles, f"{department_id} must declare at least one public role") + project, agent_id = _CANONICAL_BY_DEPARTMENT[department_id] payload = { "event_id": "evt-001", "task_id": "task-001", "department_id": department_id, "department_label": zone["label"], - "project": "Public Project", - "agent_id": "agent-001", + "project": project, + "agent_id": agent_id, "role": roles[0], "status": "active", "updated_at": "2026-08-08T11:55:00Z", @@ -97,8 +108,46 @@ def _event(self, department_id="development", **overrides): payload.update(overrides) return payload - def _project(self, events, **kwargs): - return self.campus.department_campus_projection(events, now=NOW, **kwargs) + def _heartbeats(self, events): + if not isinstance(events, list): + return [] + records = [] + seen = set() + for event in events: + if not isinstance(event, dict) or event.get("status") not in {"active", "testing"}: + continue + identity = ( + event.get("project"), + event.get("agent_id"), + event.get("task_id"), + ) + if identity in seen: + continue + seen.add(identity) + records.append({ + "project": event.get("project"), + "agent_id": event.get("agent_id"), + "run_id": event.get("task_id"), + "session_id": f"session-{len(records) + 1}", + "state": "working", + "heartbeat_at": (NOW - timedelta(seconds=15)).isoformat(), + }) + return records + + def _project(self, events, *, heartbeats=_AUTO_HEARTBEATS, **kwargs): + heartbeat_records = self._heartbeats(events) if heartbeats is _AUTO_HEARTBEATS else heartbeats + try: + return self.campus.department_campus_projection( + events, + heartbeats=heartbeat_records, + now=NOW, + **kwargs, + ) + except TypeError as exc: + self.fail( + "RED: department_campus_projection must accept explicit heartbeats " + f"({exc})" + ) def _css(self): return (ASSETS_DIR / "style.css").read_text(encoding="utf-8") @@ -148,7 +197,7 @@ def test_ac_2_read_only_surface_exposes_no_owner_or_mutation_controls(self): self.assertRegex(script, r"fetch\([^)]*/api/manager/departments") self.assertNotRegex(script, r"fetch\([^)]*/api/manager/departments[^)]*\bPOST\b") - def test_ac_4_status_text_and_motion_allowlist_are_explicit_and_accessible(self): + def test_ac_4_status_text_is_accessible_but_lifecycle_rows_are_not_live_presence(self): expected = { "queued": "в очереди", "active": "работает", @@ -161,9 +210,19 @@ def test_ac_4_status_text_and_motion_allowlist_are_explicit_and_accessible(self) script = self._script() for raw, visible in expected.items(): with self.subTest(status=raw): - projected = self._project([self._event(status=raw)]) - self.assertEqual(projected["events"][0]["status"], raw) self.assertIn(visible, html + script) + try: + projected = self._project( + [self._event(status=raw)], + heartbeats=[], + ) + except TypeError as exc: + self.fail( + "RED: lifecycle truth requires the heartbeats projection " + f"contract ({exc})" + ) + self.assertNotEqual(projected["state"], "active") + self.assertEqual(projected["events"], []) self.assertIn("active", script) self.assertIn("testing", script) for nonmoving in ("queued", "waiting", "done", "failed"): @@ -171,19 +230,15 @@ def test_ac_4_status_text_and_motion_allowlist_are_explicit_and_accessible(self) self.assertIn("aria-live", html) def test_ac_5_ec_2_lane_limit_is_clamped_to_zero_through_three_with_honest_counts(self): + departments = ("hq", "sales", "development", "design", "infrastructure") events = [ - self._event( - event_id=f"evt-{index}", - task_id=f"task-{index}", - agent_id=f"agent-{index}", - ) - for index in range(5) + self._event(department_id=department_id, event_id=f"evt-{index}", task_id=f"task-{index}") + for index, department_id in enumerate(departments) ] support = self._event( - department_id="design", + department_id="internal", event_id="evt-support", task_id="task-0", - agent_id="agent-support", ) for requested, visible, omitted in ((-4, 0, 5), (0, 0, 5), (2, 2, 3), (99, 3, 2)): with self.subTest(max_tasks=requested): @@ -192,10 +247,10 @@ def test_ac_5_ec_2_lane_limit_is_clamped_to_zero_through_three_with_honest_count self.assertEqual(projected["omitted_task_count"], omitted) self.assertLessEqual(len({item["task_id"] for item in projected["events"]}), visible) capped = self._project(events + [support], max_tasks=99) - self.assertIn("agent-support", {item["agent_id"] for item in capped["events"]}) + self.assertIn("VAULT", {item["agent_id"] for item in capped["events"]}) def test_ac_6_ec_4_registry_mismatches_are_rejected_and_support_stays_in_own_zone(self): - valid_support = self._event(department_id="design", agent_id="support-agent") + valid_support = self._event(department_id="design") projected = self._project([valid_support]) self.assertEqual(projected["events"][0]["department_id"], "design") self.assertEqual(projected["events"][0]["zone_id"], self._zones()["design"]["zone_id"]) @@ -347,13 +402,13 @@ def test_ac_12_ec_3_ec_9_freshness_and_dedup_choose_newest_with_stable_ties(self duplicate = self._event() self.assertEqual(len(self._project([duplicate, dict(duplicate)])["events"]), 1) - older = self._event(project="Older", updated_at="2026-08-08T11:50:00Z") - newer = self._event(event_id="evt-new", project="Newer", updated_at="2026-08-08T11:59:00Z") - self.assertEqual(self._project([older, newer])["events"][0]["project"], "Newer") + older = self._event(event_id="evt-old", updated_at="2026-08-08T11:50:00Z") + newer = self._event(event_id="evt-new", updated_at="2026-08-08T11:59:00Z") + self.assertEqual(self._project([older, newer])["events"][0]["event_id"], "evt-new") - first = self._event(event_id="evt-first", project="First") - tied = self._event(event_id="evt-tied", project="Second") - self.assertEqual(self._project([first, tied])["events"][0]["project"], "First") + first = self._event(event_id="evt-first") + tied = self._event(event_id="evt-tied") + self.assertEqual(self._project([first, tied])["events"][0]["event_id"], "evt-first") def test_ac_13_projection_and_idle_render_start_no_external_work(self): with ( @@ -405,12 +460,12 @@ def test_ec_6_err_6_invalid_status_time_or_ephemeral_flag_never_moves_or_display self.assertEqual(self._project([event])["events"], []) def test_ec_8_err_3_safe_text_is_bounded_and_secret_or_identity_shapes_fail_closed(self): - long_unicode = "Проект-" + ("🚀" * 1000) - projected = self._project([self._event(project=long_unicode)]) + long_unicode = "Шаг-" + ("🚀" * 1000) + projected = self._project([self._event(next_step=long_unicode)]) self.assertEqual(len(projected["events"]), 1) - safe_project = projected["events"][0]["project"] - self.assertLess(len(safe_project), len(long_unicode)) - safe_project.encode("utf-8") + safe_step = projected["events"][0]["next_step"] + self.assertLess(len(safe_step), len(long_unicode)) + safe_step.encode("utf-8") unsafe = self._event( event_id="evt-secret", @@ -453,19 +508,19 @@ def test_ac_8_ec_8_private_uris_paths_credentials_and_keys_never_reach_public_te "aws_access": "AKIACAMPUSPROBE1234X", } - for field in ("project", "next_step"): - for shape, value in high_risk.items(): - with self.subTest(field=field, shape=shape): - projected = self._project([self._event(**{field: value})]) - rendered = repr(projected) - self.assertNotIn(value, rendered) - self.assertNotIn("CAMPUS_", rendered) + for shape, value in high_risk.items(): + with self.subTest(field="next_step", shape=shape): + projected = self._project([self._event(next_step=value)]) + rendered = repr(projected) + self.assertNotIn(value, rendered) + self.assertNotIn("CAMPUS_", rendered) + with self.subTest(field="project", shape=shape): + self.assertEqual(self._project([self._event(project=value)])["events"], []) safe_prose = "Document OAuth token rotation policy with the team" - for field in ("project", "next_step"): - with self.subTest(field=field, shape="safe_prose"): - projected = self._project([self._event(**{field: safe_prose})]) - self.assertEqual(projected["events"][0][field], safe_prose) + projected = self._project([self._event(next_step=safe_prose)]) + self.assertEqual(projected["events"][0]["next_step"], safe_prose) + self.assertEqual(self._project([self._event(project=safe_prose)])["events"], []) def test_ac_8_embedded_relative_and_schemeless_sensitive_text_is_neutral(self): sensitive_values = ( @@ -478,12 +533,11 @@ def test_ac_8_embedded_relative_and_schemeless_sensitive_text_is_neutral(self): "internal.example/private?token=CAMPUS_URL_TOKEN", ) neutral = "Недоступно в публичной сводке" - for field in ("project", "next_step"): - for value in sensitive_values: - with self.subTest(field=field, value=value[:18]): - projected = self._project([self._event(**{field: value})]) + for value in sensitive_values: + with self.subTest(field="next_step", value=value[:18]): + projected = self._project([self._event(next_step=value)]) self.assertEqual(len(projected["events"]), 1) - self.assertEqual(projected["events"][0][field], neutral) + self.assertEqual(projected["events"][0]["next_step"], neutral) rendered = repr(projected) for marker in ( "/Users/mark", @@ -494,12 +548,12 @@ def test_ac_8_embedded_relative_and_schemeless_sensitive_text_is_neutral(self): r"C:\Users\Mark", ): self.assertNotIn(marker, rendered) + with self.subTest(field="project", value=value[:18]): + self.assertEqual(self._project([self._event(project=value)])["events"], []) safe = "Open the documentation and review token rotation policy" - projected = self._project([ - self._event(project=safe, next_step=safe), - ]) - self.assertEqual(projected["events"][0]["project"], safe) + projected = self._project([self._event(next_step=safe)]) + self.assertEqual(projected["events"][0]["project"], "MY DICTIONARY") self.assertEqual(projected["events"][0]["next_step"], safe) def test_ac_8_generic_secret_shaped_identity_is_rejected(self): @@ -523,34 +577,34 @@ def test_ac_8_hidden_paths_json_jwt_are_neutral_while_public_help_and_safe_ids_s "eyJhbGciOiJIUzI1NiJ9.CAMPUS_JWT_PAYLOAD.CAMPUS_JWT_SIGNATURE", ) neutral = "Недоступно в публичной сводке" - for field in ("project", "next_step"): - for value in sensitive_values: - with self.subTest(field=field, value=value[:18]): - projected = self._project([self._event(**{field: value})]) + for value in sensitive_values: + with self.subTest(field="next_step", value=value[:18]): + projected = self._project([self._event(next_step=value)]) self.assertEqual(len(projected["events"]), 1) - self.assertEqual(projected["events"][0][field], neutral) + self.assertEqual(projected["events"][0]["next_step"], neutral) self.assertNotIn(value, repr(projected)) self.assertNotIn("CAMPUS_", repr(projected)) + with self.subTest(field="project", value=value[:18]): + self.assertEqual(self._project([self._event(project=value)])["events"], []) safe_prose = "Run /help for public documentation" with self.subTest(safe_case="public-help-prose"): safe_text_event = self._project([ - self._event(project=safe_prose, next_step=safe_prose), + self._event(next_step=safe_prose), ])["events"][0] - self.assertEqual(safe_text_event["project"], safe_prose) + self.assertEqual(safe_text_event["project"], "MY DICTIONARY") self.assertEqual(safe_text_event["next_step"], safe_prose) with self.subTest(safe_case="ordinary-token-secret-words-in-identifiers"): - safe_identity_projection = self._project([ - self._event( - task_id="design-token-audit", - agent_id="secret-santa-task", - ), - ]) + safe_identity_projection = self._project([self._event(task_id="design-token-audit")]) self.assertEqual(len(safe_identity_projection["events"]), 1) safe_event = safe_identity_projection["events"][0] self.assertEqual(safe_event["task_id"], "design-token-audit") - self.assertEqual(safe_event["agent_id"], "secret-santa-task") + self.assertEqual(safe_event["agent_id"], "BUILDER") + self.assertEqual( + self._project([self._event(agent_id="secret-santa-task")])["events"], + [], + ) def test_ac_8_release_matrix_redacts_private_mounts_hosts_auth_and_email(self): sensitive_values = ( @@ -566,14 +620,15 @@ def test_ac_8_release_matrix_redacts_private_mounts_hosts_auth_and_email(self): ) neutral = "Недоступно в публичной сводке" - for field in ("project", "next_step"): - for value, marker in sensitive_values: - with self.subTest(field=field, marker=marker): - projected = self._project([self._event(**{field: value})]) + for value, marker in sensitive_values: + with self.subTest(field="next_step", marker=marker): + projected = self._project([self._event(next_step=value)]) self.assertEqual(len(projected["events"]), 1) - self.assertEqual(projected["events"][0][field], neutral) + self.assertEqual(projected["events"][0]["next_step"], neutral) self.assertNotIn(value, repr(projected)) self.assertNotIn(marker, repr(projected)) + with self.subTest(field="project", marker=marker): + self.assertEqual(self._project([self._event(project=value)])["events"], []) def test_ac_8_err_3_common_api_key_shapes_drop_public_identity_fields(self): credentials = ( @@ -596,17 +651,15 @@ def test_ac_8_err_3_common_api_key_shapes_drop_public_identity_fields(self): def test_ac_8_ec_8_control_and_bidi_characters_are_neutralized(self): unsafe = "Safe text\x00\x1b[31m\u202eCAMPUS_BIDI\u2066" forbidden = ("\x00", "\x1b", "\u202e", "\u2066") - for field in ("project", "next_step"): - with self.subTest(field=field): - projected = self._project([self._event(**{field: unsafe})]) - rendered = repr(projected) - for character in forbidden: - self.assertNotIn(character, projected["events"][0][field]) + projected = self._project([self._event(next_step=unsafe)]) + for character in forbidden: + self.assertNotIn(character, projected["events"][0]["next_step"]) + self.assertEqual(self._project([self._event(project=unsafe)])["events"], []) def test_ec_8_unicode_bound_does_not_end_inside_a_zwj_grapheme(self): - long_unicode = "Проект-" + ("👩\u200d💻" * 1000) - projected = self._project([self._event(project=long_unicode)]) - safe_project = projected["events"][0]["project"] + long_unicode = "Шаг-" + ("👩\u200d💻" * 1000) + projected = self._project([self._event(next_step=long_unicode)]) + safe_project = projected["events"][0]["next_step"] self.assertLess(len(safe_project), len(long_unicode)) safe_project.encode("utf-8") @@ -660,7 +713,7 @@ def test_ac_15_ec_11_ec_12_direction_b_shell_has_boulevard_lanes_and_two_waypoin self.assertRegex(script, r"(?:slice\(\s*0\s*,\s*3\s*\)|length\s*<\s*3)") self.assertRegex(script, r"textContent\s*=\s*[^;]*task_id") - def test_ac_15_status_only_routes_to_department_test_lab_or_github(self): + def test_ac_15_only_verified_active_or_testing_status_can_create_live_routes(self): script = self._campus_script() css = self._css() @@ -669,18 +722,20 @@ def test_ac_15_status_only_routes_to_department_test_lab_or_github(self): for token in ( "active", "testing", - "done", - "queued", - "waiting", - "failed", "department", "test-lab", - "github-station", ): self.assertIn(token, script) self.assertRegex(script, r"active[\s\S]{0,240}department") self.assertRegex(script, r"testing[\s\S]{0,240}test-lab") - self.assertRegex(script, r"done[\s\S]{0,240}github-station") + self.assertRegex( + script, + r"(?:liveStatuses|movingStatuses)[\s\S]{0,160}active[\s\S]{0,80}testing", + ) + self.assertRegex( + script, + r"(?:liveStatuses|movingStatuses)\.includes\([^)]*event\.status", + ) self.assertRegex(script, r"active[\s\S]{0,160}testing[\s\S]{0,240}moving") self.assertRegex(css, r"data-campus-moving[^{}]*true") @@ -712,20 +767,18 @@ def test_ac_5_ac_15_support_agents_share_one_route_per_task_lane(self): self._event( event_id="evt-primary", task_id="task-shared", - agent_id="agent-primary", status="active", ), self._event( department_id="design", event_id="evt-support", task_id="task-shared", - agent_id="agent-support", status="active", ), self._event( + department_id="infrastructure", event_id="evt-second", task_id="task-second", - agent_id="agent-second", status="testing", ), ] @@ -750,30 +803,27 @@ def test_ac_15_mixed_status_task_route_uses_testing_precedence_and_one_lane(self self._event( event_id="evt-shared-done", task_id="task-shared", - agent_id="agent-done", status="done", ), self._event( event_id="evt-shared-active", task_id="task-shared", - agent_id="agent-active", status="active", ), self._event( department_id="design", event_id="evt-shared-testing", task_id="task-shared", - agent_id="agent-testing", status="testing", ), self._event( - event_id="evt-two", task_id="task-two", agent_id="agent-two", + department_id="infrastructure", event_id="evt-two", task_id="task-two", ), self._event( - event_id="evt-three", task_id="task-three", agent_id="agent-three", + department_id="internal", event_id="evt-three", task_id="task-three", ), self._event( - event_id="evt-four", task_id="task-four", agent_id="agent-four", + department_id="finance", event_id="evt-four", task_id="task-four", ), ] projected = self._project(events, max_tasks=99) diff --git a/builder/tests/test_department_campus_live_truth.py b/builder/tests/test_department_campus_live_truth.py new file mode 100644 index 0000000..4017c41 --- /dev/null +++ b/builder/tests/test_department_campus_live_truth.py @@ -0,0 +1,994 @@ +from __future__ import annotations + +from datetime import datetime, timedelta, timezone +import importlib +import importlib.util +import json +import os +from pathlib import Path +import re +import shutil +import signal +import subprocess +import tempfile +import time +import unittest + + +BUILDER_DIR = Path(__file__).resolve().parents[1] +PIXEL_EVENT_PATH = BUILDER_DIR / "jarvis-pixel-agent-event" +PIPELINE_PATH = BUILDER_DIR / "jarvis-agent-pipeline" +SERVER_PATH = BUILDER_DIR / "dashboard-server-m4.py" +NOW = datetime(2026, 8, 18, 6, 0, tzinfo=timezone.utc) + + +class DepartmentCampusLiveTruthTests(unittest.TestCase): + """RED contract for `.tdd/spec-real-agent-heartbeat-truth-v1.md`.""" + + @classmethod + def setUpClass(cls): + cls.campus = importlib.import_module("dashboard_builder.department_campus") + + def _event(self, *, status: str = "active", **overrides) -> dict: + zone = self.campus.DEPARTMENT_ZONES["design"] + event = { + "event_id": "run-pixelverse-36", + "task_id": "run-pixelverse-36", + "department_id": "design", + "department_label": zone["label"], + "project": "PIXELVERSE DASHBOARD", + "agent_id": "DESIGNER", + "role": "Designer", + "status": status, + "updated_at": "2026-08-18T05:59:50Z", + "next_step": "Проверить безопасный результат", + "evidence_count": 1, + "ephemeral": True, + "zone_id": zone["zone_id"], + } + event.update(overrides) + return event + + def _heartbeat(self, **overrides) -> dict: + heartbeat = { + "project": "PIXELVERSE DASHBOARD", + "agent_id": "DESIGNER", + "run_id": "run-pixelverse-36", + "session_id": "session-pixelverse-36", + "state": "working", + "heartbeat_at": "2026-08-18T05:59:45Z", + } + heartbeat.update(overrides) + return heartbeat + + def _project(self, events: object, heartbeats: object) -> dict: + try: + return self.campus.department_campus_projection( + events, + heartbeats=heartbeats, + now=NOW, + ) + except TypeError as exc: + self.fail( + "RED: department_campus_projection must accept the verified " + f"heartbeats source ({exc})" + ) + + def _pipeline_process_fixture( + self, + root: Path, + *, + fake_claude_source: str, + run_id: str, + ) -> tuple[dict[str, str], Path, Path, Path]: + home = root / "home" + scripts = home / "scripts" + report_dir = home / "reports" + project_dir = home / "project" + scripts.mkdir(parents=True) + report_dir.mkdir() + project_dir.mkdir() + child_pid_path = home / "provider-child.pid" + event_log = home / "pixel-events.log" + + fake_pixel = scripts / "jarvis-pixel-agent-event" + fake_pixel.write_text( + """#!/bin/sh +alive=no +if [ "${1:-}" = done ] && [ -s "$FAKE_PROVIDER_CHILD_PID" ]; then + child_pid="$(cat "$FAKE_PROVIDER_CHILD_PID")" + if kill -0 "$child_pid" 2>/dev/null; then alive=yes; fi +fi +printf '%s|%s|child_alive=%s\n' "${1:-}" "${2:-}" "$alive" >> "$FAKE_PIXEL_LOG" +""", + encoding="utf-8", + ) + fake_pixel.chmod(0o755) + + fake_claude = home / "fake-claude.py" + fake_claude.write_text(fake_claude_source, encoding="utf-8") + fake_claude.chmod(0o755) + + environment = os.environ.copy() + environment.update( + { + "HOME": str(home), + "JARVIS_AGENT_ENV_FILE": str(home / "missing.env"), + "JARVIS_PROJECT_DIR": str(project_dir), + "JARVIS_PROJECT_NAME": "PIXELVERSE DASHBOARD", + "JARVIS_AGENT_REPORT_DIR": str(report_dir), + "JARVIS_AGENT_RUN_ID": run_id, + "JARVIS_AGENT_PROVIDER": "claude", + "JARVIS_CLAUDE_BIN": str(fake_claude), + "JARVIS_AGENT_SELF_HEAL_ENABLED": "0", + "JARVIS_OMNI_APPROVAL_ENABLED": "0", + "JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS": "5", + "FAKE_PROVIDER_CHILD_PID": str(child_pid_path), + "FAKE_PIXEL_LOG": str(event_log), + } + ) + return environment, child_pid_path, event_log, report_dir + + def _wait_for_path(self, file_path: Path, timeout: float = 8) -> None: + deadline = time.monotonic() + timeout + while time.monotonic() < deadline and not file_path.is_file(): + time.sleep(0.02) + self.assertTrue(file_path.is_file(), f"timed out waiting for {file_path.name}") + + def _stop_pipeline_process(self, process: subprocess.Popen) -> None: + try: + os.killpg(process.pid, signal.SIGKILL) + except ProcessLookupError: + pass + try: + process.wait(timeout=2) + except subprocess.TimeoutExpired: + pass + if process.stdout is not None: + process.stdout.close() + if process.stderr is not None: + process.stderr.close() + + def _stop_recorded_provider_group(self, child_pid_path: Path) -> None: + if not child_pid_path.is_file(): + return + try: + child_pid = int(child_pid_path.read_text(encoding="utf-8")) + provider_group = os.getpgid(child_pid) + if provider_group != os.getpgrp(): + os.killpg(provider_group, signal.SIGKILL) + else: + os.kill(child_pid, signal.SIGKILL) + except (OSError, ProcessLookupError, ValueError): + pass + + def test_ac_1_ec_1_only_exactly_fresh_working_heartbeat_creates_live_presence(self): + boundary = self._heartbeat( + heartbeat_at=(NOW - timedelta(seconds=45)).isoformat() + ) + live = self._project([self._event()], [boundary]) + + self.assertEqual(live["state"], "active") + self.assertEqual(live["visible_task_count"], 1) + self.assertEqual(len(live["events"]), 1) + self.assertEqual(live["events"][0]["project"], "PIXELVERSE DASHBOARD") + self.assertEqual(live["events"][0]["agent_id"], "DESIGNER") + + rejected = ( + self._heartbeat( + heartbeat_at=(NOW - timedelta(seconds=45, microseconds=1)).isoformat() + ), + self._heartbeat(heartbeat_at=(NOW + timedelta(microseconds=1)).isoformat()), + self._heartbeat(state="idle"), + self._heartbeat(state="waiting"), + ) + for heartbeat in rejected: + with self.subTest(heartbeat=heartbeat): + projection = self._project([self._event()], [heartbeat]) + self.assertNotEqual(projection["state"], "active") + self.assertEqual(projection["visible_task_count"], 0) + self.assertEqual(projection["events"], []) + + def test_ac_1_fresh_heartbeat_without_lifecycle_row_synthesizes_safe_live_event(self): + projection = self._project([], [self._heartbeat()]) + + self.assertEqual(projection["state"], "active") + self.assertEqual(projection["visible_task_count"], 1) + self.assertEqual(len(projection["events"]), 1) + event = projection["events"][0] + self.assertEqual(event["event_id"], "run-pixelverse-36") + self.assertEqual(event["task_id"], "run-pixelverse-36") + self.assertEqual(event["project"], "PIXELVERSE DASHBOARD") + self.assertEqual(event["agent_id"], "DESIGNER") + self.assertEqual(event["status"], "active") + self.assertNotIn("session-pixelverse-36", repr(projection)) + + def test_ac_2_ec_2_bridge_lifecycle_without_heartbeat_is_history_not_presence(self): + for status in ("queued", "active", "done", "failed"): + with self.subTest(status=status): + projection = self._project([self._event(status=status)], []) + self.assertNotEqual(projection["state"], "active") + self.assertEqual(projection["visible_task_count"], 0) + self.assertEqual(projection["events"], []) + + for terminal_status in ("done", "failed"): + with self.subTest(terminal_status=terminal_status): + projection = self._project( + [self._event(status=terminal_status)], + [self._heartbeat()], + ) + self.assertNotEqual(projection["state"], "active") + self.assertEqual(projection["events"], []) + + def test_ac_3_unknown_mismatched_unsafe_or_incomplete_heartbeat_fails_closed(self): + invalid = ( + self._heartbeat(project="UNKNOWN PROJECT"), + self._heartbeat(agent_id="BUILDER"), + self._heartbeat(run_id="../../private/run"), + self._heartbeat(session_id="token.secret-12345678901234567890"), + {key: value for key, value in self._heartbeat().items() if key != "session_id"}, + self._heartbeat(heartbeat_at="not-a-time"), + ) + for heartbeat in invalid: + with self.subTest(heartbeat=heartbeat): + projection = self._project([self._event()], [heartbeat]) + self.assertNotEqual(projection["state"], "active") + self.assertEqual(projection["events"], []) + + def test_ac_3_private_heartbeat_fields_never_enter_public_projection(self): + heartbeat = self._heartbeat( + prompt="", + task="", + project_dir="/Users/mark/private/project", + tool_output="", + credential="ghp_abcdefghijklmnopqrstuvwxyz123456", + raw={"file": "/private/tmp/raw.json"}, + ) + + projection = self._project([self._event()], [heartbeat]) + + self.assertEqual(projection["state"], "active") + rendered = repr(projection) + for forbidden in ( + "", + "", + "/Users/mark/private/project", + "", + "ghp_abcdefghijklmnopqrstuvwxyz123456", + "/private/tmp/raw.json", + "session-pixelverse-36", + ): + self.assertNotIn(forbidden, rendered) + + def test_err_1_unreadable_or_malformed_heartbeat_storage_has_no_live_events(self): + malformed_sources = (None, "{bad json", {}, {"heartbeats": "not-a-list"}) + for heartbeats in malformed_sources: + with self.subTest(heartbeats=heartbeats): + projection = self._project([self._event()], heartbeats) + self.assertNotEqual(projection["state"], "active") + self.assertEqual(projection["events"], []) + + def test_ac_1_err_1_server_reads_heartbeat_file_and_fails_closed_when_unreadable(self): + spec = importlib.util.spec_from_file_location( + "dashboard_server_live_truth_red", SERVER_PATH + ) + self.assertIsNotNone(spec) + self.assertIsNotNone(spec.loader) + server = importlib.util.module_from_spec(spec) + spec.loader.exec_module(server) + bridge_task = { + "id": "run-pixelverse-36", + "status": "running", + "agent_role": "DESIGNER", + "project": "PIXELVERSE DASHBOARD", + "created_at": "2026-08-18T05:59:00Z", + "claimed_at": "2026-08-18T05:59:50Z", + } + + with tempfile.TemporaryDirectory() as temp_dir: + unreadable = Path(temp_dir) / "missing-heartbeats.json" + try: + projected = server._department_campus_payload( + {"tasks": [bridge_task]}, + heartbeat_path=unreadable, + now=NOW, + ) + except TypeError as exc: + self.fail( + "RED: server payload must read an injected heartbeat_path " + f"and fail closed ({exc})" + ) + + self.assertNotEqual(projected["state"], "active") + self.assertEqual(projected["events"], []) + + valid_path = Path(temp_dir) / "heartbeats.json" + valid_path.write_text( + json.dumps( + { + "updatedAt": "2026-08-18T05:59:45Z", + "agents": { + "designer": { + "project": "PIXELVERSE DASHBOARD", + "agentId": "DESIGNER", + "runId": "run-pixelverse-36", + "sessionId": "session-pixelverse-36", + "state": "working", + "heartbeatAt": "2026-08-18T05:59:45Z", + } + }, + } + ), + encoding="utf-8", + ) + live = server._department_campus_payload( + {"tasks": [bridge_task]}, + heartbeat_path=valid_path, + now=NOW, + ) + + self.assertEqual(live["state"], "active") + self.assertEqual(live["visible_task_count"], 1) + self.assertEqual(live["events"][0]["agent_id"], "DESIGNER") + + def test_err_2_duplicate_or_conflicting_agent_identity_fails_that_agent_closed(self): + duplicate = self._heartbeat( + run_id="run-pixelverse-duplicate", + session_id="session-pixelverse-duplicate", + ) + second_event = self._event( + event_id="run-pixelverse-duplicate", + task_id="run-pixelverse-duplicate", + ) + + projection = self._project( + [self._event(), second_event], + [self._heartbeat(), duplicate], + ) + + self.assertNotEqual(projection["state"], "active") + self.assertEqual(projection["visible_task_count"], 0) + self.assertEqual(projection["events"], []) + + def test_ac_4_all_persistent_residents_are_static_until_verified_live_replacement(self): + self.assertEqual(len(self.campus.CAMPUS_RESIDENTS), 7) + self.assertTrue( + all(resident["wandering"] is False for resident in self.campus.CAMPUS_RESIDENTS) + ) + html = self.campus.build_department_campus_html() + self.assertNotIn("is-wandering", html) + self.assertEqual(html.count("ожидает задач"), 7) + + def test_ac_5_heartbeat_command_refreshes_copy_without_pixel_tool_event(self): + with tempfile.TemporaryDirectory() as temp_dir: + home = Path(temp_dir) + details_path = home / ".pixel-agents" / "jarvis-agent-details.json" + details_path.parent.mkdir(parents=True) + before = { + "updatedAt": "2026-08-18T05:50:00.000Z", + "agents": { + "designer": { + "role": "designer", + "state": "working", + "current": "Safe public status", + "task": "Safe task copy", + "project": "PIXELVERSE DASHBOARD", + "agentId": "DESIGNER", + "runId": "run-pixelverse-36", + "sessionId": "session-pixelverse-36", + "heartbeatAt": "2026-08-18T05:50:00.000Z", + "history": [{"at": "2026-08-18T05:50:00.000Z", "type": "start"}], + } + }, + } + details_path.write_text(json.dumps(before), encoding="utf-8") + environment = os.environ.copy() + environment.update( + { + "HOME": str(home), + "JARVIS_PIXEL_PROJECT_NAME": "PIXELVERSE DASHBOARD", + "JARVIS_PIXEL_AGENT_ID": "DESIGNER", + "JARVIS_PIXEL_RUN_ID": "run-pixelverse-36", + "JARVIS_PIXEL_SESSION_ID": "session-pixelverse-36", + } + ) + + bundled_node = ( + Path.home() / ".local/node-v22.22.3-darwin-arm64/bin/node" + ) + node = shutil.which("node") or ( + str(bundled_node) if bundled_node.is_file() else None + ) + self.assertIsNotNone( + node, + "RED environment: Node is required to execute the producer contract", + ) + result = subprocess.run( + [node, str(PIXEL_EVENT_PATH), "heartbeat", "designer"], + text=True, + capture_output=True, + env=environment, + timeout=10, + check=False, + ) + + self.assertEqual(result.returncode, 0, result.stderr) + after = json.loads(details_path.read_text(encoding="utf-8")) + record = after["agents"]["designer"] + for field in ("state", "current", "task", "project", "agentId", "runId", "sessionId", "history"): + self.assertEqual(record[field], before["agents"]["designer"][field]) + self.assertNotEqual(record["heartbeatAt"], before["agents"]["designer"]["heartbeatAt"]) + self.assertFalse(details_path.with_suffix(".json.tmp").exists()) + + def test_ac_5_producer_and_pipeline_define_atomic_bounded_heartbeat_lifecycle(self): + producer = PIXEL_EVENT_PATH.read_text(encoding="utf-8") + pipeline = PIPELINE_PATH.read_text(encoding="utf-8") + + for field in ("agentId", "runId", "sessionId", "heartbeatAt"): + with self.subTest(producer_field=field): + self.assertIn(field, producer) + self.assertIn("case 'heartbeat'", producer) + self.assertRegex(producer, r"writeFileSync\([^\n]*tmp") + self.assertRegex(producer, r"renameSync\([^\n]*tmp") + + self.assertIn("JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS", pipeline) + self.assertIn("pixel_event heartbeat", pipeline) + self.assertRegex(pipeline, r"while[^\n]*(?:kill -0|provider)") + self.assertRegex(pipeline, r"wait[^\n]*(?:heartbeat|HEARTBEAT)") + self.assertRegex(pipeline, r"pixel_event (?:done|end)") + + def test_ac_6_ec_3_two_independent_heartbeat_writers_preserve_both_agents(self): + bundled_node = Path.home() / ".local/node-v22.22.3-darwin-arm64/bin/node" + node = shutil.which("node") or ( + str(bundled_node) if bundled_node.is_file() else None + ) + self.assertIsNotNone( + node, + "RED environment: Node is required to execute the producer contract", + ) + + with tempfile.TemporaryDirectory() as temp_dir: + home = Path(temp_dir) + details_path = home / ".pixel-agents" / "jarvis-agent-details.json" + details_path.parent.mkdir(parents=True) + old_heartbeat = "2026-08-18T05:50:00.000Z" + details_path.write_text( + json.dumps( + { + "updatedAt": old_heartbeat, + "agents": { + "designer": { + "role": "designer", + "state": "working", + "project": "PIXELVERSE DASHBOARD", + "agentId": "DESIGNER", + "runId": "run-designer", + "sessionId": "session-designer", + "heartbeatAt": old_heartbeat, + }, + "infrastructure": { + "role": "infrastructure", + "state": "working", + "project": "JARVIS", + "agentId": "INFRASTRUCTURE", + "runId": "run-infrastructure", + "sessionId": "session-infrastructure", + "heartbeatAt": old_heartbeat, + }, + }, + } + ), + encoding="utf-8", + ) + + # Force both unprotected writers to finish their read before either + # can write. A correct inter-process lock serializes the reads; the + # first writer's bounded barrier then expires and releases the lock. + barrier_dir = home / "race-barrier" + barrier_dir.mkdir() + preload = home / "race-read-barrier.mjs" + preload.write_text( + """ +import fs from 'node:fs'; +import path from 'node:path'; + +const originalRead = fs.readFileSync.bind(fs); +const sleeper = new Int32Array(new SharedArrayBuffer(4)); +fs.readFileSync = function (filePath, ...args) { + const value = originalRead(filePath, ...args); + if (path.resolve(String(filePath)) !== path.resolve(process.env.RACE_DETAILS_PATH)) { + return value; + } + fs.writeFileSync(path.join(process.env.RACE_BARRIER_DIR, process.env.RACE_MARKER), 'read'); + const deadline = Date.now() + 600; + while (Date.now() < deadline) { + if (fs.readdirSync(process.env.RACE_BARRIER_DIR).length >= 2) break; + Atomics.wait(sleeper, 0, 0, 5); + } + return value; +}; +""".strip(), + encoding="utf-8", + ) + + processes = [] + identities = ( + ( + "designer", + "PIXELVERSE DASHBOARD", + "DESIGNER", + "run-designer", + "session-designer", + ), + ( + "infrastructure", + "JARVIS", + "INFRASTRUCTURE", + "run-infrastructure", + "session-infrastructure", + ), + ) + for role, project, agent_id, run_id, session_id in identities: + environment = os.environ.copy() + environment.update( + { + "HOME": str(home), + "NODE_OPTIONS": f"--import={preload}", + "RACE_DETAILS_PATH": str(details_path), + "RACE_BARRIER_DIR": str(barrier_dir), + "RACE_MARKER": role, + "JARVIS_PIXEL_PROJECT_NAME": project, + "JARVIS_PIXEL_AGENT_ID": agent_id, + "JARVIS_PIXEL_RUN_ID": run_id, + "JARVIS_PIXEL_SESSION_ID": session_id, + } + ) + processes.append( + subprocess.Popen( + [node, str(PIXEL_EVENT_PATH), "heartbeat", role], + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=environment, + ) + ) + + results = [process.communicate(timeout=5) for process in processes] + for process, (_, stderr) in zip(processes, results): + self.assertEqual(process.returncode, 0, stderr) + + after = json.loads(details_path.read_text(encoding="utf-8")) + self.assertEqual(set(after["agents"]), {"designer", "infrastructure"}) + for role in ("designer", "infrastructure"): + with self.subTest(role=role): + self.assertNotEqual( + after["agents"][role]["heartbeatAt"], + old_heartbeat, + "one concurrent read-modify-write discarded the other agent", + ) + + def test_ac_6_abandoned_lock_is_bounded_and_never_loses_existing_agents(self): + producer = PIXEL_EVENT_PATH.read_text(encoding="utf-8") + self.assertIn("lock", producer.lower(), "producer has no lock contract") + self.assertIsNotNone( + re.search(r"(?:['\"]wx['\"]|O_EXCL|mkdirSync)", producer), + "producer lock acquisition is not exclusive", + ) + self.assertIsNotNone( + re.search(r"(?i)(?:deadline|timeout|maxWait|lockWait)", producer), + "producer lock acquisition is not bounded", + ) + self.assertIsNotNone( + re.search(r"(?s)finally.*(?:unlinkSync|rmdirSync|rmSync)", producer), + "producer lock release is not guaranteed", + ) + + def test_ac_6_stale_recovery_never_removes_a_concurrent_replacement_lock(self): + bundled_node = Path.home() / ".local/node-v22.22.3-darwin-arm64/bin/node" + node = shutil.which("node") or ( + str(bundled_node) if bundled_node.is_file() else None + ) + self.assertIsNotNone(node, "Node is required for stale-lock recovery smoke") + + with tempfile.TemporaryDirectory() as temp_dir: + home = Path(temp_dir) + details_path = home / ".pixel-agents" / "jarvis-agent-details.json" + lock_path = Path(f"{details_path}.lock") + details_path.parent.mkdir(parents=True) + original_heartbeat = "2026-08-18T05:50:00.000Z" + before = { + "updatedAt": original_heartbeat, + "agents": { + "designer": { + "role": "designer", + "state": "working", + "project": "PIXELVERSE DASHBOARD", + "agentId": "DESIGNER", + "runId": "run-designer", + "sessionId": "session-designer", + "heartbeatAt": original_heartbeat, + } + }, + } + details_path.write_text( + json.dumps(before), + encoding="utf-8", + ) + lock_path.mkdir() + (lock_path / "owner.json").write_text( + json.dumps({"pid": 2_000_000_000, "token": "stale-owner"}), + encoding="utf-8", + ) + stale_time = time.time() - 30 + os.utime(lock_path, (stale_time, stale_time)) + + swap_marker = home / "replacement-installed" + preload = home / "swap-before-recursive-remove.mjs" + preload.write_text( + """ +import fs from 'node:fs'; +import path from 'node:path'; + +const originalRm = fs.rmSync.bind(fs); +let swapped = false; +fs.rmSync = function (target, ...args) { + if (!swapped && path.resolve(String(target)) === path.resolve(process.env.RACE_LOCK_PATH)) { + swapped = true; + fs.renameSync(process.env.RACE_LOCK_PATH, `${process.env.RACE_LOCK_PATH}.displaced`); + fs.mkdirSync(process.env.RACE_LOCK_PATH); + fs.writeFileSync( + path.join(process.env.RACE_LOCK_PATH, 'owner.json'), + JSON.stringify({pid: Number(process.env.RACE_REPLACEMENT_PID), token: 'replacement-owner'}), + ); + fs.writeFileSync(process.env.RACE_SWAP_MARKER, 'installed'); + } + return originalRm(target, ...args); +}; +""".strip(), + encoding="utf-8", + ) + environment = os.environ.copy() + environment.update( + { + "HOME": str(home), + "NODE_OPTIONS": f"--import={preload}", + "RACE_LOCK_PATH": str(lock_path), + "RACE_SWAP_MARKER": str(swap_marker), + "RACE_REPLACEMENT_PID": str(os.getpid()), + "JARVIS_PIXEL_PROJECT_NAME": "PIXELVERSE DASHBOARD", + "JARVIS_PIXEL_AGENT_ID": "DESIGNER", + "JARVIS_PIXEL_RUN_ID": "run-designer", + "JARVIS_PIXEL_SESSION_ID": "session-designer", + } + ) + result = subprocess.run( + [node, str(PIXEL_EVENT_PATH), "heartbeat", "designer"], + text=True, + capture_output=True, + env=environment, + timeout=6, + check=False, + ) + + after = json.loads(details_path.read_text(encoding="utf-8")) + if swap_marker.is_file(): + self.assertTrue( + lock_path.is_dir(), + "stale recovery recursively deleted the concurrently replaced lock", + ) + replacement = json.loads( + (lock_path / "owner.json").read_text(encoding="utf-8") + ) + self.assertEqual(replacement["token"], "replacement-owner") + elif result.returncode == 0: + self.assertNotEqual( + after["agents"]["designer"]["heartbeatAt"], + original_heartbeat, + "safe stale-lock recovery returned success without the heartbeat", + ) + else: + self.assertEqual( + after, + before, + "bounded fail-closed recovery partially changed heartbeat storage", + ) + + def test_ac_7_pipeline_declares_process_group_cleanup_contract(self): + pipeline = PIPELINE_PATH.read_text(encoding="utf-8") + for pattern, message in ( + (r"ACTIVE_PROVIDER_(?:PGID|PROCESS_GROUP)", "provider group is not tracked"), + (r"kill[^\n]*(?:PGID|PROCESS_GROUP)", "provider group is not terminated"), + ( + r"wait[^\n]*(?:PGID|PROCESS_GROUP|ACTIVE_PROVIDER)", + "provider group cleanup is not waited", + ), + ): + with self.subTest(pattern=pattern): + self.assertIsNotNone(re.search(pattern, pipeline), message) + + @unittest.skipUnless(shutil.which("zsh"), "zsh is required for process-tree smoke") + def test_ac_7_ec_4_signal_waits_for_provider_process_group_before_done(self): + with tempfile.TemporaryDirectory() as temp_dir: + home = Path(temp_dir) + scripts = home / "scripts" + report_dir = home / "reports" + project_dir = home / "project" + scripts.mkdir() + report_dir.mkdir() + project_dir.mkdir() + child_pid_path = home / "provider-child.pid" + event_log = home / "pixel-events.log" + + fake_pixel = scripts / "jarvis-pixel-agent-event" + fake_pixel.write_text( + """#!/bin/sh +alive=no +if [ "${1:-}" = done ] && [ -s "$FAKE_PROVIDER_CHILD_PID" ]; then + child_pid="$(cat "$FAKE_PROVIDER_CHILD_PID")" + if kill -0 "$child_pid" 2>/dev/null; then alive=yes; fi +fi +printf '%s|%s|child_alive=%s\n' "${1:-}" "${2:-}" "$alive" >> "$FAKE_PIXEL_LOG" +""", + encoding="utf-8", + ) + fake_pixel.chmod(0o755) + + fake_claude = home / "fake-claude.py" + fake_claude.write_text( + """#!/usr/bin/env python3 +import os +import subprocess +import sys + +child = subprocess.Popen([ + sys.executable, + '-c', + 'import os, time; open(os.environ["FAKE_PROVIDER_CHILD_PID"], "w").write(str(os.getpid())); time.sleep(60)', +]) +child.wait() +print('provider complete') +""", + encoding="utf-8", + ) + fake_claude.chmod(0o755) + + environment = os.environ.copy() + environment.update( + { + "HOME": str(home), + "JARVIS_AGENT_ENV_FILE": str(home / "missing.env"), + "JARVIS_PROJECT_DIR": str(project_dir), + "JARVIS_PROJECT_NAME": "PIXELVERSE DASHBOARD", + "JARVIS_AGENT_REPORT_DIR": str(report_dir), + "JARVIS_AGENT_RUN_ID": "signal-cleanup-red", + "JARVIS_AGENT_PROVIDER": "claude", + "JARVIS_CLAUDE_BIN": str(fake_claude), + "JARVIS_AGENT_SELF_HEAL_ENABLED": "0", + "JARVIS_OMNI_APPROVAL_ENABLED": "0", + "JARVIS_PIXEL_HEARTBEAT_INTERVAL_SECONDS": "5", + "FAKE_PROVIDER_CHILD_PID": str(child_pid_path), + "FAKE_PIXEL_LOG": str(event_log), + } + ) + process = subprocess.Popen( + ["zsh", str(PIPELINE_PATH), "исправь конкретную ошибку dashboard"], + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=environment, + start_new_session=True, + ) + try: + deadline = time.monotonic() + 8 + while time.monotonic() < deadline and not child_pid_path.is_file(): + if process.poll() is not None: + break + time.sleep(0.02) + self.assertTrue( + child_pid_path.is_file(), + "fake provider descendant never started", + ) + + process.send_signal(signal.SIGTERM) + process.wait(timeout=8) + + events = event_log.read_text(encoding="utf-8") + terminal_events = [ + line for line in events.splitlines() if line.startswith("done|") + ] + self.assertTrue(terminal_events, events) + self.assertTrue( + all("child_alive=no" in line for line in terminal_events), + "terminal Pixel done was emitted while provider work survived", + ) + self.assertEqual( + list(report_dir.glob(".provider-output.*")), + [], + "interrupted provider output must not remain on disk", + ) + finally: + try: + os.killpg(process.pid, signal.SIGKILL) + except ProcessLookupError: + pass + try: + process.wait(timeout=2) + except subprocess.TimeoutExpired: + pass + if process.stdout is not None: + process.stdout.close() + if process.stderr is not None: + process.stderr.close() + + @unittest.skipUnless(shutil.which("zsh"), "zsh is required for process-tree smoke") + def test_ac_7_early_signal_after_spawn_cannot_escape_unregistered_provider(self): + with tempfile.TemporaryDirectory() as temp_dir: + root = Path(temp_dir) + environment, child_pid_path, event_log, report_dir = ( + self._pipeline_process_fixture( + root, + run_id="early-signal-red", + fake_claude_source="""#!/usr/bin/env python3 +import os, subprocess, sys +child = subprocess.Popen([ + sys.executable, '-c', + 'import os, time; open(os.environ["FAKE_PROVIDER_CHILD_PID"], "w").write(str(os.getpid())); time.sleep(60)', +]) +child.wait() +""", + ) + ) + spawn_marker = root / "spawned-before-registration" + environment["EARLY_PROVIDER_SPAWN_FILE"] = str(spawn_marker) + source = PIPELINE_PATH.read_text(encoding="utf-8") + needle = ' zsh "$PIPELINE_SCRIPT" "$TASK" > "$provider_output" 2>&1 &\n ACTIVE_PROVIDER_PID=$!' + injected = ( + ' zsh "$PIPELINE_SCRIPT" "$TASK" > "$provider_output" 2>&1 &\n' + ' while [[ ! -s "$FAKE_PROVIDER_CHILD_PID" ]]; do sleep 0.01; done\n' + ' print -rn -- "$!" > "$EARLY_PROVIDER_SPAWN_FILE"\n' + ' sleep 5\n' + ' ACTIVE_PROVIDER_PID=$!' + ) + self.assertEqual(source.count(needle), 1, "provider spawn seam changed") + test_pipeline = root / "jarvis-agent-pipeline-early-signal" + test_pipeline.write_text(source.replace(needle, injected), encoding="utf-8") + test_pipeline.chmod(0o755) + + process = subprocess.Popen( + ["zsh", str(test_pipeline), "исправь конкретную ошибку dashboard"], + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=environment, + start_new_session=True, + ) + try: + self._wait_for_path(spawn_marker) + process.send_signal(signal.SIGTERM) + process.wait(timeout=8) + events = event_log.read_text(encoding="utf-8") + terminal_events = [ + line for line in events.splitlines() if line.startswith("done|") + ] + self.assertTrue(terminal_events, events) + self.assertTrue( + all("child_alive=no" in line for line in terminal_events), + "signal landed before provider PID/group registration and work escaped", + ) + self.assertEqual(list(report_dir.glob(".provider-*")), []) + finally: + self._stop_recorded_provider_group(child_pid_path) + self._stop_pipeline_process(process) + + @unittest.skipUnless(shutil.which("zsh"), "zsh is required for process-tree smoke") + def test_ac_7_normal_leader_exit_cannot_leave_descendant_live_before_done(self): + with tempfile.TemporaryDirectory() as temp_dir: + root = Path(temp_dir) + environment, child_pid_path, event_log, _ = self._pipeline_process_fixture( + root, + run_id="leader-exit-red", + fake_claude_source="""#!/usr/bin/env python3 +import os, subprocess, sys, time +subprocess.Popen([ + sys.executable, '-c', + 'import os, time; open(os.environ["FAKE_PROVIDER_CHILD_PID"], "w").write(str(os.getpid())); time.sleep(60)', +], stdin=subprocess.DEVNULL, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) +deadline = time.monotonic() + 3 +while time.monotonic() < deadline and not os.path.exists(os.environ['FAKE_PROVIDER_CHILD_PID']): + time.sleep(0.01) +print('provider leader exited normally') +""", + ) + process = subprocess.Popen( + ["zsh", str(PIPELINE_PATH), "исправь конкретную ошибку dashboard"], + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=environment, + start_new_session=True, + ) + try: + self._wait_for_path(child_pid_path) + deadline = time.monotonic() + 8 + terminal_events = [] + while time.monotonic() < deadline: + if event_log.is_file(): + terminal_events = [ + line + for line in event_log.read_text(encoding="utf-8").splitlines() + if line.startswith("done|") + ] + if terminal_events: + break + time.sleep(0.02) + self.assertTrue(terminal_events, "pipeline never emitted terminal done") + self.assertIn( + "child_alive=no", + terminal_events[0], + "provider leader exited but its working descendant survived done", + ) + finally: + self._stop_recorded_provider_group(child_pid_path) + self._stop_pipeline_process(process) + + @unittest.skipUnless(shutil.which("zsh"), "zsh is required for process-tree smoke") + def test_ac_7_cleanup_confirms_group_disappearance_after_kill_before_done(self): + with tempfile.TemporaryDirectory() as temp_dir: + root = Path(temp_dir) + environment, child_pid_path, event_log, _ = self._pipeline_process_fixture( + root, + run_id="post-kill-red", + fake_claude_source="""#!/usr/bin/env python3 +import os, subprocess, sys +child = subprocess.Popen([ + sys.executable, '-c', + 'import os, signal, time; signal.signal(signal.SIGTERM, signal.SIG_IGN); open(os.environ["FAKE_PROVIDER_CHILD_PID"], "w").write(str(os.getpid())); time.sleep(60)', +]) +child.wait() +""", + ) + source = PIPELINE_PATH.read_text(encoding="utf-8") + seam = "unsetopt BG_NICE\n" + delayed_kill = """unsetopt BG_NICE +kill() { + if [[ "${1:-}" == "-KILL" ]]; then + ( sleep 1; command kill "$@" ) & + return 0 + fi + command kill "$@" +} +""" + self.assertEqual(source.count(seam), 1, "kill seam changed") + test_pipeline = root / "jarvis-agent-pipeline-delayed-kill" + test_pipeline.write_text( + source.replace(seam, delayed_kill), encoding="utf-8" + ) + test_pipeline.chmod(0o755) + process = subprocess.Popen( + ["zsh", str(test_pipeline), "исправь конкретную ошибку dashboard"], + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + env=environment, + start_new_session=True, + ) + try: + self._wait_for_path(child_pid_path) + process.send_signal(signal.SIGTERM) + process.wait(timeout=10) + events = event_log.read_text(encoding="utf-8") + terminal_events = [ + line for line in events.splitlines() if line.startswith("done|") + ] + self.assertTrue(terminal_events, events) + self.assertTrue( + all("child_alive=no" in line for line in terminal_events), + "cleanup emitted done before delayed KILL removed the group", + ) + finally: + self._stop_recorded_provider_group(child_pid_path) + self._stop_pipeline_process(process) + + +if __name__ == "__main__": + unittest.main() diff --git a/builder/tests/test_department_campus_projects.py b/builder/tests/test_department_campus_projects.py index 0daa38d..4771979 100644 --- a/builder/tests/test_department_campus_projects.py +++ b/builder/tests/test_department_campus_projects.py @@ -61,8 +61,37 @@ def _event(self, project="MY DICTIONARY", department_id="development", agent_id= event.update(overrides) return event - def _projection(self, events): - return self.campus.department_campus_projection(events, now=NOW) + def _projection(self, events, *, heartbeats=None): + heartbeat_records = heartbeats + if heartbeat_records is None: + heartbeat_records = [] + seen = set() + for event in events if isinstance(events, list) else []: + if not isinstance(event, dict) or event.get("status") not in {"active", "testing"}: + continue + identity = (event.get("project"), event.get("agent_id"), event.get("task_id")) + if identity in seen: + continue + seen.add(identity) + heartbeat_records.append({ + "project": event.get("project"), + "agent_id": event.get("agent_id"), + "run_id": event.get("task_id"), + "session_id": f"session-project-{len(heartbeat_records) + 1}", + "state": "working", + "heartbeat_at": "2026-08-10T11:59:45Z", + }) + try: + return self.campus.department_campus_projection( + events, + heartbeats=heartbeat_records, + now=NOW, + ) + except TypeError as exc: + self.fail( + "RED: project projection must accept explicit heartbeats " + f"({exc})" + ) def test_ac_1_ec_4_registry_is_the_exact_canonical_public_project_set(self): registry = self._registry() diff --git a/builder/tests/test_department_campus_server.py b/builder/tests/test_department_campus_server.py index 25d6948..6f11261 100644 --- a/builder/tests/test_department_campus_server.py +++ b/builder/tests/test_department_campus_server.py @@ -3,7 +3,9 @@ from datetime import datetime, timedelta, timezone import importlib import importlib.util +import json from pathlib import Path +import tempfile import unittest from unittest.mock import patch @@ -59,15 +61,15 @@ def setUp(self): f"({type(exc).__name__}: {exc})" ) - def _event(self, department_id="development", **overrides): + def _event(self, department_id="design", **overrides): zone = self.campus.DEPARTMENT_ZONES[department_id] payload = { "event_id": "evt-server-1", "task_id": "task-server-1", "department_id": department_id, "department_label": zone["label"], - "project": "Public Project", - "agent_id": "agent-server-1", + "project": "PIXELVERSE DASHBOARD", + "agent_id": "DESIGNER", "role": zone["roles"][0], "status": "active", "updated_at": "2026-08-08T11:59:00Z", @@ -79,6 +81,42 @@ def _event(self, department_id="development", **overrides): payload.update(overrides) return payload + def _heartbeat(self, *, project="PIXELVERSE DASHBOARD", agent_id="DESIGNER", run_id="task-server-1"): + return { + "project": project, + "agentId": agent_id, + "runId": run_id, + "sessionId": f"session-{agent_id.lower()}-server", + "state": "working", + "heartbeatAt": "2026-08-08T11:59:45Z", + } + + def _payload(self, bridge_data, *, heartbeats, owner_view=False): + with tempfile.TemporaryDirectory() as tmp: + heartbeat_path = Path(tmp) / "jarvis-agent-details.json" + heartbeat_path.write_text( + json.dumps({ + "updatedAt": "2026-08-08T11:59:45Z", + "agents": { + f"record-{index}": record + for index, record in enumerate(heartbeats) + }, + }), + encoding="utf-8", + ) + try: + return self.server._department_campus_payload( + bridge_data, + heartbeat_path=heartbeat_path, + now=NOW, + owner_view=owner_view, + ) + except TypeError as exc: + self.fail( + "RED: server payload must accept explicit heartbeat_path " + f"({exc})" + ) + def _snapshot(self, events, *, updated_at="2026-08-08T11:59:30Z", **metadata_overrides): metadata = { "event": "status", @@ -95,19 +133,22 @@ def _snapshot(self, events, *, updated_at="2026-08-08T11:59:30Z", **metadata_ove "messages": [{"body": ""}], } - def test_ac_3_err_5_only_newest_same_metadata_verified_manager_snapshot_is_used(self): + def test_ac_3_err_5_manager_snapshots_without_heartbeat_are_not_live(self): older = self._snapshot( - [self._event(event_id="evt-old", project="Older Snapshot")], + [self._event(event_id="evt-old")], updated_at="2026-08-08T11:50:00Z", ) newer = self._snapshot( - [self._event(event_id="evt-new", project="Newest Snapshot")], + [self._event(event_id="evt-new")], updated_at="2026-08-08T11:59:30Z", source_agent="main-manager", ) - projected = self.server._department_campus_payload({"tasks": [older, newer]}, now=NOW) + projected = self._payload( + {"tasks": [older, newer]}, + heartbeats=[self._heartbeat()], + ) self.assertEqual(projected["state"], "active") - self.assertEqual([item["project"] for item in projected["events"]], ["Newest Snapshot"]) + self.assertEqual([item["event_id"] for item in projected["events"]], ["evt-new"]) ordinary = self._snapshot([self._event()]) ordinary["metadata"].pop("source_agent") @@ -122,26 +163,23 @@ def test_ac_3_err_5_only_newest_same_metadata_verified_manager_snapshot_is_used( {"tasks": [{"metadata": {}}, {"metadata": {"source_agent": "MAIN MANAGER"}}]}, ): with self.subTest(bridge_data=bridge_data): - rejected = self.server._department_campus_payload(bridge_data, now=NOW) + rejected = self._payload(bridge_data, heartbeats=[self._heartbeat()]) self.assertEqual(rejected["state"], "empty") self.assertEqual(rejected["events"], []) - def test_ac_3_source_agent_provenance_requires_a_real_string(self): + def test_ac_3_source_agent_provenance_alone_never_proves_live_work(self): for source in ("MAIN MANAGER", "main-manager", "main_manager"): with self.subTest(kind="valid_string", source=source): snapshot = self._snapshot( - [self._event(project="Verified Manager")], + [self._event()], source_agent=source, ) - projected = self.server._department_campus_payload( + projected = self._payload( {"tasks": [snapshot]}, - now=NOW, + heartbeats=[self._heartbeat()], ) self.assertEqual(projected["state"], "active") - self.assertEqual( - [event["project"] for event in projected["events"]], - ["Verified Manager"], - ) + self.assertEqual(projected["events"][0]["project"], "PIXELVERSE DASHBOARD") for source in ({"main": "manager"}, ["main", "manager"]): with self.subTest(kind="forged_non_string", source=source): @@ -149,15 +187,15 @@ def test_ac_3_source_agent_provenance_requires_a_real_string(self): [self._event(project="FORGED NON STRING SOURCE")], source_agent=source, ) - projected = self.server._department_campus_payload( + projected = self._payload( {"tasks": [snapshot]}, - now=NOW, + heartbeats=[self._heartbeat()], ) self.assertEqual(projected["state"], "empty") self.assertEqual(projected["events"], []) self.assertNotIn("FORGED", repr(projected)) - def test_ac_fresh_canonical_bridge_dispatch_projects_one_safe_active_owner(self): + def test_ac_fresh_canonical_bridge_dispatch_without_heartbeat_is_not_live(self): task = { "id": "bridge-task-23", "status": "running", @@ -176,26 +214,14 @@ def test_ac_fresh_canonical_bridge_dispatch_projects_one_safe_active_owner(self) }, } - projected = self.server._department_campus_payload({"tasks": [task]}, now=NOW) + projected = self._payload({"tasks": [task]}, heartbeats=[]) self.assertEqual(tuple(projected), TOP_LEVEL_FIELDS) - self.assertEqual(projected["state"], "active") - self.assertEqual(projected["visible_task_count"], 1) + self.assertNotEqual(projected["state"], "active") + self.assertEqual(projected["visible_task_count"], 0) self.assertEqual(projected["omitted_task_count"], 0) self.assertEqual(projected["privacy"], "public_projection") - self.assertEqual(len(projected["events"]), 1) - event = projected["events"][0] - self.assertEqual(tuple(event), PUBLIC_EVENT_FIELDS) - self.assertEqual(event["task_id"], "bridge-task-23") - self.assertEqual(event["department_id"], "infrastructure") - self.assertEqual(event["department_label"], "Infrastructure") - self.assertEqual(event["project"], "JARVIS") - self.assertEqual(event["agent_id"], "INFRASTRUCTURE") - self.assertEqual(event["role"], "Infrastructure Engineer") - self.assertEqual(event["status"], "active") - self.assertEqual(event["updated_at"], "2026-08-08T11:59:00Z") - self.assertTrue(event["ephemeral"]) - self.assertEqual(event["zone_id"], "campus-zone-infrastructure") + self.assertEqual(projected["events"], []) rendered = repr(projected) for private_value in ( "", @@ -206,7 +232,7 @@ def test_ac_fresh_canonical_bridge_dispatch_projects_one_safe_active_owner(self) ): self.assertNotIn(private_value, rendered) - def test_bridge_fallback_maps_lifecycle_timestamps_and_fails_closed(self): + def test_bridge_fallback_lifecycle_never_becomes_live_without_heartbeat(self): base = { "id": "bridge-lifecycle", "agent_role": "TESTER", @@ -215,22 +241,13 @@ def test_bridge_fallback_maps_lifecycle_timestamps_and_fails_closed(self): "claimed_at": "2026-08-08T11:58:00Z", "completed_at": "2026-08-08T11:59:00Z", } - cases = ( - ("pending", "queued", "2026-08-08T11:57:00Z"), - ("claimed", "active", "2026-08-08T11:58:00Z"), - ("done", "done", "2026-08-08T11:59:00Z"), - ("failed", "failed", "2026-08-08T11:59:00Z"), - ) - for status, expected_status, expected_time in cases: + for status in ("pending", "claimed", "done", "failed"): with self.subTest(status=status): task = {**base, "id": f"bridge-{status}", "status": status} - projected = self.server._department_campus_payload( - {"tasks": [task]}, - now=NOW, - ) - self.assertEqual(projected["state"], "active") - self.assertEqual(projected["events"][0]["status"], expected_status) - self.assertEqual(projected["events"][0]["updated_at"], expected_time) + projected = self._payload({"tasks": [task]}, heartbeats=[]) + self.assertNotEqual(projected["state"], "active") + self.assertEqual(projected["visible_task_count"], 0) + self.assertEqual(projected["events"], []) pipeline_task = { **base, @@ -238,15 +255,9 @@ def test_bridge_fallback_maps_lifecycle_timestamps_and_fails_closed(self): "status": "done", "agent_role": "SUPERVISOR_BUILDER_TESTER", } - pipeline_projection = self.server._department_campus_payload( - {"tasks": [pipeline_task]}, - now=NOW, - ) - self.assertEqual(pipeline_projection["state"], "active") - self.assertEqual( - pipeline_projection["events"][0]["agent_id"], - "INFRASTRUCTURE", - ) + pipeline_projection = self._payload({"tasks": [pipeline_task]}, heartbeats=[]) + self.assertNotEqual(pipeline_projection["state"], "active") + self.assertEqual(pipeline_projection["events"], []) rejected = ( {**base, "status": "cancelled"}, @@ -256,45 +267,39 @@ def test_bridge_fallback_maps_lifecycle_timestamps_and_fails_closed(self): ) for task in rejected: with self.subTest(rejected=task): - projected = self.server._department_campus_payload( - {"tasks": [task]}, - now=NOW, - ) + projected = self._payload({"tasks": [task]}, heartbeats=[]) self.assertEqual(projected["state"], "empty") self.assertEqual(projected["events"], []) - def test_live_bridge_snapshot_uses_terminal_or_claimed_timestamp_when_updated_at_is_absent(self): - """Bridge task rows expose lifecycle timestamps, not an updated_at field.""" + + def test_bridge_snapshot_lifecycle_timestamp_without_heartbeat_is_not_live(self): for lifecycle_field in ("completed_at", "claimed_at", "created_at"): with self.subTest(lifecycle_field=lifecycle_field): snapshot = self._snapshot( - [self._event(project="MY DICTIONARY")], + [self._event()], updated_at="2026-08-08T11:59:30Z", ) snapshot.pop("updated_at") snapshot[lifecycle_field] = "2026-08-08T11:59:30Z" - projected = self.server._department_campus_payload( + projected = self._payload( {"tasks": [snapshot]}, - now=NOW, + heartbeats=[self._heartbeat()], ) self.assertEqual(projected["state"], "active") self.assertEqual(projected["visible_task_count"], 1) - self.assertEqual( - [event["project"] for event in projected["events"]], - ["MY DICTIONARY"], - ) + self.assertEqual(projected["events"][0]["project"], "PIXELVERSE DASHBOARD") def test_err_1_non_object_bridge_or_non_list_pixel_events_is_unavailable_without_partial_data(self): malformed_sources = (None, [], "bad", 7, True) for bridge_data in malformed_sources: with self.subTest(bridge_data=bridge_data): - projected = self.server._department_campus_payload(bridge_data, now=NOW) + projected = self._payload(bridge_data, heartbeats=[]) self.assertEqual(projected["state"], "unavailable") self.assertEqual(projected["events"], []) self.assertEqual(projected["visible_task_count"], 0) malformed_events = self._snapshot({"not": "a list"}) - projected = self.server._department_campus_payload({"tasks": [malformed_events]}, now=NOW) + projected = self._payload({"tasks": [malformed_events]}, heartbeats=[]) self.assertEqual(projected["state"], "unavailable") self.assertEqual(projected["events"], []) @@ -303,13 +308,13 @@ def test_ac_11_ec_7_stale_or_malformed_refresh_returns_no_previous_agents(self): [self._event(updated_at=(NOW - timedelta(minutes=31)).isoformat())], updated_at=(NOW - timedelta(minutes=31)).isoformat(), ) - projected = self.server._department_campus_payload({"tasks": [stale]}, now=NOW) + projected = self._payload({"tasks": [stale]}, heartbeats=[self._heartbeat()]) self.assertEqual(projected["state"], "stale") self.assertEqual(projected["events"], []) self.assertEqual(projected["visible_task_count"], 0) malformed = self._snapshot([self._event(updated_at="not-a-time")]) - projected = self.server._department_campus_payload({"tasks": [malformed]}, now=NOW) + projected = self._payload({"tasks": [malformed]}, heartbeats=[self._heartbeat()]) self.assertNotEqual(projected["state"], "active") self.assertEqual(projected["events"], []) @@ -326,24 +331,39 @@ def bridge_request(method, path, payload=None): return {"tasks": [self._snapshot([self._event(prompt="")])]} original_payload = self.server._department_campus_payload + with tempfile.TemporaryDirectory() as tmp: + heartbeat_path = Path(tmp) / "jarvis-agent-details.json" + heartbeat_path.write_text( + json.dumps({ + "updatedAt": "2026-08-08T11:59:45Z", + "agents": {"designer": self._heartbeat()}, + }), + encoding="utf-8", + ) - def payload_at_fixed_now(bridge_data, *, now=None): - return original_payload(bridge_data, now=NOW) + def payload_at_fixed_now(bridge_data, *, now=None, **kwargs): + return original_payload( + bridge_data, + heartbeat_path=heartbeat_path, + now=NOW, + **kwargs, + ) - with ( - patch.object(self.server, "_bridge_request", side_effect=bridge_request), - patch.object( - self.server, - "_department_campus_payload", - side_effect=payload_at_fixed_now, - ), - ): - handler.do_GET() + with ( + patch.object(self.server, "_bridge_request", side_effect=bridge_request), + patch.object( + self.server, + "_department_campus_payload", + side_effect=payload_at_fixed_now, + ), + ): + handler.do_GET() self.assertEqual(response["status"], 200) self.assertEqual(requested, [("GET", "/api/tasks?limit=24&include_messages=1", None)]) self.assertEqual(tuple(response["payload"]), TOP_LEVEL_FIELDS) self.assertEqual(response["payload"]["privacy"], "public_projection") + self.assertEqual(response["payload"]["state"], "active") self.assertEqual(tuple(response["payload"]["events"][0]), PUBLIC_EVENT_FIELDS) self.assertNotIn("", repr(response["payload"])) self.assertNotIn("", repr(response["payload"])) diff --git a/builder/tests/test_department_campus_visual_integration.py b/builder/tests/test_department_campus_visual_integration.py index 4088ea1..6f48ad7 100644 --- a/builder/tests/test_department_campus_visual_integration.py +++ b/builder/tests/test_department_campus_visual_integration.py @@ -438,11 +438,12 @@ def test_ac_9_configured_roster_is_visible_without_claiming_live_work(self): self.assertIn("const activeAgentCount = new Set", self.script) self.assertIn("в команде", self.script) - def test_ac_10_idle_residents_walk_inside_rooms_and_reduced_motion_stops_them(self): + def test_ac_10_idle_residents_stay_static_and_reduced_motion_remains_a_backstop(self): self.assertIn("@keyframes campusResidentWander", self.css) - self.assertRegex( - self.css, - r"\.campus-resident\.is-wandering[^{]*\{[^}]*animation:", + self.assertNotIn("is-wandering", self.campus_html) + self.assertTrue( + all(resident["wandering"] is False for resident in self.campus.CAMPUS_RESIDENTS), + "idle residents must not imply work through ambient walking", ) self.assertRegex( self.css, diff --git a/builder/tests/test_department_campus_work_summary.py b/builder/tests/test_department_campus_work_summary.py index cafe703..6075c3a 100644 --- a/builder/tests/test_department_campus_work_summary.py +++ b/builder/tests/test_department_campus_work_summary.py @@ -3,6 +3,7 @@ from datetime import datetime, timezone import importlib import importlib.util +import json import os from pathlib import Path import re @@ -111,6 +112,66 @@ def _manager_snapshot(self, event: dict, **metadata_overrides) -> dict: "metadata": metadata, } + def _heartbeat_document(self, tasks: list[dict]) -> dict: + records = {} + registry = { + record["project"]: record["agent_id"] + for record in self.campus.CAMPUS_PROJECTS + } + for task in tasks: + metadata = task.get("metadata") if isinstance(task, dict) else None + pixel_events = metadata.get("pixel_events") if isinstance(metadata, dict) else None + candidates = pixel_events if isinstance(pixel_events, list) else [task] + for item in candidates: + if not isinstance(item, dict): + continue + project = item.get("project") + agent_id = item.get("agent_id") or registry.get(project) + run_id = item.get("task_id") or item.get("id") + if not all(isinstance(value, str) for value in (project, agent_id, run_id)): + continue + records[agent_id.lower()] = { + "project": project, + "agentId": agent_id, + "runId": run_id, + "sessionId": f"session-{agent_id.lower()}-summary", + "state": "working", + "heartbeatAt": datetime.now(timezone.utc).isoformat(), + } + return { + "updatedAt": datetime.now(timezone.utc).isoformat(), + "agents": records, + } + + def _with_heartbeat_payload(self, tasks: list[dict], callback): + original_payload = self.server._department_campus_payload + with tempfile.TemporaryDirectory() as tmp: + heartbeat_path = Path(tmp) / "jarvis-agent-details.json" + heartbeat_path.write_text( + json.dumps(self._heartbeat_document(tasks)), + encoding="utf-8", + ) + + def payload_with_heartbeat(data, **kwargs): + try: + return original_payload( + data, + heartbeat_path=heartbeat_path, + **kwargs, + ) + except TypeError as exc: + self.fail( + "RED: endpoint projection must read explicit heartbeat storage " + f"({exc})" + ) + + with patch.object( + self.server, + "_department_campus_payload", + side_effect=payload_with_heartbeat, + ): + return callback() + def _endpoint_payload(self, tasks: list[dict], *, owner: bool) -> dict: response: dict = {} handler = self.server.Handler.__new__(self.server.Handler) @@ -120,10 +181,13 @@ def _endpoint_payload(self, tasks: list[dict], *, owner: bool) -> dict: handler._json_response = lambda status, payload: response.update( status=status, payload=payload ) - with patch.object( - self.server, "_bridge_request", return_value={"tasks": tasks} - ): - handler.do_GET() + def request(): + with patch.object( + self.server, "_bridge_request", return_value={"tasks": tasks} + ): + handler.do_GET() + + self._with_heartbeat_payload(tasks, request) self.assertEqual(response.get("status"), 200) return response["payload"] @@ -148,10 +212,13 @@ def _public_departments_get( handler._json_response = lambda status, payload: response.update( status=status, payload=payload ) - with patch.object( - self.server, "_bridge_request", return_value={"tasks": tasks} - ): - handler.do_GET() + def request(): + with patch.object( + self.server, "_bridge_request", return_value={"tasks": tasks} + ): + handler.do_GET() + + self._with_heartbeat_payload(tasks, request) self.assertEqual(response.get("status"), 200) return response["payload"] @@ -353,6 +420,8 @@ def test_ac07_existing_or_env_dashboard_token_authorizes_without_mutation(self): provided_token="owner-token", ) + self.assertEqual(payload["state"], "active") + self.assertEqual(len(payload["events"]), 1) self.assertEqual(payload["events"][0]["issue_number"], 34) self.assertIn("work_summary", payload["events"][0]) if token_source == "file": @@ -411,16 +480,23 @@ def test_ec02_summary_is_whitespace_normalized_unicode_safe_and_bounded(self): def test_err03_existing_empty_cap_focus_escape_and_read_only_contracts_remain(self): now = datetime(2026, 8, 14, 12, 0, tzinfo=timezone.utc) - zone = self.campus.DEPARTMENT_ZONES["development"] + identities = ( + ("hq", "MAIN MANAGER", "COORDINATOR"), + ("sales", "AI STUDIO", "RESEARCHER"), + ("development", "MY DICTIONARY", "BUILDER"), + ("design", "PIXELVERSE DASHBOARD", "DESIGNER"), + ) events = [] - for index in range(4): + heartbeats = [] + for index, (department_id, project, agent_id) in enumerate(identities): + zone = self.campus.DEPARTMENT_ZONES[department_id] events.append({ "event_id": f"evt-regression-{index}", "task_id": f"task-regression-{index}", - "department_id": "development", + "department_id": department_id, "department_label": zone["label"], - "project": "MY DICTIONARY", - "agent_id": f"agent-regression-{index}", + "project": project, + "agent_id": agent_id, "role": zone["roles"][0], "status": "active", "updated_at": "2026-08-14T11:59:00Z", @@ -429,8 +505,30 @@ def test_err03_existing_empty_cap_focus_escape_and_read_only_contracts_remain(se "ephemeral": True, "zone_id": zone["zone_id"], }) - capped = self.campus.department_campus_projection(events, now=now) - empty = self.campus.department_campus_projection([], now=now) + heartbeats.append({ + "project": project, + "agent_id": agent_id, + "run_id": f"task-regression-{index}", + "session_id": f"session-regression-{index}", + "state": "working", + "heartbeat_at": "2026-08-14T11:59:45Z", + }) + try: + capped = self.campus.department_campus_projection( + events, + heartbeats=heartbeats, + now=now, + ) + empty = self.campus.department_campus_projection( + [], + heartbeats=[], + now=now, + ) + except TypeError as exc: + self.fail( + "RED: cap regression requires explicit heartbeat projection " + f"({exc})" + ) self.assertEqual(capped["visible_task_count"], 3) self.assertEqual(capped["omitted_task_count"], 1) self.assertEqual(empty["state"], "empty")