Skip to content

Commit 1308406

Browse files
authored
fix: preserve reverse broker handoff transport (#10573)
1 parent cc776d1 commit 1308406

10 files changed

Lines changed: 145 additions & 21 deletions

File tree

crates/homeboy-core/src/test_support.rs

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1115,6 +1115,18 @@ fn handle_reverse_broker_request(
11151115
.expect("claim broker job");
11161116
return ok(json!({ "claim": claim }));
11171117
}
1118+
if request.method == "POST" && request.path == "/runner/jobs/reconcile" {
1119+
let reconciled = store
1120+
.reconcile_expired_remote_runner_claims_for_runner(
1121+
chrono::Utc::now().timestamp_millis().max(0) as u64,
1122+
Some(runner_id),
1123+
)
1124+
.expect("reconcile broker jobs");
1125+
return ok(json!({
1126+
"reconciled_count": reconciled.len(),
1127+
"reconciled": reconciled,
1128+
}));
1129+
}
11181130
// Reverse-runner file transfer (`RunnerFileChannel::BrokerHttp`) posts
11191131
// `/files/{mkdir,upload,download}` for the same host the fixture runs on, so
11201132
// the fixture performs the real filesystem operation. Without these the

crates/homeboy-lab-runner/src/connection.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1828,14 +1828,14 @@ pub(crate) fn reverse_broker_job_snapshot_at(
18281828
let job_data = broker_http::get_json(
18291829
&client,
18301830
broker_url,
1831-
&format!("/runner/jobs/{job_id}"),
1831+
&format!("/jobs/{job_id}"),
18321832
"fetch reverse runner broker job",
18331833
token.as_deref(),
18341834
)?;
18351835
let events_data = broker_http::get_json(
18361836
&client,
18371837
broker_url,
1838-
&format!("/runner/jobs/{job_id}/events"),
1838+
&format!("/jobs/{job_id}/events"),
18391839
"fetch reverse runner broker job events",
18401840
token.as_deref(),
18411841
)?;

crates/homeboy-lab-runner/src/evidence/mirror.rs

Lines changed: 41 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -254,11 +254,13 @@ pub fn mirror_reverse_broker_evidence(
254254
let terminal_result: homeboy_core::api_jobs::RemoteRunnerJobResult =
255255
serde_json::from_value(result.clone()).map_err(|error| {
256256
Error::validation_invalid_argument(
257-
"result.observation_run_details",
258-
"reverse runner terminal result is not a valid typed observation detail contract",
259-
Some(error.to_string()),
260-
None,
261-
)
257+
"result.observation_run_details",
258+
format!(
259+
"reverse runner terminal result is not a valid typed observation detail contract: {error}"
260+
),
261+
None,
262+
None,
263+
)
262264
})?;
263265
let runs;
264266
let local_artifacts = if job.status == JobStatus::Failed {
@@ -512,6 +514,40 @@ pub fn mirror_daemon_job_progress(
512514
.map(|run| run.run)
513515
}
514516

517+
pub fn mirror_reverse_broker_job_progress(
518+
runner: &Runner,
519+
broker_url: &str,
520+
cwd: &str,
521+
command: &[String],
522+
job: &Job,
523+
run_id: Option<&str>,
524+
) -> Result<RunRecord> {
525+
let store = ObservationStore::open_initialized()?;
526+
let run = mirror_job_run(
527+
&store,
528+
runner,
529+
cwd,
530+
command,
531+
job,
532+
&[],
533+
&json!({}),
534+
run_id,
535+
None,
536+
)?
537+
.run;
538+
record_reverse_broker_metadata(
539+
&store,
540+
run,
541+
runner,
542+
broker_url,
543+
job,
544+
&[],
545+
&json!({}),
546+
Vec::new(),
547+
None,
548+
)
549+
}
550+
515551
/// Records that the controller can no longer observe an accepted runner job.
516552
/// The remote job may still exist, but the controller-side lifecycle is terminal
517553
/// and includes the polling diagnostic instead of leaving a stale running mirror.

crates/homeboy-lab-runner/src/evidence/mod.rs

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,8 @@ pub use homeboy_core::api_jobs::RunnerJobLogSnapshot;
2222
pub use mirror::runner_job_log_snapshot_for_session;
2323
pub use mirror::{
2424
controller_artifact_metadata, mirror_connected_runner_run, mirror_daemon_evidence,
25-
mirror_daemon_job_progress, mirror_reverse_broker_evidence, mirrored_runner_job_identity,
26-
refresh_mirrored_daemon_evidence, runner_job_log_snapshot, terminalize_mirrored_daemon_job,
25+
mirror_daemon_job_progress, mirror_reverse_broker_evidence, mirror_reverse_broker_job_progress,
26+
mirrored_runner_job_identity, refresh_mirrored_daemon_evidence, runner_job_log_snapshot,
27+
terminalize_mirrored_daemon_job,
2728
};
2829
pub(crate) use util::{local_job_run_id, runner_exec_run_label};

crates/homeboy-lab-runner/src/evidence/tests/mod.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1566,8 +1566,8 @@ fn reverse_broker_refresh_uses_persisted_transport_after_store_reopen() {
15661566
broker.join().expect("broker requests")
15671567
},
15681568
vec![
1569-
format!("/runner/jobs/{}", job.id),
1570-
format!("/runner/jobs/{}/events", job.id),
1569+
format!("/jobs/{}", job.id),
1570+
format!("/jobs/{}/events", job.id),
15711571
"/runner/jobs/reconcile".to_string(),
15721572
]
15731573
);

crates/homeboy-lab-runner/src/execution/broker.rs

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -194,7 +194,16 @@ pub(super) fn exec_via_reverse_broker(
194194
}
195195
}
196196
let persisted_run_id = mirror_evidence
197-
.then(|| persist_lab_offload_handoff_run(runner, &cwd, &command, &job, run_id.as_deref()))
197+
.then(|| {
198+
persist_lab_offload_handoff_run(
199+
runner,
200+
&cwd,
201+
&command,
202+
&job,
203+
run_id.as_deref(),
204+
Some(broker_url),
205+
)
206+
})
198207
.flatten();
199208
validate_generic_exec_mirror_run_id(
200209
run_id_owns_generic_exec,

crates/homeboy-lab-runner/src/execution/daemon.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -221,7 +221,9 @@ pub(super) fn exec_via_daemon(
221221
}
222222
}
223223
let persisted_run_id = mirror_evidence
224-
.then(|| persist_lab_offload_handoff_run(runner, &cwd, &command, &job, run_id.as_deref()))
224+
.then(|| {
225+
persist_lab_offload_handoff_run(runner, &cwd, &command, &job, run_id.as_deref(), None)
226+
})
225227
.flatten();
226228
validate_generic_exec_mirror_run_id(
227229
run_id_owns_generic_exec,

crates/homeboy-lab-runner/src/execution/handoff.rs

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ use homeboy_core::error::{Error, Result};
1010
use crate::daemon_http_get::parse_daemon_response_json;
1111

1212
use super::super::broker_http;
13-
use super::super::evidence::mirror_daemon_job_progress;
13+
use super::super::evidence::{mirror_daemon_job_progress, mirror_reverse_broker_job_progress};
1414
use super::super::{load, status, Runner, RunnerTunnelMode};
1515

1616
#[allow(unused_imports)]
@@ -393,8 +393,15 @@ pub(super) fn persist_lab_offload_handoff_run(
393393
command: &[String],
394394
job: &Job,
395395
run_id: Option<&str>,
396+
reverse_broker_url: Option<&str>,
396397
) -> Option<String> {
397-
match mirror_daemon_job_progress(runner, cwd, command, job, &[], run_id) {
398+
let mirrored = match reverse_broker_url {
399+
Some(broker_url) => {
400+
mirror_reverse_broker_job_progress(runner, broker_url, cwd, command, job, run_id)
401+
}
402+
None => mirror_daemon_job_progress(runner, cwd, command, job, &[], run_id),
403+
};
404+
match mirrored {
398405
Ok(run) => Some(run.id),
399406
Err(err) => {
400407
eprintln!(

crates/homeboy-lab-runner/src/execution/tests/handoff.rs

Lines changed: 51 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -381,6 +381,7 @@ fn lab_offload_handoff_persists_run_when_job_is_accepted() {
381381
&["homeboy".to_string(), "trace".to_string()],
382382
&job,
383383
None,
384+
None,
384385
)
385386
.expect("persist handoff run");
386387

@@ -399,6 +400,38 @@ fn lab_offload_handoff_persists_run_when_job_is_accepted() {
399400
});
400401
}
401402

403+
#[test]
404+
fn reverse_handoff_persists_broker_transport_for_post_disconnect_refresh() {
405+
homeboy_core::test_support::with_isolated_home(|_| {
406+
let runner = ssh_runner();
407+
let job = running_job();
408+
let command = vec!["homeboy".to_string(), "trace".to_string()];
409+
let run_id = persist_lab_offload_handoff_run(
410+
&runner,
411+
"/srv/homeboy/project",
412+
&command,
413+
&job,
414+
None,
415+
Some("http://127.0.0.1:4321"),
416+
)
417+
.expect("persist reverse handoff run");
418+
419+
let store = homeboy_core::observation::ObservationStore::open_initialized().expect("store");
420+
let run = store
421+
.get_run(&run_id)
422+
.expect("get run")
423+
.expect("persisted reverse handoff run");
424+
assert_eq!(
425+
run.metadata_json["lab"]["reverse_broker"]["broker_url"],
426+
"http://127.0.0.1:4321"
427+
);
428+
assert_eq!(
429+
run.metadata_json["lab"]["reverse_broker"]["job_id"],
430+
job.id.to_string()
431+
);
432+
});
433+
}
434+
402435
#[test]
403436
fn accepted_job_that_disappears_persists_a_terminal_controller_failure() {
404437
homeboy_core::test_support::with_isolated_home(|_| {
@@ -410,9 +443,15 @@ fn accepted_job_that_disappears_persists_a_terminal_controller_failure() {
410443
"runtime".to_string(),
411444
"refresh".to_string(),
412445
];
413-
let run_id =
414-
persist_lab_offload_handoff_run(&runner, "/srv/homeboy/project", &command, &job, None)
415-
.expect("accepted handoff mirror");
446+
let run_id = persist_lab_offload_handoff_run(
447+
&runner,
448+
"/srv/homeboy/project",
449+
&command,
450+
&job,
451+
None,
452+
None,
453+
)
454+
.expect("accepted handoff mirror");
416455

417456
let err = terminal_runner_poll_failure(
418457
&runner,
@@ -480,9 +519,15 @@ fn transient_daemon_transport_drop_keeps_the_durable_job_recoverable() {
480519
"runtime".to_string(),
481520
"refresh".to_string(),
482521
];
483-
let run_id =
484-
persist_lab_offload_handoff_run(&runner, "/srv/homeboy/project", &command, &job, None)
485-
.expect("accepted handoff mirror");
522+
let run_id = persist_lab_offload_handoff_run(
523+
&runner,
524+
"/srv/homeboy/project",
525+
&command,
526+
&job,
527+
None,
528+
None,
529+
)
530+
.expect("accepted handoff mirror");
486531

487532
// A transport-layer drop: the daemon endpoint became unreachable while
488533
// polling. This is the shape `runner_daemon_health_failure` classifies

tests/reverse_cook_queue_acceptance.rs

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -317,6 +317,18 @@ fn detached_cook_accepts_reverse_capacity_queue_and_worker_completes_once() {
317317
.count(),
318318
1,
319319
);
320+
let terminal_result = broker
321+
.store
322+
.events(completed.id)
323+
.expect("broker events")
324+
.into_iter()
325+
.find(|event| event.kind == JobEventKind::Result)
326+
.and_then(|event| event.data)
327+
.expect("broker terminal result event");
328+
assert!(
329+
terminal_result.get("exit_code").is_some(),
330+
"broker terminal result preserves the typed payload: {terminal_result}"
331+
);
320332

321333
let (_, duplicate_code) =
322334
homeboy::runner::run_reverse_worker(homeboy::runner::ReverseRunnerWorkerOptions {

0 commit comments

Comments
 (0)