Skip to content

Commit 9cf5762

Browse files
committed
fix: report live Lab cook progress accurately
1 parent bfe5504 commit 9cf5762

5 files changed

Lines changed: 326 additions & 7 deletions

File tree

src/commands/runner/status.rs

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -918,6 +918,15 @@ pub(super) fn runner_status_operator_hints(report: &RunnerStatusReport) -> Vec<S
918918
"Active-job status for `{}` is unavailable: {reason}. Treat active_job_count=0 as unknown, not idle.",
919919
report.runner_id
920920
));
921+
} else if let Some(error) = report
922+
.active_job_error
923+
.as_ref()
924+
.filter(|error| error.code == "active_job_view_inconsistent")
925+
{
926+
hints.push(format!(
927+
"Active-job view for `{}` is inconsistent: {}",
928+
report.runner_id, error.message
929+
));
921930
}
922931
if report.stale_runner_job_count > 0 {
923932
hints.push(format!(

src/core/agent_task_lifecycle/lifecycle_ops.rs

Lines changed: 132 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -500,6 +500,22 @@ pub(crate) fn reconcile_runner_job_snapshot(
500500
let metadata = reconciled.ensure_metadata_object();
501501
metadata.insert("runner_job_status".to_string(), json!(snapshot.job.status));
502502
metadata.insert("runner_job_last_seen_at".to_string(), json!(last_seen_at));
503+
metadata.insert("runner_job_events".to_string(), json!(snapshot.events));
504+
metadata.insert("phase".to_string(), json!("executing"));
505+
metadata.insert(
506+
"phase_activity".to_string(),
507+
json!("provider/executor process is active"),
508+
);
509+
metadata.insert("provider_state".to_string(), json!("active"));
510+
if let Some(provider) = metadata
511+
.get("provider_rotation")
512+
.and_then(|rotation| rotation.get("entries"))
513+
.and_then(Value::as_array)
514+
.and_then(|entries| entries.first())
515+
{
516+
metadata.insert("active_provider".to_string(), provider.clone());
517+
}
518+
merge_live_provider_handles(&mut reconciled, &snapshot.events);
503519
store::write_record(&reconciled)?;
504520
}
505521
crate::core::api_jobs::JobStatus::Succeeded
@@ -530,6 +546,52 @@ pub(crate) fn reconcile_runner_job_snapshot(
530546
Ok(())
531547
}
532548

549+
fn merge_live_provider_handles(
550+
record: &mut AgentTaskRunRecord,
551+
events: &[crate::core::api_jobs::JobEvent],
552+
) {
553+
for handle in events.iter().filter_map(|event| {
554+
event
555+
.data
556+
.as_ref()
557+
.and_then(|data| {
558+
data.pointer("/metadata/provider_handle")
559+
.or_else(|| data.get("provider_handle"))
560+
})
561+
.and_then(provider_handle_from_value)
562+
}) {
563+
if record
564+
.provider_handles
565+
.iter()
566+
.any(|existing| existing.provider_run_id == handle.run_id)
567+
{
568+
continue;
569+
}
570+
record.provider_handles.push(AgentTaskRunProviderHandle {
571+
kind: handle.kind,
572+
task_id: handle.task_id,
573+
backend: handle.backend,
574+
provider_run_id: handle.run_id,
575+
stream_uri: handle.stream_uri,
576+
state: Some(AgentTaskState::Running),
577+
metadata: handle.metadata,
578+
});
579+
}
580+
if !record.provider_handles.is_empty() {
581+
record.lifecycle.provider_runtime = record
582+
.provider_handles
583+
.iter()
584+
.map(provider_runtime_for_handle)
585+
.collect();
586+
record.lifecycle.external_runtime_ids = record
587+
.lifecycle
588+
.provider_runtime
589+
.iter()
590+
.flat_map(|runtime| runtime.external_runtime_ids.clone())
591+
.collect();
592+
}
593+
}
594+
533595
fn validate_runner_job_snapshot(
534596
record: &AgentTaskRunRecord,
535597
snapshot: &crate::core::runner::RunnerJobLogSnapshot,
@@ -845,8 +907,9 @@ pub fn record_lab_offload_phase(
845907
durable_plan,
846908
)?;
847909
record.updated_at = Some(now_timestamp());
910+
let phase_started_at = record.updated_at.clone().unwrap_or_else(now_timestamp);
848911
let metadata = record.ensure_metadata_object();
849-
metadata.insert("phase".to_string(), json!(phase));
912+
record_lab_offload_phase_metadata(metadata, phase, &phase_started_at);
850913
metadata.insert("provider_state".to_string(), json!("pending"));
851914
if let Some(remote_workspace) = remote_workspace {
852915
metadata.insert("remote_workspace".to_string(), json!(remote_workspace));
@@ -875,8 +938,9 @@ pub fn record_lab_offload_phase_executions(
875938
.filter(|id| !id.trim().is_empty())
876939
.collect();
877940
record.updated_at = Some(now_timestamp());
941+
let phase_started_at = record.updated_at.clone().unwrap_or_else(now_timestamp);
878942
let metadata = record.ensure_metadata_object();
879-
metadata.insert("phase".to_string(), json!(phase));
943+
record_lab_offload_phase_metadata(metadata, phase, &phase_started_at);
880944
metadata.insert(
881945
"materialization_execution_ids".to_string(),
882946
json!(execution_ids),
@@ -889,6 +953,44 @@ pub fn record_lab_offload_phase_executions(
889953
Ok(record)
890954
}
891955

956+
fn record_lab_offload_phase_metadata(
957+
metadata: &mut serde_json::Map<String, Value>,
958+
phase: &str,
959+
started_at: &str,
960+
) {
961+
let previous_phase = metadata
962+
.get("phase")
963+
.and_then(Value::as_str)
964+
.map(str::to_string);
965+
if previous_phase.as_deref() != Some(phase) {
966+
if let Some(previous_phase) = previous_phase {
967+
if let Some(entry) = metadata
968+
.get_mut("phase_history")
969+
.and_then(Value::as_array_mut)
970+
.and_then(|entries| {
971+
entries.iter_mut().rev().find(|entry| {
972+
entry.get("phase").and_then(Value::as_str) == Some(previous_phase.as_str())
973+
&& entry.get("ended_at").is_none()
974+
})
975+
})
976+
{
977+
entry["ended_at"] = json!(started_at);
978+
}
979+
}
980+
metadata
981+
.entry("phase_history".to_string())
982+
.or_insert_with(|| json!([]))
983+
.as_array_mut()
984+
.expect("phase history is an array")
985+
.push(json!({ "phase": phase, "started_at": started_at }));
986+
}
987+
metadata.insert("phase".to_string(), json!(phase));
988+
metadata.insert(
989+
"phase_activity".to_string(),
990+
json!(format!("Homeboy {phase}")),
991+
);
992+
}
993+
892994
pub fn record_detached_lab_run(input: DetachedLabRunRecord<'_>) -> Result<AgentTaskRunRecord> {
893995
let run_id = sanitize_run_id(input.run_id);
894996
let plan = detached_lab_plan(&run_id, &input);
@@ -1209,14 +1311,19 @@ pub fn retry(run_id: &str, requested_run_id: Option<&str>) -> Result<AgentTaskRu
12091311
}
12101312

12111313
pub fn logs(run_id: &str) -> Result<AgentTaskRunLog> {
1212-
let run_id = resolve_run_id(run_id)?;
1213-
let record = store::read_record(&run_id)?;
1314+
// Status reconciliation fetches the live daemon snapshot for a bound Lab
1315+
// child, making executor progress visible before the child is terminal.
1316+
let record = status(run_id)?;
1317+
let run_id = record.run_id.clone();
12141318
let (events, artifact_refs) = match store::read_aggregate(&run_id) {
12151319
Ok(aggregate) => {
12161320
let refs = artifact_refs_for_outcomes(&aggregate.outcomes);
12171321
(aggregate.events, refs)
12181322
}
1219-
Err(_) => (queued_events(&record.tasks), record.artifact_refs.clone()),
1323+
Err(_) => (
1324+
runner_job_progress_events(&record).unwrap_or_else(|| queued_events(&record.tasks)),
1325+
record.artifact_refs.clone(),
1326+
),
12201327
};
12211328
let normalized_events = normalize_progress_events(&run_id, &events, &artifact_refs);
12221329
Ok(AgentTaskRunLog {
@@ -1227,6 +1334,26 @@ pub fn logs(run_id: &str) -> Result<AgentTaskRunLog> {
12271334
})
12281335
}
12291336

1337+
fn runner_job_progress_events(record: &AgentTaskRunRecord) -> Option<Vec<AgentTaskProgressEvent>> {
1338+
let events = record.metadata.get("runner_job_events")?.as_array()?;
1339+
let task_id = record
1340+
.tasks
1341+
.first()
1342+
.map(|task| task.task_id.clone())
1343+
.unwrap_or_else(|| record.run_id.clone());
1344+
Some(
1345+
events
1346+
.iter()
1347+
.map(|event| AgentTaskProgressEvent {
1348+
task_id: task_id.clone(),
1349+
state: AgentTaskState::Running,
1350+
attempt: 0,
1351+
message: serde_json::to_string(event).ok(),
1352+
})
1353+
.collect(),
1354+
)
1355+
}
1356+
12301357
pub fn artifacts(run_id: &str) -> Result<AgentTaskRunArtifacts> {
12311358
let run_id = resolve_run_id(run_id)?;
12321359
let record = store::read_record(&run_id)?;

src/core/agent_task_lifecycle/tests.rs

Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ use crate::core::agent_task_scheduler::{
1212
AgentTaskAggregate, AgentTaskAggregateStatus, AgentTaskAggregateTotals,
1313
AGENT_TASK_AGGREGATE_SCHEMA,
1414
};
15+
use crate::core::api_jobs::{JobEvent, JobEventKind};
1516
use crate::test_support::with_isolated_home;
1617

1718
#[test]
@@ -444,6 +445,47 @@ fn controller_proxy_records_pre_execution_phase_progress() {
444445
let loaded = status("agent-task-pre-execution").expect("status resolves during setup");
445446
assert_eq!(loaded.metadata["phase"], "hydrating");
446447
assert_eq!(loaded.metadata["provider_state"], "pending");
448+
let phases = loaded.metadata["phase_history"]
449+
.as_array()
450+
.expect("phase history");
451+
assert_eq!(phases.len(), 2);
452+
assert_eq!(phases[0]["phase"], "materializing");
453+
assert!(phases[0].get("started_at").is_some());
454+
assert!(phases[0].get("ended_at").is_some());
455+
assert_eq!(phases[1]["phase"], "hydrating");
456+
assert!(phases[1].get("started_at").is_some());
457+
});
458+
}
459+
460+
#[test]
461+
fn logs_expose_mirrored_live_runner_events_before_terminal_aggregate() {
462+
with_isolated_home(|_| {
463+
let command = vec!["homeboy".to_string(), "agent-task".to_string()];
464+
let mut record = record_detached_lab_run(DetachedLabRunRecord {
465+
run_id: "live-runner-events",
466+
runner_id: "homeboy-lab",
467+
runner_job_id: "job-live",
468+
remote_workspace: "/runner/workspace/homeboy",
469+
remote_command: &command,
470+
})
471+
.expect("running proxy");
472+
record.metadata["runner_job_events"] = json!([JobEvent {
473+
sequence: 1,
474+
job_id: uuid::Uuid::new_v4(),
475+
kind: JobEventKind::Progress,
476+
timestamp_ms: 42,
477+
message: Some("provider started".to_string()),
478+
data: Some(json!({"provider": "openai/gpt-5.6-terra"})),
479+
}]);
480+
store::write_record(&record).expect("persist mirrored event");
481+
482+
let log = logs("live-runner-events").expect("live logs resolve");
483+
484+
assert_eq!(log.events.len(), 1);
485+
assert!(log.events[0]
486+
.message
487+
.as_deref()
488+
.is_some_and(|message| message.contains("provider started")));
447489
});
448490
}
449491

@@ -645,6 +687,58 @@ fn reachable_running_child_clears_disconnected_liveness_and_refreshes_heartbeat(
645687
});
646688
}
647689

690+
#[test]
691+
fn running_child_snapshot_persists_provider_handle_and_live_log_progress() {
692+
with_isolated_home(|_| {
693+
let command = vec!["homeboy".to_string(), "agent-task".to_string()];
694+
let mut record = record_detached_lab_run(DetachedLabRunRecord {
695+
run_id: "agent-task-live-provider",
696+
runner_id: "homeboy-lab",
697+
runner_job_id: "00000000-0000-0000-0000-000000000123",
698+
remote_workspace: "/runner/workspace/repo",
699+
remote_command: &command,
700+
})
701+
.expect("running proxy");
702+
let mut snapshot = terminal_child_snapshot(&succeeded_aggregate(&test_plan()));
703+
snapshot.job.status = crate::core::api_jobs::JobStatus::Running;
704+
snapshot.events = vec![crate::core::api_jobs::JobEvent {
705+
sequence: 1,
706+
job_id: snapshot.job.id,
707+
kind: crate::core::api_jobs::JobEventKind::Progress,
708+
timestamp_ms: 2,
709+
message: Some("provider dispatch accepted".to_string()),
710+
data: Some(json!({
711+
"metadata": {
712+
"provider_handle": AgentTaskExecutionHandle {
713+
kind: crate::core::agent_task::AgentTaskExecutionHandleKind::ProviderRun,
714+
task_id: "task-a".to_string(),
715+
backend: "openai/gpt-5.6-terra".to_string(),
716+
run_id: "provider-run-live".to_string(),
717+
stream_uri: Some("provider://runs/provider-run-live/events".to_string()),
718+
metadata: json!({"progress": "accepted"}),
719+
}
720+
}
721+
})),
722+
}];
723+
724+
reconcile_runner_job_snapshot(&mut record, &snapshot).expect("live reconciliation");
725+
726+
assert_eq!(record.metadata["phase"], "executing");
727+
assert_eq!(record.metadata["provider_state"], "active");
728+
assert_eq!(record.provider_handles.len(), 1);
729+
assert_eq!(
730+
record.provider_handles[0].provider_run_id,
731+
"provider-run-live"
732+
);
733+
let log = logs(&record.run_id).expect("live logs");
734+
assert_eq!(log.events.len(), 1);
735+
assert!(log.events[0]
736+
.message
737+
.as_deref()
738+
.is_some_and(|message| message.contains("provider dispatch accepted")));
739+
});
740+
}
741+
648742
#[test]
649743
fn terminal_child_projection_rejects_mismatched_persisted_run_identity() {
650744
with_isolated_home(|_| {

0 commit comments

Comments
 (0)