Skip to content

Commit 8cae0e4

Browse files
authored
fix(lab): poll reverse staging through broker (#10589)
1 parent 60d3b49 commit 8cae0e4

2 files changed

Lines changed: 38 additions & 1 deletion

File tree

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

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2363,11 +2363,28 @@ impl ProductionLabStagingOperations {
23632363
events: observed.events,
23642364
});
23652365
}
2366+
let broker_url = Self::status_for_recipe_transport(request)?
2367+
.session
2368+
.and_then(|session| session.broker_url)
2369+
.ok_or_else(|| {
2370+
Error::validation_invalid_argument(
2371+
"runner",
2372+
"durable reverse Lab staging requires its accepted broker endpoint",
2373+
Some(request.recipe.runner_id.clone()),
2374+
None,
2375+
)
2376+
})?;
23662377
loop {
23672378
if cancellation.is_some_and(LabStagingCancellationToken::is_cancelled) {
23682379
return Err(Self::cancellation_error());
23692380
}
2370-
match crate::runner_job_log_snapshot(&request.recipe.runner_id, runner_job_id) {
2381+
match crate::connection::reverse_broker_job_snapshot_at(
2382+
&broker_url,
2383+
&request.recipe.runner_id,
2384+
runner_job_id,
2385+
)
2386+
.map(|(job, events)| homeboy_core::api_jobs::RunnerJobLogSnapshot { job, events })
2387+
{
23712388
Ok(snapshot) if snapshot.job.status.is_terminal() => return Ok(snapshot),
23722389
Ok(_) => std::thread::sleep(Duration::from_millis(200)),
23732390
Err(mut error) => {

tests/reverse_cook_queue_acceptance.rs

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -329,6 +329,26 @@ fn detached_cook_accepts_reverse_capacity_queue_and_worker_completes_once() {
329329
terminal_result.get("exit_code").is_some(),
330330
"broker terminal result preserves the typed payload: {terminal_result}"
331331
);
332+
let broker_events: serde_json::Value =
333+
reqwest::blocking::get(format!("{}/jobs/{}/events", broker.url(), completed.id))
334+
.expect("fetch broker events over HTTP")
335+
.json()
336+
.expect("parse broker events response");
337+
let broker_terminal_result = broker_events
338+
.pointer("/data/body/events")
339+
.and_then(serde_json::Value::as_array)
340+
.and_then(|events| {
341+
events.iter().rev().find_map(|event| {
342+
(event["kind"] == serde_json::json!("result")).then(|| event["data"].clone())
343+
})
344+
})
345+
.expect("broker HTTP response retains terminal result");
346+
serde_json::from_value::<homeboy::core::api_jobs::RemoteRunnerJobResult>(
347+
broker_terminal_result.clone(),
348+
)
349+
.unwrap_or_else(|error| {
350+
panic!("broker HTTP terminal result must retain its typed contract: {error}\nresult={broker_terminal_result}")
351+
});
332352

333353
let (_, duplicate_code) =
334354
homeboy::runner::run_reverse_worker(homeboy::runner::ReverseRunnerWorkerOptions {

0 commit comments

Comments
 (0)