From 3b4065ecc6cbf88f796b7a756ec4f40cda8cdb66 Mon Sep 17 00:00:00 2001 From: Yvette Carlisle Date: Fri, 12 Jun 2026 08:47:33 +0800 Subject: [PATCH] {"schema":"decodex/commit/1","summary":"Surface Program Intake state in status and dashboard","authority":"XY-941","related":["XY-936"]} --- apps/decodex/src/execution_program.rs | 15 + apps/decodex/src/orchestrator.rs | 4 +- .../src/orchestrator/agent_evidence.rs | 2 - apps/decodex/src/orchestrator/git_ops.rs | 2 +- .../src/orchestrator/operator_dashboard.html | 211 ++++++++- .../src/orchestrator/program_reconciler.rs | 3 +- apps/decodex/src/orchestrator/status.rs | 127 ++++- apps/decodex/src/orchestrator/tests.rs | 3 +- .../tests/operator/status/agent_evidence.rs | 2 - .../tests/operator/status/dashboard.rs | 26 +- .../tests/operator/status/text.rs | 446 +++++++++++++++++- apps/decodex/src/orchestrator/types.rs | 334 ++++++++++++- docs/reference/operator-control-plane.md | 17 +- 13 files changed, 1172 insertions(+), 20 deletions(-) diff --git a/apps/decodex/src/execution_program.rs b/apps/decodex/src/execution_program.rs index ee2aa8a7a..53f7949d5 100644 --- a/apps/decodex/src/execution_program.rs +++ b/apps/decodex/src/execution_program.rs @@ -370,6 +370,16 @@ pub(crate) enum ExecutionQueueLabelAction { /// Remove the service queue label. Remove, } +impl ExecutionQueueLabelAction { + /// Stable machine-readable queue-label action. + pub(crate) fn as_str(self) -> &'static str { + match self { + Self::Apply => "apply", + Self::Retain => "retain", + Self::Remove => "remove", + } + } +} /// Conflict-domain key for one program node. #[derive(Clone, Debug, Eq, Hash, PartialEq, Deserialize, Serialize)] @@ -1726,6 +1736,11 @@ fn lifecycle_state_for( { return ExecutionProgramNodeLifecycleState::NeedsAttention; } + if let Some(issue) = node.linear_issue() + && issue.has_active_label + { + return ExecutionProgramNodeLifecycleState::Active; + } match state { ExecutionReadinessState::NotReady | ExecutionReadinessState::Paused => diff --git a/apps/decodex/src/orchestrator.rs b/apps/decodex/src/orchestrator.rs index a413d5f42..b51d079c1 100644 --- a/apps/decodex/src/orchestrator.rs +++ b/apps/decodex/src/orchestrator.rs @@ -13,7 +13,7 @@ pub(crate) use lane_control::{ #[cfg(unix)] use std::os::fd::AsRawFd; use std::{ cmp::Ordering, - collections::{HashMap, HashSet}, + collections::{self, BTreeMap, BTreeSet, HashMap, HashSet}, env, error::Error, fmt::{self, Display, Formatter}, @@ -40,7 +40,7 @@ use time::{OffsetDateTime, format_description::well_known::Rfc3339}; use crate::{agent, default_branch_sync, git_credentials, maintenance, state}; #[rustfmt::skip] -use crate::{agent::{ACTIVE_RUN_IDLE_TIMEOUT, AppServerCapabilityPreflightFailure, AppServerDynamicToolFailure, AppServerHomePreflightFailure, AppServerPhaseGoalFailure, AppServerProcessEnv, AppServerRunRequest, AppServerRunResult, AppServerTransportFailure, AppServerTurnFailure, ISSUE_DELIVERY_CLOSEOUT_COMPLETE_TOOL_NAME, ISSUE_LABEL_ADD_TOOL_NAME, ISSUE_PROGRESS_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_HANDOFF_TOOL_NAME, ISSUE_REVIEW_REPAIR_COMPLETE_TOOL_NAME, ISSUE_TERMINAL_FINALIZE_TOOL_NAME, ISSUE_TRANSITION_TOOL_NAME, DecodexRunContext, DecodexToolBridge, PhaseGoalController, PhaseGoalKind, PhaseGoalSpec, PhaseGoalTransition, ReviewExecutionMode, ReviewHandoffContext, ReviewHandoffWritebackFailed, ReviewPolicyStopReason, ReviewPolicyStopRequested, RunCompletionDisposition, TrackerToolBridge, TurnContinuationGuard}, config::{ReviewLevel, ServiceConfig}, execution_program::{ExecutionProgramOperatorSummary, ExecutionProgramReadinessContext, ExecutionWorkflowPolicy}, git_credentials::GitCredentialSource, github, prelude::{Result, eyre}, state::{ChildAgentActivityBucket, ChildAgentActivitySummary, CodexAccountActivitySummary, ExecutionProgramRecord, LoopGuardrailCheckpoint, LoopGuardrailCheckpointInput, ProjectRegistration, ProjectRunStatus, ProtocolActivitySummary, RUN_OPERATION_AGENT_RUN, RUN_OPERATION_APP_SERVER_PREFLIGHT, RUN_OPERATION_GIT_CREDENTIALS, RUN_OPERATION_IDLE, RUN_OPERATION_RECONCILIATION, RUN_OPERATION_REPO_GATE, RUN_OPERATION_REVIEW_WRITEBACK, RUN_OPERATION_WAITING_EXTERNAL, ReviewHandoffMarker, ReviewOrchestrationMarker, RunActivityMarker, RunAttempt, StateStore, WorktreeMapping}, tracker::{IssueTracker, TrackerComment, TrackerIssue, linear::LinearClient, records}, workflow::{WorkflowDocument, WorkflowExecution}, worktree::{WorktreeManager, WorktreeSpec}}; +use crate::{agent::{ACTIVE_RUN_IDLE_TIMEOUT, AppServerCapabilityPreflightFailure, AppServerDynamicToolFailure, AppServerHomePreflightFailure, AppServerPhaseGoalFailure, AppServerProcessEnv, AppServerRunRequest, AppServerRunResult, AppServerTransportFailure, AppServerTurnFailure, ISSUE_DELIVERY_CLOSEOUT_COMPLETE_TOOL_NAME, ISSUE_LABEL_ADD_TOOL_NAME, ISSUE_PROGRESS_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_HANDOFF_TOOL_NAME, ISSUE_REVIEW_REPAIR_COMPLETE_TOOL_NAME, ISSUE_TERMINAL_FINALIZE_TOOL_NAME, ISSUE_TRANSITION_TOOL_NAME, DecodexRunContext, DecodexToolBridge, PhaseGoalController, PhaseGoalKind, PhaseGoalSpec, PhaseGoalTransition, ReviewExecutionMode, ReviewHandoffContext, ReviewHandoffWritebackFailed, ReviewPolicyStopReason, ReviewPolicyStopRequested, RunCompletionDisposition, TrackerToolBridge, TurnContinuationGuard}, config::{ReviewLevel, ServiceConfig}, execution_program::{ExecutionNodeEvaluation, ExecutionProgramEvaluation, ExecutionProgramOperatorSummary, ExecutionProgramReadinessContext, ExecutionQueueLabelAction, ExecutionWorkflowPolicy}, git_credentials::GitCredentialSource, github, prelude::{Result, eyre}, state::{ChildAgentActivityBucket, ChildAgentActivitySummary, CodexAccountActivitySummary, ExecutionProgramRecord, LoopGuardrailCheckpoint, LoopGuardrailCheckpointInput, ProjectRegistration, ProjectRunStatus, ProtocolActivitySummary, RUN_OPERATION_AGENT_RUN, RUN_OPERATION_APP_SERVER_PREFLIGHT, RUN_OPERATION_GIT_CREDENTIALS, RUN_OPERATION_IDLE, RUN_OPERATION_RECONCILIATION, RUN_OPERATION_REPO_GATE, RUN_OPERATION_REVIEW_WRITEBACK, RUN_OPERATION_WAITING_EXTERNAL, ReviewHandoffMarker, ReviewOrchestrationMarker, RunActivityMarker, RunAttempt, StateStore, WorktreeMapping}, tracker::{IssueTracker, TrackerComment, TrackerIssue, linear::LinearClient, records}, workflow::{WorkflowDocument, WorkflowExecution}, worktree::{WorktreeManager, WorktreeSpec}}; use harness_improvement::{ HarnessImprovementCandidateSummary, HarnessOutcomeKind, harness_improvement_candidates_from_private_events, record_harness_outcome_best_effort, diff --git a/apps/decodex/src/orchestrator/agent_evidence.rs b/apps/decodex/src/orchestrator/agent_evidence.rs index b8a77bcf2..c32e46d0a 100644 --- a/apps/decodex/src/orchestrator/agent_evidence.rs +++ b/apps/decodex/src/orchestrator/agent_evidence.rs @@ -1,5 +1,3 @@ -use std::collections::{self, BTreeMap}; - const AGENT_HANDOFF_INDEX_SCHEMA: &str = "decodex.agent_handoff_index/1"; const AGENT_BLOCKER_SNAPSHOT_SCHEMA: &str = "decodex.blocker_snapshot/1"; const AGENT_RUN_CAPSULE_SCHEMA: &str = "decodex.run_capsule/1"; diff --git a/apps/decodex/src/orchestrator/git_ops.rs b/apps/decodex/src/orchestrator/git_ops.rs index 3aed90bba..67c7effe8 100644 --- a/apps/decodex/src/orchestrator/git_ops.rs +++ b/apps/decodex/src/orchestrator/git_ops.rs @@ -139,7 +139,7 @@ mod repo_gate_failure { } } -use std::{collections::BTreeSet, process::Output}; +use std::process::Output; use repo_gate_failure::{RepoGateFailure, RepoGateFailureDisposition, RepoGateFailureKind}; use crate::workflow::ResolvedRepoGate; diff --git a/apps/decodex/src/orchestrator/operator_dashboard.html b/apps/decodex/src/orchestrator/operator_dashboard.html index fa017627c..77c7b370b 100644 --- a/apps/decodex/src/orchestrator/operator_dashboard.html +++ b/apps/decodex/src/orchestrator/operator_dashboard.html @@ -3603,6 +3603,18 @@

Running Lanes

+
+
+
+

Program Intake

+
+

+
+
+
+
+
+
@@ -3786,6 +3798,7 @@

Run History

projects: document.getElementById("projects-panel"), accountPool: document.getElementById("account-pool-panel"), active: document.getElementById("active-panel"), + programs: document.getElementById("programs-panel"), queue: document.getElementById("queue-panel"), recent: document.getElementById("recent-panel"), review: document.getElementById("review-panel"), @@ -3804,6 +3817,8 @@

Run History

queuedMeta: document.getElementById("queued-meta"), activeRuns: document.getElementById("active-runs"), activeRunsMeta: document.getElementById("active-runs-meta"), + executionPrograms: document.getElementById("execution-programs"), + programsMeta: document.getElementById("programs-meta"), recentRuns: document.getElementById("recent-runs"), recentRunsMeta: document.getElementById("recent-runs-meta"), reviewQueue: document.getElementById("review-queue"), @@ -3849,12 +3864,12 @@

Run History

const detailAnimationTimers = new WeakMap(); const DASHBOARD_LAYOUT = { - primary: ["accountPool", "projects", "active", "queue", "review", "worktrees", "recent"], + primary: ["accountPool", "projects", "active", "programs", "queue", "review", "worktrees", "recent"], }; const DASHBOARD_SECTION_GROUPS = [ { marker: "control", panels: ["accountPool"] }, { marker: "projects", panels: ["projects"] }, - { marker: "execution", panels: ["active", "queue"] }, + { marker: "execution", panels: ["active", "programs", "queue"] }, { marker: "aftercare", panels: ["review", "worktrees", "recent"] }, ]; const COPY = { @@ -8656,6 +8671,7 @@

Run History

const activeRuns = snapshot?.active_runs ?? []; const recentRuns = snapshot?.recent_runs ?? []; const historyRuns = sessionHistoryRuns(snapshot); + const executionPrograms = snapshot?.execution_programs ?? []; const postReviewLanes = snapshot?.post_review_lanes ?? []; const postReviewIssueKeys = new Set(postReviewLanes.flatMap(issueIdentityKeys)); const activeRunByIssue = new Map(); @@ -8878,6 +8894,17 @@

Run History

const runningAttentionCount = attentionItems.filter((item) => item.scope === "Running").length; const liveRuns = activeRuns.filter(runCountsAsRunning).length; const intakeAttentionCount = queueBacklogCandidates.filter(queuedCandidateNeedsAttention).length; + const programAttentionCount = executionPrograms.filter((program) => + ["attention", "blocked", "stale"].includes(program.status), + ).length; + const programReadyCount = executionPrograms.reduce( + (total, program) => total + Number(program.ready_count || 0), + 0, + ); + const programQueuedCount = executionPrograms.reduce( + (total, program) => total + Number(program.queued_count || 0), + 0, + ); return { projects: snapshot?.projects ?? [], @@ -8905,6 +8932,10 @@

Run History

reviewBlockerCount, cleanupCount, runningAttentionCount, + executionPrograms, + programAttentionCount, + programReadyCount, + programQueuedCount, sessionHistoryRuns: historyRuns, worktrees, postReviewLanes, @@ -9523,6 +9554,181 @@

Run History

return parts.join(" · "); } + function programMetaText(snapshot, derived) { + if (!snapshot) { + return "0 programs"; + } + + const parts = [pluralize(derived.executionPrograms.length, "program")]; + + if (derived.programReadyCount) { + parts.push(`${derived.programReadyCount} ready`); + } + if (derived.programQueuedCount) { + parts.push(`${derived.programQueuedCount} queued`); + } + if (derived.programAttentionCount) { + parts.push( + derived.programAttentionCount === 1 + ? "1 needs attention" + : `${derived.programAttentionCount} need attention`, + ); + } + + return parts.join(" · "); + } + + function toneForProgram(program) { + switch (program.status) { + case "ready": + case "completed": + return "tone-ready"; + case "queued": + return "tone-queue"; + case "active": + return "tone-run"; + case "blocked": + case "attention": + case "stale": + return "tone-blocked"; + case "held": + return "tone-wait"; + default: + return "tone-muted"; + } + } + + function programMappedIssues(program) { + const issues = program.mapped_issue_identifiers ?? []; + + return issues.length ? issues.join(", ") : "NONE"; + } + + function programQueueActionSummary(program) { + const actions = [ + ["apply", program.queue_label_apply_count], + ["retain", program.queue_label_retain_count], + ["remove", program.queue_label_remove_count], + ].filter(([, count]) => Number(count || 0) > 0); + + return actions.length + ? actions.map(([label, count]) => `${label} ${count}`).join(", ") + : "none"; + } + + function programProgressFacts(program) { + return [ + ["Ready", String(program.ready_count ?? 0)], + ["Queued", String(program.queued_count ?? 0)], + ["Active", String(program.active_count ?? 0)], + ["Blocked", String(program.blocked_count ?? 0)], + ["Held", String(program.held_count ?? 0)], + ["Attention", String(program.needs_attention_count ?? 0)], + ["Completed", String(program.completed_count ?? 0)], + ["Stale", String(program.stale_count ?? 0)], + ["Eligible", String(program.queue_label_eligible_count ?? 0)], + ]; + } + + function programNodeReasons(node) { + const reasons = node.reasons ?? []; + if (reasons.length) { + return reasons.join("; "); + } + + const reasonCodes = node.reason_codes ?? []; + return reasonCodes.length ? reasonCodes.map(displayToken).join(", ") : "none"; + } + + function programNodeIssue(node) { + return node.issue_identifier || "unmapped"; + } + + function renderProgramNodeReadbacks(program) { + const readbacks = program.node_readbacks ?? []; + if (!readbacks.length) { + return ""; + } + + const detailKey = `program:${program.program_id}:nodes`; + + return ` +
+ ${escapeHtml(pluralize(readbacks.length, "node diagnostic"))} +
+ ${readbacks + .map((node) => { + const reasonCodes = (node.reason_codes ?? []).map(displayToken).join(", ") || "none"; + return ` + ${field("Issue", programNodeIssue(node))} + ${field("Lifecycle", displayToken(node.lifecycle_state || "unknown"))} + ${field("Readiness", displayToken(node.readiness_state || "unknown"))} + ${field("Issue state", node.issue_state || "none")} + ${field("Queue action", node.queue_label_action || "none")} + ${field("Reason codes", reasonCodes)} + ${field("Reasons", programNodeReasons(node))} + ${field("Next action", node.next_action || "none")} + `; + }) + .join("")} +
+
+ `; + } + + function renderExecutionPrograms(snapshot, derived) { + const programs = derived.executionPrograms ?? []; + setPanelMeta( + nodes.programsMeta, + programMetaText(snapshot, derived), + derived.programAttentionCount ? "attention" : "", + ); + + if (!programs.length) { + renderRoutineEmptyList(nodes.executionPrograms); + return; + } + + renderStableList( + nodes.executionPrograms, + programs + .map((program) => { + const tone = toneForProgram(program); + const mappedIssues = programMappedIssues(program); + const warning = program.readback_warning + ? inlineStatusFact("Warning", displayToken(program.readback_warning)) + : ""; + return ` +
+
+
+
+ ${escapeHtml(displayToken(program.intake_kind || "program"))} + ${escapeHtml(program.program_id)} +
+

${escapeHtml(program.public_summary || program.program_id)}

+
+
+
+ ${statusLabel(displayToken(program.status || "unknown"), tone)} + ${inlineStatusFact("Queue", programQueueActionSummary(program))} + ${warning} +
+
+ ${cardField("Mapped issues", mappedIssues, mappedIssues === "NONE" ? "is-muted" : "")} + ${cardField("Source contract", program.source_contract_id || "NONE", program.source_contract_id ? "" : "is-muted")} + ${programProgressFacts(program) + .map(([label, value]) => cardField(label, value)) + .join("")} +
+ ${renderProgramNodeReadbacks(program)} +
+ `; + }) + .join(""), + ); + } + function renderQueuedCandidates(container, items) { if (!items.length) { renderRoutineEmptyList(container); @@ -10637,6 +10843,7 @@

${escapeHtml(worktree.branch_name)}

renderProjects(snapshot, derived); renderAccountPool(snapshot); renderActiveRuns(snapshot, derived); + renderExecutionPrograms(snapshot, derived); renderQueuedCandidates( nodes.queuedCandidates, derived.queueBacklogCandidates, diff --git a/apps/decodex/src/orchestrator/program_reconciler.rs b/apps/decodex/src/orchestrator/program_reconciler.rs index f0b75a33a..ca9c671ce 100644 --- a/apps/decodex/src/orchestrator/program_reconciler.rs +++ b/apps/decodex/src/orchestrator/program_reconciler.rs @@ -1,6 +1,5 @@ use crate::execution_program::{ - ExecutionDependencySnapshot, ExecutionLinearIssueMapping, ExecutionNodeEvaluation, - ExecutionProgram, ExecutionQueueLabelAction, + ExecutionDependencySnapshot, ExecutionLinearIssueMapping, ExecutionProgram, }; use crate::execution_program::ExecutionConflictDomain; diff --git a/apps/decodex/src/orchestrator/status.rs b/apps/decodex/src/orchestrator/status.rs index df0c1b2e0..f1c1d526b 100644 --- a/apps/decodex/src/orchestrator/status.rs +++ b/apps/decodex/src/orchestrator/status.rs @@ -669,10 +669,16 @@ fn operator_execution_program_statuses( state_store: &StateStore, ) -> crate::prelude::Result> { let policy = ExecutionWorkflowPolicy::from_workflow(project.service_id(), workflow)?; - let context = ExecutionProgramReadinessContext::new(); + let records = state_store.list_execution_programs(project.service_id())?; + let context = operator_execution_program_readiness_context( + project.service_id(), + workflow, + state_store, + &records, + )?; let mut statuses = Vec::new(); - for record in state_store.list_execution_programs(project.service_id())? { + for record in records { let evaluation = if let Some(source_contract_id) = record.source_contract_id() { let Some(contract) = state_store.decision_contract(project.service_id(), source_contract_id)? else { @@ -689,6 +695,7 @@ fn operator_execution_program_statuses( statuses.push(OperatorExecutionProgramStatus::from_summary( &record, evaluation.operator_summary(), + &evaluation, )); } @@ -697,6 +704,84 @@ fn operator_execution_program_statuses( Ok(statuses) } +fn operator_execution_program_readiness_context( + service_id: &str, + workflow: &WorkflowDocument, + state_store: &StateStore, + records: &[ExecutionProgramRecord], +) -> crate::prelude::Result { + let dependency_snapshots = operator_execution_program_dependency_snapshots(records)?; + let occupied_conflict_domains = + operator_execution_program_occupied_conflict_domains(service_id, workflow, state_store, records)?; + + Ok(ExecutionProgramReadinessContext::new() + .with_dependency_snapshots(dependency_snapshots) + .with_occupied_conflict_domains(occupied_conflict_domains)) +} + +fn operator_execution_program_dependency_snapshots( + records: &[ExecutionProgramRecord], +) -> crate::prelude::Result> { + let mut snapshots = BTreeMap::new(); + + for record in records { + for node in record.program().nodes() { + let Some(issue) = node.linear_issue() else { + continue; + }; + + insert_dependency_snapshot(&mut snapshots, node.node_id(), issue.issue_state())?; + insert_dependency_snapshot(&mut snapshots, issue.issue_identifier(), issue.issue_state())?; + } + } + + Ok(snapshots.into_values().collect()) +} + +fn operator_execution_program_occupied_conflict_domains( + service_id: &str, + workflow: &WorkflowDocument, + state_store: &StateStore, + records: &[ExecutionProgramRecord], +) -> crate::prelude::Result> { + let retained_issue_ids = state_store + .list_worktrees(service_id)? + .into_iter() + .map(|worktree| worktree.issue_id().to_owned()) + .collect::>(); + let mut occupied = Vec::new(); + let mut seen = std::collections::BTreeSet::new(); + + for record in records { + for node in record.program().nodes() { + let Some(issue) = node.linear_issue() else { + continue; + }; + let retained_nonterminal = + retained_issue_ids.contains(issue.issue_id()) + && !state_name_is_terminal(issue.issue_state(), workflow); + let issue_occupies_domain = issue.has_active_label() + || issue.has_needs_attention_label() + || retained_nonterminal + || state_store.issue_has_active_shared_claim(service_id, issue.issue_id())?; + + if !issue_occupies_domain { + continue; + } + + for domain in node.conflict_domains() { + let key = format!("{}:{}", domain.kind().as_str(), domain.key()); + + if seen.insert(key) { + occupied.push(domain.clone()); + } + } + } + } + + Ok(occupied) +} + fn hydrate_live_operator_external_observers( context: LiveOperatorStatusObserverContext<'_, T>, snapshot: &mut OperatorStatusSnapshot, @@ -6570,11 +6655,16 @@ fn append_rendered_execution_programs(output: &mut String, snapshot: &OperatorSt .readback_warning .as_ref() .map_or_else(String::new, |warning| format!(" readback_warning={warning}")); + let intake_kind = program.intake_kind.as_deref().unwrap_or("unknown"); + let public_summary = program.public_summary.as_deref().unwrap_or("none"); output.push_str(&format!( - "- program_id: {} source_contract_id: {} nodes={} planned={} mapped={} ready={} queued={} blocked={} held={} active={} attention={} completed={} stale={} superseded={} queue_label_eligible={} mapped_issues={}{}\n", + "- program_id: {} status={} source_contract_id: {} intake_kind={} summary=\"{}\" nodes={} planned={} mapped={} ready={} queued={} blocked={} held={} active={} attention={} completed={} stale={} superseded={} queue_label_eligible={} queue_actions=apply:{} retain:{} remove:{} mapped_issues={}{}\n", program.program_id, + program.status, program.source_contract_id.as_deref().unwrap_or("none"), + intake_kind, + public_summary, program.node_count, program.planned_count, program.mapped_count, @@ -6588,9 +6678,40 @@ fn append_rendered_execution_programs(output: &mut String, snapshot: &OperatorSt program.stale_count, program.superseded_count, program.queue_label_eligible_count, + program.queue_label_apply_count, + program.queue_label_retain_count, + program.queue_label_remove_count, mapped_issues, readback_warning, )); + + for node in &program.node_readbacks { + let issue_identifier = node.issue_identifier.as_deref().unwrap_or("unmapped"); + let issue_state = node.issue_state.as_deref().unwrap_or("none"); + let queue_label_action = node.queue_label_action.as_deref().unwrap_or("none"); + let reason_codes = if node.reason_codes.is_empty() { + String::from("none") + } else { + node.reason_codes.join(",") + }; + let reasons = if node.reasons.is_empty() { + String::from("none") + } else { + node.reasons.join(" | ") + }; + + output.push_str(&format!( + " - node: issue={} issue_state={} lifecycle={} readiness={} queue_label_action={} reason_codes={} reasons=\"{}\" next_action=\"{}\"\n", + issue_identifier, + issue_state, + node.lifecycle_state, + node.readiness_state, + queue_label_action, + reason_codes, + reasons, + node.next_action, + )); + } } } diff --git a/apps/decodex/src/orchestrator/tests.rs b/apps/decodex/src/orchestrator/tests.rs index fe39dfed7..86c6c8828 100644 --- a/apps/decodex/src/orchestrator/tests.rs +++ b/apps/decodex/src/orchestrator/tests.rs @@ -33,8 +33,9 @@ use crate::{orchestrator::RepoGatePhaseGoalController, tracker::records}; use crate::config::{ReviewLevel, ServiceConfig}; #[rustfmt::skip] use crate::github; +use crate::loop_contract::DecisionContract; #[rustfmt::skip] -use crate::orchestrator::{self, ActiveChildRunContext, ActiveRunDisposition, ActiveRunReconciliation, ActiveWorkflowOverride, AgentEvidenceSource, AuthorityBoundaryCheckInput, AuthorityBoundaryDisposition, AuthorityDecisionRequestInput, ChildExitRetryContext, ChildRunRef, ControlPlaneProjectTick, CONTINUATION_PENDING_RUN_STATUS, DaemonRunChild, DaemonTickRuntimeContext, DashboardEventHub, EvidenceRequest, GhPullRequestReviewStateInspector, ISSUE_DELIVERY_CLOSEOUT_COMPLETE_TOOL_NAME, ISSUE_LABEL_ADD_TOOL_NAME, ISSUE_PROGRESS_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_HANDOFF_TOOL_NAME, ISSUE_REVIEW_REPAIR_COMPLETE_TOOL_NAME, ISSUE_TERMINAL_FINALIZE_TOOL_NAME, ISSUE_TRANSITION_TOOL_NAME, IssueDispatchMode, IssueRunPlan, IssueTurnContinuationGuard, ManualAttentionRequested, OPERATOR_DASHBOARD_ALIAS_ENDPOINT_PATH, OPERATOR_DASHBOARD_ENDPOINT_PATH, OperatorCodexAccountControlStatus, OperatorExecutionProgramStatus, OperatorGitHubCliAuthority, OperatorProjectStatus, OperatorStatusSnapshot, PostReviewLaneClassification, PostReviewLaneDecision, PostReviewLaneSnapshot, PreferredRunIdentity, PrepareIssueRunContext, PublishedOperatorSnapshot, PullRequestCommitConnection, PullRequestCommitNode, PullRequestCommitPayload, PullRequestIssueCommentConnection, PullRequestIssueCommentState, PullRequestIssueCommentsNode, PullRequestPageInfo, PullRequestReactionGroup, PullRequestReactionUsersConnection, PullRequestActor, PullRequestRepository, PullRequestRepositoryOwner, PullRequestReviewConnection, PullRequestIssueCommentNode, PullRequestReviewNode, PullRequestReviewRequestConnection, PullRequestReviewState, PullRequestReviewStateInspector, PullRequestReviewStateNode, PullRequestReviewStateRepository, PullRequestReviewSummaryState, PullRequestReviewThreadConnection, PullRequestReviewThreadNode, PullRequestStatusCheckRollup, RecoveredRuntimeState, RetainedPartialProgress, RetainedReviewRunIdentity, RetryComment, RetryDispatchDecision, RetryEntry, RetryKind, RetryQueue, RunCompletionDisposition, RunSummary, RepoGateFailure, TERMINAL_GUARD_MARKER_FILE, TERMINAL_GUARDED_RUN_STATUS, TRACKER_RATE_LIMIT_WARNING, TargetIssueRunContext, EXTERNAL_REVIEW_ACTOR_LOGIN, EXTERNAL_REVIEW_PASS_PHRASE, EXTERNAL_REVIEW_REQUEST_BODY}; +use crate::orchestrator::{self, ActiveChildRunContext, ActiveRunDisposition, ActiveRunReconciliation, ActiveWorkflowOverride, AgentEvidenceSource, AuthorityBoundaryCheckInput, AuthorityBoundaryDisposition, AuthorityDecisionRequestInput, ChildExitRetryContext, ChildRunRef, ControlPlaneProjectTick, CONTINUATION_PENDING_RUN_STATUS, DaemonRunChild, DaemonTickRuntimeContext, DashboardEventHub, EvidenceRequest, GhPullRequestReviewStateInspector, ISSUE_DELIVERY_CLOSEOUT_COMPLETE_TOOL_NAME, ISSUE_LABEL_ADD_TOOL_NAME, ISSUE_PROGRESS_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_CHECKPOINT_TOOL_NAME, ISSUE_REVIEW_HANDOFF_TOOL_NAME, ISSUE_REVIEW_REPAIR_COMPLETE_TOOL_NAME, ISSUE_TERMINAL_FINALIZE_TOOL_NAME, ISSUE_TRANSITION_TOOL_NAME, IssueDispatchMode, IssueRunPlan, IssueTurnContinuationGuard, ManualAttentionRequested, OPERATOR_DASHBOARD_ALIAS_ENDPOINT_PATH, OPERATOR_DASHBOARD_ENDPOINT_PATH, OperatorCodexAccountControlStatus, OperatorExecutionProgramNodeStatus, OperatorExecutionProgramStatus, OperatorGitHubCliAuthority, OperatorProjectStatus, OperatorStatusSnapshot, PostReviewLaneClassification, PostReviewLaneDecision, PostReviewLaneSnapshot, PreferredRunIdentity, PrepareIssueRunContext, PublishedOperatorSnapshot, PullRequestCommitConnection, PullRequestCommitNode, PullRequestCommitPayload, PullRequestIssueCommentConnection, PullRequestIssueCommentState, PullRequestIssueCommentsNode, PullRequestPageInfo, PullRequestReactionGroup, PullRequestReactionUsersConnection, PullRequestActor, PullRequestRepository, PullRequestRepositoryOwner, PullRequestReviewConnection, PullRequestIssueCommentNode, PullRequestReviewNode, PullRequestReviewRequestConnection, PullRequestReviewState, PullRequestReviewStateInspector, PullRequestReviewStateNode, PullRequestReviewStateRepository, PullRequestReviewSummaryState, PullRequestReviewThreadConnection, PullRequestReviewThreadNode, PullRequestStatusCheckRollup, RecoveredRuntimeState, RetainedPartialProgress, RetainedReviewRunIdentity, RetryComment, RetryDispatchDecision, RetryEntry, RetryKind, RetryQueue, RunCompletionDisposition, RunSummary, RepoGateFailure, TERMINAL_GUARD_MARKER_FILE, TERMINAL_GUARDED_RUN_STATUS, TRACKER_RATE_LIMIT_WARNING, TargetIssueRunContext, EXTERNAL_REVIEW_ACTOR_LOGIN, EXTERNAL_REVIEW_PASS_PHRASE, EXTERNAL_REVIEW_REQUEST_BODY}; #[rustfmt::skip] use crate::prelude::Result; #[rustfmt::skip] diff --git a/apps/decodex/src/orchestrator/tests/operator/status/agent_evidence.rs b/apps/decodex/src/orchestrator/tests/operator/status/agent_evidence.rs index 3e36714ad..7287244cd 100644 --- a/apps/decodex/src/orchestrator/tests/operator/status/agent_evidence.rs +++ b/apps/decodex/src/orchestrator/tests/operator/status/agent_evidence.rs @@ -2,8 +2,6 @@ use orchestrator::HarnessOutcomeKind; use orchestrator::HarnessOutcomeRecordInput; use orchestrator::PrivateEvidenceReadback; -use crate::loop_contract::DecisionContract; - #[test] fn agent_evidence_snapshot_writes_index_blockers_capsules_and_event_stream() { let temp_dir = TempDir::new().expect("temp dir should create"); diff --git a/apps/decodex/src/orchestrator/tests/operator/status/dashboard.rs b/apps/decodex/src/orchestrator/tests/operator/status/dashboard.rs index e9b2ca229..a20a1015c 100644 --- a/apps/decodex/src/orchestrator/tests/operator/status/dashboard.rs +++ b/apps/decodex/src/orchestrator/tests/operator/status/dashboard.rs @@ -43,6 +43,30 @@ fn operator_dashboard_surfaces_loop_status_fields() { assert!(response.contains("loopStatusFacts(attention.loop_status)")); } +#[test] +fn operator_dashboard_surfaces_program_intake_panel() { + let response = dashboard_response(); + + assert!(response.contains("id=\"programs-panel\"")); + assert!(response.contains("id=\"programs-meta\"")); + assert!(response.contains("id=\"execution-programs\"")); + assert!(response.contains("

Program Intake

")); + assert!(response.contains("programs: document.getElementById(\"programs-panel\")")); + assert!(response.contains("executionPrograms: document.getElementById(\"execution-programs\")")); + assert!(response.contains("programsMeta: document.getElementById(\"programs-meta\")")); + assert!(response.contains("function renderExecutionPrograms(snapshot, derived)")); + assert!(response.contains("function renderProgramNodeReadbacks(program)")); + assert!(response.contains("function programQueueActionSummary(program)")); + assert!(response.contains("program.node_readbacks ?? []")); + assert!(response.contains("program.queue_label_apply_count")); + assert!(response.contains("program.queue_label_remove_count")); + assert!(response.contains("renderExecutionPrograms(snapshot, derived);")); + assert!(response.contains("primary: [\"accountPool\", \"projects\", \"active\", \"programs\", \"queue\", \"review\", \"worktrees\", \"recent\"]")); + assert!(response.contains("{ marker: \"execution\", panels: [\"active\", \"programs\", \"queue\"] }")); + assert!(!response.contains("data-program-edit")); + assert!(!response.contains("data-program-mutate")); +} + #[test] fn operator_dashboard_background_wash_stays_viewport_fixed() { let response = dashboard_response(); @@ -544,7 +568,7 @@ fn operator_dashboard_renders_account_usage_controls() { assert!(!response.contains("queue-group-count")); assert!(response.contains("nodes.projectTitle.textContent = \"Decodex\"")); assert!(!response.contains("Decodex Operator")); - assert!(response.contains("primary: [\"accountPool\", \"projects\", \"active\", \"queue\", \"review\", \"worktrees\", \"recent\"]")); + assert!(response.contains("primary: [\"accountPool\", \"projects\", \"active\", \"programs\", \"queue\", \"review\", \"worktrees\", \"recent\"]")); assert!(!response.contains("#account-pool-panel {")); assert!(!response.contains("No accounts")); assert!(response.contains("#active-panel {\n\t\t\t\tbackground: transparent;")); diff --git a/apps/decodex/src/orchestrator/tests/operator/status/text.rs b/apps/decodex/src/orchestrator/tests/operator/status/text.rs index 4c4b594bf..30f75d65c 100644 --- a/apps/decodex/src/orchestrator/tests/operator/status/text.rs +++ b/apps/decodex/src/orchestrator/tests/operator/status/text.rs @@ -1,5 +1,11 @@ use state::ReviewPolicyCheckpointInput; +use crate::execution_program::{ + ExecutionLinearIssueMapping, ExecutionProgram, ExecutionProgramDependency, + ExecutionProgramNode, ExecutionProgramNodeStage, ExecutionQueueIntent, +}; +use crate::loop_contract::{DecisionPromotion, DecisionPromotionActorKind}; + #[test] fn operator_status_text_surfaces_github_cli_authority() { let snapshot = OperatorStatusSnapshot { @@ -181,7 +187,10 @@ fn operator_status_text_surfaces_execution_program_summary() { history_lanes: Vec::new(), execution_programs: vec![OperatorExecutionProgramStatus { program_id: String::from("program-853"), + status: String::from("blocked"), source_contract_id: Some(String::from("contract-852")), + intake_kind: Some(String::from("goal_intake")), + public_summary: Some(String::from("Resolve promoted program work.")), node_count: 3, planned_count: 0, mapped_count: 0, @@ -195,7 +204,24 @@ fn operator_status_text_surfaces_execution_program_summary() { stale_count: 0, superseded_count: 0, queue_label_eligible_count: 1, + queue_label_apply_count: 1, + queue_label_retain_count: 0, + queue_label_remove_count: 1, mapped_issue_identifiers: vec![String::from("XY-853")], + node_readbacks: vec![OperatorExecutionProgramNodeStatus { + lifecycle_state: String::from("blocked"), + readiness_state: String::from("blocked"), + issue_identifier: Some(String::from("XY-853")), + issue_state: Some(String::from("Todo")), + queue_label_action: Some(String::from("remove")), + reason_codes: vec![String::from("dependency_not_terminal")], + reasons: vec![String::from( + "a dependency has not reached a required terminal state", + )], + next_action: String::from( + "Complete the dependency issue or refresh the Execution Program dependency plan if this remains stale.", + ), + }], readback_warning: None, }], worktrees: Vec::new(), @@ -206,8 +232,426 @@ fn operator_status_text_surfaces_execution_program_summary() { assert!(rendered.contains("Execution programs: 1")); assert!(rendered.contains("Execution Programs")); assert!(rendered.contains( - "program_id: program-853 source_contract_id: contract-852 nodes=3 planned=0 mapped=0 ready=1 queued=0 blocked=1 held=0 active=0 attention=0 completed=1 stale=0 superseded=0 queue_label_eligible=1 mapped_issues=XY-853" + "program_id: program-853 status=blocked source_contract_id: contract-852 intake_kind=goal_intake summary=\"Resolve promoted program work.\" nodes=3 planned=0 mapped=0 ready=1 queued=0 blocked=1 held=0 active=0 attention=0 completed=1 stale=0 superseded=0 queue_label_eligible=1 queue_actions=apply:1 retain:0 remove:1 mapped_issues=XY-853" )); + assert!(rendered.contains( + "node: issue=XY-853 issue_state=Todo lifecycle=blocked readiness=blocked queue_label_action=remove reason_codes=dependency_not_terminal reasons=\"a dependency has not reached a required terminal state\" next_action=\"Complete the dependency issue or refresh the Execution Program dependency plan if this remains stale.\"" + )); +} + +#[test] +fn operator_status_json_accepts_legacy_program_rows_without_new_fields() { + let snapshot: OperatorStatusSnapshot = serde_json::from_value(serde_json::json!({ + "project_id": "decodex", + "run_limit": 10, + "warnings": [], + "warning_details": [], + "connector_backoffs": [], + "projects": [], + "account_control": { + "mode": "balanced", + "account_selector": null, + }, + "accounts": [], + "active_runs": [], + "recent_runs": [], + "history_lanes": [], + "execution_programs": [{ + "program_id": "legacy-program", + "source_contract_id": null, + "node_count": 1, + "planned_count": 0, + "mapped_count": 0, + "ready_count": 1, + "queued_count": 0, + "blocked_count": 0, + "held_count": 0, + "active_count": 0, + "needs_attention_count": 0, + "completed_count": 0, + "stale_count": 0, + "superseded_count": 0, + "queue_label_eligible_count": 1, + "mapped_issue_identifiers": ["XY-853"], + "readback_warning": null, + }], + "queued_candidates": [], + "worktrees": [], + "post_review_lanes": [], + })) + .expect("legacy operator snapshot should deserialize"); + let program = snapshot.execution_programs.first().expect("program should deserialize"); + + assert_eq!(program.status, "unknown"); + assert_eq!(program.intake_kind, None); + assert_eq!(program.public_summary, None); + assert_eq!(program.queue_label_apply_count, 0); + assert_eq!(program.queue_label_retain_count, 0); + assert_eq!(program.queue_label_remove_count, 0); + assert!(program.node_readbacks.is_empty()); +} + +#[test] +fn operator_status_snapshot_surfaces_program_intake_and_node_readbacks() { + let (_temp_dir, config, workflow) = temp_project_layout(); + let state_store = StateStore::open_in_memory().expect("state store should open"); + + seed_program_readback_status(&state_store, &config); + + let snapshot = build_program_readback_snapshot(&config, &workflow, &state_store); + let program = snapshot.execution_programs.first().expect("program should surface"); + let program_json = program_readback_json(&snapshot); + + assert_program_readback_summary(program); + assert_program_readback_json(&program_json); + assert_program_node_readbacks(program, &program_json); +} + +fn seed_program_readback_status(state_store: &StateStore, config: &ServiceConfig) { + let program = ExecutionProgram::from_issue_batch_intake( + "program-status-readback", + config.service_id(), + "program-status-fingerprint", + "Coordinate status readback work.", + vec![ + status_program_node( + "node-ready", + "issue-ready", + "PUB-941", + "Todo", + ExecutionQueueIntent::ReadyToQueue, + false, + ), + status_program_node( + "node-queued", + "issue-queued", + "PUB-942", + "Todo", + ExecutionQueueIntent::Queued, + true, + ), + status_program_node_with_dependency( + "node-blocked", + "issue-blocked", + "PUB-943", + "Todo", + "PUB-944", + true, + ), + status_program_node( + "node-dependency", + "issue-dependency", + "PUB-944", + "Todo", + ExecutionQueueIntent::NotReady, + false, + ), + status_program_active_node("node-active", "issue-active", "PUB-945", "In Progress"), + ], + ) + .expect("program should build"); + + state_store + .upsert_execution_program(config.service_id(), program) + .expect("program should persist"); +} + +fn build_program_readback_snapshot( + config: &ServiceConfig, + workflow: &WorkflowDocument, + state_store: &StateStore, +) -> OperatorStatusSnapshot { + let tracker = FakeTracker::new(Vec::new()); + + orchestrator::build_live_operator_status_snapshot( + &tracker, + config, + workflow, + state_store, + 10, + ) + .expect("status snapshot should build") +} + +fn assert_program_readback_summary(program: &OperatorExecutionProgramStatus) { + assert_eq!(program.program_id, "program-status-readback"); + assert_eq!(program.status, "blocked"); + assert_eq!(program.intake_kind.as_deref(), Some("issue_batch_intake")); + assert_eq!(program.public_summary.as_deref(), Some("Coordinate status readback work.")); + assert_eq!(program.ready_count, 1); + assert_eq!(program.queued_count, 1); + assert_eq!(program.blocked_count, 1); + assert_eq!(program.held_count, 2); + assert_eq!(program.active_count, 1); + assert_eq!(program.stale_count, 0); + assert_eq!(program.queue_label_eligible_count, 2); + assert_eq!(program.queue_label_apply_count, 1); + assert_eq!(program.queue_label_retain_count, 1); + assert_eq!(program.queue_label_remove_count, 1); + assert_eq!( + program.mapped_issue_identifiers, + vec![ + String::from("PUB-941"), + String::from("PUB-942"), + String::from("PUB-943"), + String::from("PUB-944"), + String::from("PUB-945"), + ] + ); +} + +fn program_readback_json(snapshot: &OperatorStatusSnapshot) -> Value { + let snapshot_json = serde_json::to_value(snapshot).expect("snapshot should serialize"); + + snapshot_json["execution_programs"] + .as_array() + .expect("execution programs should serialize as an array") + .first() + .expect("program should serialize") + .clone() +} + +fn assert_program_readback_json(program_json: &Value) { + assert_eq!(program_json["program_id"], "program-status-readback"); + assert_eq!(program_json["status"], "blocked"); + assert_eq!(program_json["intake_kind"], "issue_batch_intake"); + assert_eq!(program_json["public_summary"], "Coordinate status readback work."); + assert_eq!(program_json["ready_count"], 1); + assert_eq!(program_json["queued_count"], 1); + assert_eq!(program_json["active_count"], 1); + assert_eq!(program_json["held_count"], 2); + assert_eq!(program_json["queue_label_apply_count"], 1); + assert_eq!(program_json["queue_label_retain_count"], 1); + assert_eq!(program_json["queue_label_remove_count"], 1); + assert!(program_json.get("contract").is_none()); + assert!(program_json.get("graph").is_none()); +} + +fn assert_program_node_readbacks( + program: &OperatorExecutionProgramStatus, + program_json: &Value, +) { + let node_by_issue = program + .node_readbacks + .iter() + .filter_map(|node| node.issue_identifier.as_deref().map(|issue| (issue, node))) + .collect::>(); + let ready_node = node_by_issue.get("PUB-941").expect("ready node should render"); + let queued_node = node_by_issue.get("PUB-942").expect("queued node should render"); + let blocked_node = node_by_issue.get("PUB-943").expect("blocked node should render"); + let held_node = node_by_issue.get("PUB-944").expect("held node should render"); + let active_node = node_by_issue.get("PUB-945").expect("active node should render"); + + assert_eq!(ready_node.queue_label_action.as_deref(), Some("apply")); + assert_eq!(queued_node.queue_label_action.as_deref(), Some("retain")); + assert_eq!(blocked_node.lifecycle_state, "blocked"); + assert_eq!(blocked_node.queue_label_action.as_deref(), Some("remove")); + assert!(blocked_node.reason_codes.contains(&String::from("dependency_not_terminal"))); + assert_eq!( + blocked_node.reasons, + vec![String::from("a dependency has not reached a required terminal state")] + ); + assert!(blocked_node.next_action.contains("Execution Program dependency plan")); + assert_eq!(held_node.lifecycle_state, "mapped"); + assert!(held_node.reason_codes.contains(&String::from("queue_intent_not_ready"))); + assert_eq!(active_node.lifecycle_state, "active"); + assert!(active_node.reason_codes.contains(&String::from("mapped_issue_active_label_present"))); + + let node_json = program_json["node_readbacks"] + .as_array() + .expect("node readbacks should serialize as an array") + .iter() + .find(|node| node["issue_identifier"] == "PUB-945") + .expect("active node json should serialize"); + + assert_eq!(node_json["lifecycle_state"], "active"); + assert_eq!(node_json["readiness_state"], "blocked"); + assert_eq!(node_json["queue_label_action"], serde_json::Value::Null); +} + +#[test] +fn operator_status_json_surfaces_missing_contract_program_recovery() { + let (_temp_dir, config, workflow) = temp_project_layout(); + let state_store = StateStore::open_in_memory().expect("state store should open"); + let contract = accepted_status_decision_contract_fixture(); + let program = ExecutionProgram::from_accepted_contract( + "program-missing-contract", + config.service_id(), + &contract, + vec![status_program_node( + "node-stale", + "issue-stale", + "PUB-946", + "Todo", + ExecutionQueueIntent::ReadyToQueue, + false, + )], + ) + .expect("program should build"); + + state_store + .upsert_execution_program(config.service_id(), program) + .expect("program should persist"); + + let tracker = FakeTracker::new(Vec::new()); + let snapshot = orchestrator::build_live_operator_status_snapshot( + &tracker, + &config, + &workflow, + &state_store, + 10, + ) + .expect("status snapshot should build"); + let program = snapshot.execution_programs.first().expect("program should surface"); + + assert_eq!(program.program_id, "program-missing-contract"); + assert_eq!(program.status, "stale"); + assert_eq!(program.source_contract_id.as_deref(), Some(contract.contract_id())); + assert_eq!(program.intake_kind.as_deref(), Some("goal_intake")); + assert_eq!(program.stale_count, 1); + assert_eq!(program.readback_warning.as_deref(), Some("source_decision_contract_missing")); + assert_eq!(program.mapped_issue_identifiers, vec![String::from("PUB-946")]); + + let node = program.node_readbacks.first().expect("stale node should render"); + + assert_eq!(node.lifecycle_state, "stale"); + assert_eq!(node.readiness_state, "stale"); + assert_eq!(node.issue_identifier.as_deref(), Some("PUB-946")); + assert_eq!(node.reason_codes, vec![String::from("source_decision_contract_missing")]); + assert!(node.next_action.contains("Decision Contract")); + + let snapshot_json = serde_json::to_value(&snapshot).expect("snapshot should serialize"); + let program_json = snapshot_json["execution_programs"] + .as_array() + .expect("execution programs should serialize as an array") + .first() + .expect("program should serialize"); + + assert_eq!(program_json["status"], "stale"); + assert_eq!(program_json["readback_warning"], "source_decision_contract_missing"); + assert_eq!(program_json["node_readbacks"][0]["reason_codes"][0], "source_decision_contract_missing"); + assert_eq!( + program_json["node_readbacks"][0]["next_action"], + "Restore or supersede the source Decision Contract before queueing this program." + ); + assert!(program_json.get("contract").is_none()); + assert!(program_json.get("decision_contract").is_none()); +} + +fn status_program_node( + node_id: &str, + issue_id: &str, + issue_identifier: &str, + issue_state: &str, + queue_intent: ExecutionQueueIntent, + program_owned_queue_label: bool, +) -> ExecutionProgramNode { + let mapping = status_program_issue_mapping( + issue_id, + issue_identifier, + issue_state, + program_owned_queue_label, + ); + + ExecutionProgramNode::new( + node_id, + ExecutionProgramNodeStage::Runtime, + format!("Resolve {issue_identifier}."), + queue_intent, + ) + .expect("node should build") + .with_acceptance_expectations([format!("{issue_identifier} acceptance is explicit.")]) + .expect("acceptance should attach") + .with_validation_expectations([String::from("Run focused validation.")]) + .expect("validation should attach") + .with_linear_issue(mapping) + .expect("mapping should attach") +} + +fn status_program_node_with_dependency( + node_id: &str, + issue_id: &str, + issue_identifier: &str, + issue_state: &str, + dependency_identifier: &str, + program_owned_queue_label: bool, +) -> ExecutionProgramNode { + status_program_node( + node_id, + issue_id, + issue_identifier, + issue_state, + ExecutionQueueIntent::ReadyToQueue, + program_owned_queue_label, + ) + .with_dependencies([ExecutionProgramDependency::new(dependency_identifier) + .expect("dependency should build")]) + .expect("dependency should attach") +} + +fn status_program_issue_mapping( + issue_id: &str, + issue_identifier: &str, + issue_state: &str, + program_owned_queue_label: bool, +) -> ExecutionLinearIssueMapping { + let mapping = ExecutionLinearIssueMapping::new(issue_id, issue_identifier, issue_state) + .expect("mapping should build"); + + if program_owned_queue_label { + mapping.with_program_owned_queue_label(true) + } else { + mapping + } +} + +fn status_program_active_node( + node_id: &str, + issue_id: &str, + issue_identifier: &str, + issue_state: &str, +) -> ExecutionProgramNode { + let mapping = status_program_issue_mapping(issue_id, issue_identifier, issue_state, false) + .with_active_label(true); + + ExecutionProgramNode::new( + node_id, + ExecutionProgramNodeStage::Runtime, + format!("Resolve {issue_identifier}."), + ExecutionQueueIntent::ReadyToQueue, + ) + .expect("node should build") + .with_acceptance_expectations([format!("{issue_identifier} acceptance is explicit.")]) + .expect("acceptance should attach") + .with_validation_expectations([String::from("Run focused validation.")]) + .expect("validation should attach") + .with_linear_issue(mapping) + .expect("mapping should attach") +} + +fn accepted_status_decision_contract_fixture() -> DecisionContract { + let mut contract: DecisionContract = serde_json::from_str(include_str!( + concat!( + env!("CARGO_MANIFEST_DIR"), + "/fixtures/decision_contract/research_x_latent_contract.json" + ) + )) + .expect("decision contract fixture should deserialize"); + + contract + .promote( + DecisionPromotion::new( + "operator", + DecisionPromotionActorKind::User, + "2026-06-09T10:00:00Z", + "conversation", + Some(String::from("User accepted the program boundary.")), + ) + .expect("promotion should build"), + ) + .expect("contract should promote"); + + contract } #[test] diff --git a/apps/decodex/src/orchestrator/types.rs b/apps/decodex/src/orchestrator/types.rs index 75ce96db2..d3e9eb16a 100644 --- a/apps/decodex/src/orchestrator/types.rs +++ b/apps/decodex/src/orchestrator/types.rs @@ -854,6 +854,7 @@ impl LoopGuardrailStopRequested { } } } + impl Display for LoopGuardrailStopRequested { fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result { let source = self.source_error_class.as_deref().unwrap_or("none"); @@ -1291,7 +1292,13 @@ struct OperatorProjectStatus { #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] struct OperatorExecutionProgramStatus { program_id: String, + #[serde(default = "operator_execution_program_unknown_status")] + status: String, source_contract_id: Option, + + intake_kind: Option, + + public_summary: Option, node_count: usize, planned_count: usize, mapped_count: usize, @@ -1305,17 +1312,33 @@ struct OperatorExecutionProgramStatus { stale_count: usize, superseded_count: usize, queue_label_eligible_count: usize, + #[serde(default)] + queue_label_apply_count: usize, + #[serde(default)] + queue_label_retain_count: usize, + #[serde(default)] + queue_label_remove_count: usize, mapped_issue_identifiers: Vec, + #[serde(default)] + node_readbacks: Vec, readback_warning: Option, } impl OperatorExecutionProgramStatus { fn from_summary( record: &ExecutionProgramRecord, summary: ExecutionProgramOperatorSummary, + evaluation: &ExecutionProgramEvaluation, ) -> Self { + let (queue_label_apply_count, queue_label_retain_count, queue_label_remove_count) = + operator_execution_program_queue_label_action_counts(evaluation); + let program_intake_plan = record.program().program_intake_plan(); + Self { - program_id: summary.program_id, + status: operator_execution_program_status(&summary, record.program().nodes().len(), None), + program_id: summary.program_id.clone(), source_contract_id: record.source_contract_id().map(str::to_owned), + intake_kind: program_intake_plan.map(|plan| plan.intake_kind().as_str().to_owned()), + public_summary: program_intake_plan.map(|plan| plan.public_summary().to_owned()), node_count: record.program().nodes().len(), planned_count: summary.planned_count, mapped_count: summary.mapped_count, @@ -1329,7 +1352,16 @@ impl OperatorExecutionProgramStatus { stale_count: summary.stale_count, superseded_count: summary.superseded_count, queue_label_eligible_count: summary.queue_label_eligible_count, + queue_label_apply_count, + queue_label_retain_count, + queue_label_remove_count, mapped_issue_identifiers: summary.mapped_issue_identifiers, + node_readbacks: evaluation + .nodes() + .iter() + .filter(|node| operator_execution_program_node_should_render(node)) + .map(operator_execution_program_node_readback) + .collect(), readback_warning: None, } } @@ -1339,7 +1371,16 @@ impl OperatorExecutionProgramStatus { Self { program_id: record.program_id().to_owned(), + status: String::from("stale"), source_contract_id: record.source_contract_id().map(str::to_owned), + intake_kind: record + .program() + .program_intake_plan() + .map(|plan| plan.intake_kind().as_str().to_owned()), + public_summary: record + .program() + .program_intake_plan() + .map(|plan| plan.public_summary().to_owned()), node_count, planned_count: 0, mapped_count: 0, @@ -1353,12 +1394,28 @@ impl OperatorExecutionProgramStatus { stale_count: node_count, superseded_count: 0, queue_label_eligible_count: 0, - mapped_issue_identifiers: Vec::new(), + queue_label_apply_count: 0, + queue_label_retain_count: 0, + queue_label_remove_count: 0, + mapped_issue_identifiers: operator_execution_program_mapped_issue_identifiers(record), + node_readbacks: operator_execution_program_missing_contract_nodes(record), readback_warning: Some(String::from("source_decision_contract_missing")), } } } +#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] +struct OperatorExecutionProgramNodeStatus { + lifecycle_state: String, + readiness_state: String, + issue_identifier: Option, + issue_state: Option, + queue_label_action: Option, + reason_codes: Vec, + reasons: Vec, + next_action: String, +} + #[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)] struct OperatorGitHubCliAuthority { command_path: String, @@ -2165,6 +2222,279 @@ pub(crate) fn record_authority_decision_request_private_event( ) } +fn operator_execution_program_status( + summary: &ExecutionProgramOperatorSummary, + node_count: usize, + readback_warning: Option<&str>, +) -> String { + if readback_warning.is_some() || summary.stale_count > 0 || summary.superseded_count > 0 { + String::from("stale") + } else if summary.needs_attention_count > 0 { + String::from("attention") + } else if summary.blocked_count > 0 { + String::from("blocked") + } else if summary.active_count > 0 { + String::from("active") + } else if summary.queued_count > 0 { + String::from("queued") + } else if summary.ready_count > 0 { + String::from("ready") + } else if node_count > 0 && summary.completed_count == node_count { + String::from("completed") + } else if summary.held_count > 0 { + String::from("held") + } else { + String::from("idle") + } +} + +fn operator_execution_program_unknown_status() -> String { + String::from("unknown") +} + +fn operator_execution_program_queue_label_action_counts( + evaluation: &ExecutionProgramEvaluation, +) -> (usize, usize, usize) { + let mut apply_count = 0; + let mut retain_count = 0; + let mut remove_count = 0; + + for node in evaluation.nodes() { + match node.queue_label_action() { + Some(ExecutionQueueLabelAction::Apply) => apply_count += 1, + Some(ExecutionQueueLabelAction::Retain) => { + retain_count += 1 + }, + Some(ExecutionQueueLabelAction::Remove) => { + remove_count += 1 + }, + None => {}, + } + } + + (apply_count, retain_count, remove_count) +} + +fn operator_execution_program_node_should_render( + node: &ExecutionNodeEvaluation, +) -> bool { + node.queue_label_action().is_some() + || matches!( + node.lifecycle_state(), + crate::execution_program::ExecutionProgramNodeLifecycleState::Active + | crate::execution_program::ExecutionProgramNodeLifecycleState::Blocked + | crate::execution_program::ExecutionProgramNodeLifecycleState::Mapped + | crate::execution_program::ExecutionProgramNodeLifecycleState::NeedsAttention + | crate::execution_program::ExecutionProgramNodeLifecycleState::Planned + | crate::execution_program::ExecutionProgramNodeLifecycleState::Stale + | crate::execution_program::ExecutionProgramNodeLifecycleState::Superseded + ) +} + +fn operator_execution_program_node_readback( + node: &ExecutionNodeEvaluation, +) -> OperatorExecutionProgramNodeStatus { + let reason_codes = operator_execution_program_reason_codes(node.reasons()); + let reasons = node + .reasons() + .iter() + .map(|reason| operator_execution_program_public_reason(reason)) + .collect::>(); + let issue = node.linear_issue(); + + OperatorExecutionProgramNodeStatus { + lifecycle_state: node.lifecycle_state().as_str().to_owned(), + readiness_state: node.state().as_str().to_owned(), + issue_identifier: issue.map(|issue| issue.issue_identifier().to_owned()), + issue_state: issue.map(|issue| issue.issue_state().to_owned()), + queue_label_action: node.queue_label_action().map(|action| action.as_str().to_owned()), + next_action: operator_execution_program_node_next_action(node, &reason_codes), + reason_codes, + reasons, + } +} + +fn operator_execution_program_missing_contract_nodes( + record: &ExecutionProgramRecord, +) -> Vec { + record + .program() + .nodes() + .iter() + .map(|node| { + let issue = node.linear_issue(); + + OperatorExecutionProgramNodeStatus { + lifecycle_state: String::from("stale"), + readiness_state: String::from("stale"), + issue_identifier: issue.map(|issue| issue.issue_identifier().to_owned()), + issue_state: issue.map(|issue| issue.issue_state().to_owned()), + queue_label_action: None, + reason_codes: vec![String::from("source_decision_contract_missing")], + reasons: vec![String::from("source Decision Contract is missing")], + next_action: String::from( + "Restore or supersede the source Decision Contract before queueing this program.", + ), + } + }) + .collect() +} + +fn operator_execution_program_mapped_issue_identifiers( + record: &ExecutionProgramRecord, +) -> Vec { + let mut identifiers = record + .program() + .nodes() + .iter() + .filter_map(|node| node.linear_issue().map(|issue| issue.issue_identifier().to_owned())) + .collect::>(); + + identifiers.sort(); + identifiers.dedup(); + + identifiers +} + +fn operator_execution_program_reason_codes(reasons: &[String]) -> Vec { + let mut seen = BTreeSet::new(); + + for reason in reasons { + seen.insert(operator_execution_program_reason_code(reason).to_owned()); + } + + seen.into_iter().collect() +} + +fn operator_execution_program_reason_code(reason: &str) -> &'static str { + if reason == "node no longer matches the accepted Decision Contract" { + "accepted_contract_mismatch" + } else if reason == "node queue intent is not-ready" { + "queue_intent_not_ready" + } else if reason == "node queue intent is paused" { + "queue_intent_paused" + } else if reason == "node already has an active lane" { + "active_lane_present" + } else if reason == "node queue intent is terminal" { + "queue_intent_terminal" + } else if reason == "node is ready for normal Linear issue execution" { + "ready_for_linear_execution" + } else if reason == "node has no acceptance expectations" { + "acceptance_expectations_missing" + } else if reason == "node has no validation expectations" { + "validation_expectations_missing" + } else if reason.starts_with("dependency `") { + "dependency_not_terminal" + } else if reason.starts_with("conflict domain `") { + "conflict_domain_occupied" + } else if reason == "node has no normal Linear issue mapping" { + "linear_issue_mapping_missing" + } else if reason.contains(" is already terminal in `") { + "mapped_issue_terminal" + } else if reason.contains(" is not in a startable state") { + "mapped_issue_not_startable" + } else if reason.contains(" already carries `") { + "mapped_issue_active_label_present" + } else if reason.contains(" carries `decodex:manual-only`") { + "mapped_issue_manual_only" + } else if reason.contains(" carries `decodex:needs-attention`") { + "mapped_issue_needs_attention" + } else if reason.contains(" has open tracker dependency blockers") { + "mapped_issue_open_blockers" + } else if reason.contains(" is missing a generic dispatch briefing") { + "mapped_issue_dispatch_briefing_missing" + } else { + "program_readiness_blocked" + } +} + +fn operator_execution_program_public_reason(reason: &str) -> String { + if reason.starts_with("conflict domain `") { + String::from("another active or retained program node occupies this conflict domain") + } else if reason.starts_with("dependency `") { + String::from("a dependency has not reached a required terminal state") + } else { + reason.to_owned() + } +} + +fn operator_execution_program_node_next_action( + node: &ExecutionNodeEvaluation, + reason_codes: &[String], +) -> String { + if matches!( + node.lifecycle_state(), + crate::execution_program::ExecutionProgramNodeLifecycleState::Stale + | crate::execution_program::ExecutionProgramNodeLifecycleState::Superseded + ) { + return String::from( + "Refresh or supersede the accepted Decision Contract before queueing this program.", + ); + } + if reason_codes.iter().any(|code| code == "dependency_not_terminal") { + return String::from( + "Complete the dependency issue or refresh the Execution Program dependency plan if this remains stale.", + ); + } + if matches!( + node.lifecycle_state(), + crate::execution_program::ExecutionProgramNodeLifecycleState::NeedsAttention + ) + || reason_codes.iter().any(|code| code == "mapped_issue_needs_attention") + { + return String::from( + "Resolve the mapped issue's needs-attention stop before queueing this node.", + ); + } + if matches!( + node.lifecycle_state(), + crate::execution_program::ExecutionProgramNodeLifecycleState::Active + ) { + return String::from( + "Wait for the active lane or recover its retained state before queueing this node.", + ); + } + if matches!( + node.queue_label_action(), + Some(crate::execution_program::ExecutionQueueLabelAction::Remove) + ) { + return String::from( + "Allow the program reconciler to remove its owned queue label on the next Linear scan.", + ); + } + if matches!( + node.queue_label_action(), + Some(crate::execution_program::ExecutionQueueLabelAction::Apply) + ) { + return String::from( + "Wait for the next Linear scan or request a targeted scan to apply the service queue label.", + ); + } + if matches!( + node.queue_label_action(), + Some(crate::execution_program::ExecutionQueueLabelAction::Retain) + ) { + return String::from("Issue is already queued by the program reconciler."); + } + if matches!( + node.lifecycle_state(), + crate::execution_program::ExecutionProgramNodeLifecycleState::Planned + | crate::execution_program::ExecutionProgramNodeLifecycleState::Mapped + ) { + return String::from("Map, promote, or unpause this intake node before queueing it."); + } + if matches!( + node.lifecycle_state(), + crate::execution_program::ExecutionProgramNodeLifecycleState::Blocked + ) { + return String::from( + "Repair mapped issue blockers, briefing, or program readiness before retrying.", + ); + } + + String::from("No operator action required.") +} + fn validate_authority_boundary_check_input( input: &AuthorityBoundaryCheckInput<'_>, ) -> Result<()> { diff --git a/docs/reference/operator-control-plane.md b/docs/reference/operator-control-plane.md index 2dd7a418c..886c03642 100644 --- a/docs/reference/operator-control-plane.md +++ b/docs/reference/operator-control-plane.md @@ -41,7 +41,12 @@ Decodex currently runs as a local, single-machine control plane: queue-label-eligible nodes plus mapped issue identifiers. Optional planned, mapped, active, and superseded counts may appear when the runtime has that detail. Low-level node edges, graph operations, and queue-label ownership records remain internal - runtime state. + runtime state. Status JSON may include sparse node readbacks for queue-label + decisions and held, blocked, stale, active, or attention-bound nodes: mapped issue + identifier, issue state, lifecycle/readiness state, queue-label action, public-safe + reason codes, public-safe reasons, and a recovery next action. It must not expose + Decision Contract payloads, raw graph edges, local paths, credentials, raw logs, or + private runtime events. - Persisted Execution Programs are reconciled during normal Linear scan ticks before issue selection. Ready, startable, program-owned mapped nodes can receive the service queue label automatically; blocked, stale, paused, terminal, @@ -231,6 +236,15 @@ the command falls back to a direct local runtime read and emits always bypasses the cached snapshot and rebuilds status with fresh Linear/GitHub observers. +Operator JSON snapshots include `execution_programs[]` for Program Intake and +Execution Program readback. Each program row carries public intake kind/summary, +source contract id when present, the compact summary counts, mapped issue +identifiers, queue-label action counts (`apply`, `retain`, `remove`), and sparse +`node_readbacks[]` for queue decisions or nodes that need operator context. Dependency +diagnostics use `dependency_not_terminal` with a next action to complete the +dependency issue or refresh the Execution Program dependency plan when a stale +dependency program is the real blocker. + For the lane-control rollout, active-lane UI posture is observe-only. The dashboard renders active-lane state, protocol activity, liveness, private-evidence references, local run-control capability metadata, and local acknowledgement/account controls, but @@ -272,6 +286,7 @@ protocol activity durable outside the local operator surface. | `Accounts` | Shared Codex account pool and usage table from `~/.codex/decodex/accounts.jsonl` when `[codex.accounts]` is enabled for a project. Account identity can be obscured from the `Account` column header eye without changing the underlying snapshot. The row weight column shows the capacity multiplier used for pool usage estimates: `pro` accounts count as `20x`, and all other plans count as `1x`. Usage probes read Codex `/wham/usage` for window capacity and `/wham/profiles/me` for profile token stats such as lifetime tokens, peak daily tokens, longest task, streaks, and daily token activity. Selecting an account writes the global `[codex.accounts].fixed_account` selector in `~/.codex/decodex/config.toml`; clearing it returns all new account-pool runs to balanced account selection. Account display-name rerolls write `[codex.account_names.offsets]` in the same global config so Decodex App and the dashboard share the privacy-preserving names. Theme, sort, and identity-visibility preferences are client-local presentation state. The selector is global and does not pin a project to an account. | | `Projects` | Fleet-level project table. The section-level filter toggles between active project work and the full registry. Location is its own compact path column and can be obscured from the location header eye. `Activity` shows a relative timestamp or `-`; `Work` is `running/waiting/attention`. It should not duplicate per-lane details already shown below. | | `Running Lanes` | Active leased or live-executing issue lanes. A lane here is currently owned by this local control plane, or a live process/thread/protocol marker still explains active execution even when the queue lease is not held. It shows issue identity, phase, operation, attempt, queue lease state, execution liveness, thread/protocol status, child-agent activity when captured, phase-goal status when app-server reports it, timing, branch, and worktree. | +| `Program Intake` | Read-only Program Intake and Execution Program progress. It shows each active program's public intake summary, status, mapped issue identifiers, summary counts, queue-label action counts, and sparse node diagnostics for queue decisions or held/blocked/stale/attention nodes. It does not expose graph editing, raw node-edge mutation controls, Decision Contract payloads, or private runtime evidence. | | `Intake Queue` | Queued tracker issues before execution. Candidates are classified as `ready`, capacity-waiting, claimed without a matching local lane, blocked, or closed/stale. Repeated identical open dependency blockers surface as `dependency_program_stale` after the guardrail threshold so operators can distinguish a stale Execution Program/dependency plan from a newly blocked queue item. A blocked queued candidate can still show an attached `.worktrees/XY-*` path when the queue owns the attention state; if that worktree has tracked changes after stalled reconciliation, failure writeback, or retries, the candidate is partial retained progress and not just a generic stalled or retry-budget hold. Human-required authority stops expose their compact decision request fields here: `phase = human_required`, reason, boundary, `decision_request_id`, and `next_action`. When queued attention still maps to a run/attempt, it also carries the same compact loop status used by running lanes. Running lanes are not repeated as normal intake work. | | `Review & Landing` | Retained PR lanes after review handoff. This section owns post-review repair, wait-for-review, ready-to-land, closeout, cleanup, and blocked retained-lane visibility. Retained lanes expose compact loop status for their bound handoff run/attempt so operators can see review repair checkpoint state, architecture recovery stops, and boundary/human-required disposition without direct SQLite inspection. | | `Recovery Worktrees` | Retained local worktrees that are not currently owned by `Running Lanes`, `Review & Landing`, or queued attention in `Intake Queue`. This is the cleanup or recovery inbox for recovered paths, retained PR leftovers, and cleanup-only local worktrees. Empty is the normal healthy state. |