diff --git a/CHANGELOG.md b/CHANGELOG.md index 2044eb001..bb252f668 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,7 +5,11 @@ All notable changes to Agent Relay will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] +## [Unreleased - Patch] + +### Fixed + +- Live broker workers missing from the Relaycast reconnect inventory are now restored from their existing agent identity, so a node-control reconnect no longer makes a still-running terminal permanently unreachable. ## [11.6.9] - 2026-08-16 diff --git a/crates/broker/src/runtime/event_loop.rs b/crates/broker/src/runtime/event_loop.rs index ba983f0fb..b6d9428a6 100644 --- a/crates/broker/src/runtime/event_loop.rs +++ b/crates/broker/src/runtime/event_loop.rs @@ -231,6 +231,10 @@ pub(crate) struct BrokerRuntime { pub(super) fleet_delivery_book: FleetDeliveryBook, pub(super) fleet_max_agents: u32, pub(super) fleet_inventory: HashMap, + /// Per-worker retry deadlines for failed Relaycast identity lookups while + /// rebuilding the reconnect inventory. + pub(super) fleet_inventory_reconcile_retry_after: + HashMap, pub(super) sdk_out_tx: mpsc::Sender>, pub(super) worker_event_rx: mpsc::Receiver, pub(super) worker_events_open: bool, diff --git a/crates/broker/src/runtime/fleet.rs b/crates/broker/src/runtime/fleet.rs index 1e84bf869..e61d5fe1d 100644 --- a/crates/broker/src/runtime/fleet.rs +++ b/crates/broker/src/runtime/fleet.rs @@ -11,6 +11,7 @@ use crate::{ TerminalControlCommand, TerminalControlEvent, TerminalFromCloud, TerminalMode, TerminalToCloud, }, + worker::LiveFleetInventoryCandidate, }; use base64::{engine::general_purpose::STANDARD as BASE64, Engine as _}; @@ -21,6 +22,18 @@ const TERMINAL_INPUT_MAX_BASE64_BYTES: usize = TERMINAL_INPUT_MAX_BYTES * 4 / 3 const TERMINAL_SNAPSHOT_TIMEOUT: Duration = Duration::from_secs(10); const TERMINAL_INPUT_ACK_TIMEOUT: Duration = Duration::from_secs(5); const TERMINAL_INPUT_MAX_IN_FLIGHT_PER_SESSION: usize = 16; +// Reconciliation runs from the broker's single event loop. Keep a transient +// Relaycast outage or a large inventory gap from monopolizing a maintenance +// tick; deferred workers are revisited on later ticks. +const FLEET_INVENTORY_RECONCILE_BATCH_SIZE: usize = 2; +const FLEET_INVENTORY_RECONCILE_LOOKUP_TIMEOUT: Duration = Duration::from_secs(2); +const FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF: Duration = Duration::from_secs(60); + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(super) struct FleetInventoryRetry { + generation: Uuid, + retry_after: Instant, +} // Relaycast currently limits a node to 32 terminal sessions. Keep that many // slots free from high-volume frames so every affected session can still get a // terminal.closed notification when the output lane applies backpressure. @@ -1693,6 +1706,150 @@ pub(super) async fn record_fleet_inventory_agent( publish_fleet_inventory_snapshot(fleet_control_tx, fleet_inventory).await; } +/// Restore reconnect-inventory entries for broker-owned workers that are still +/// live locally. A successful process launch and a successful inventory write +/// happen on separate asynchronous paths, so the inventory is a projection of +/// the live worker registry rather than an independently authoritative list. +/// +/// This deliberately resolves an existing Relaycast identity by name instead +/// of calling `register_agent_token`: reconciliation must never rotate a live +/// worker's credential merely to rebuild the reconnect snapshot. +fn schedule_fleet_inventory_retry( + retry_after: &mut HashMap, + name: WorkerName, + generation: Uuid, + now: Instant, +) { + retry_after.insert( + name, + FleetInventoryRetry { + generation, + retry_after: now + FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF, + }, + ); +} + +pub(super) async fn reconcile_fleet_inventory_with_live_workers( + fleet_control_tx: &mpsc::Sender, + relaycast_http: &RelaycastHttpClient, + fleet_delivery_book: &mut FleetDeliveryBook, + fleet_inventory: &mut HashMap, + retry_after: &mut HashMap, + live_workers: Vec, + now: Instant, +) -> usize { + let live_worker_generations: HashMap<_, _> = live_workers + .iter() + .map(|worker| (worker.name.clone(), worker.generation)) + .collect(); + // A worker that exits or has since been restored needs no retained retry + // state. A restarted same-name worker has a different generation and must not + // inherit the old process's retry deadline. + retry_after.retain(|name, retry| { + live_worker_generations.get(name) == Some(&retry.generation) + && !fleet_inventory.contains_key(name) + }); + let missing_workers: Vec<_> = live_workers + .into_iter() + .filter(|worker| { + !fleet_inventory.contains_key(&worker.name) + && retry_after + .get(&worker.name) + .is_none_or(|retry| retry.retry_after <= now) + }) + .take(FLEET_INVENTORY_RECONCILE_BATCH_SIZE) + .collect(); + if missing_workers.is_empty() { + return 0; + } + + let Some(relay) = relaycast_http.relay_client() else { + // Relaycast is optional for local broker mode. Do not turn every + // maintenance tick into a warning when a client is not configured. + return 0; + }; + + let mut repaired = 0; + for LiveFleetInventoryCandidate { + name, + session_ref, + generation, + } in missing_workers + { + let agent = match timeout( + FLEET_INVENTORY_RECONCILE_LOOKUP_TIMEOUT, + relay.get_agent(name.as_str()), + ) + .await + { + Ok(Ok(agent)) => agent, + Ok(Err(error)) => { + tracing::warn!( + worker = %name, + error = %error, + retry_after_secs = FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF.as_secs(), + "could not resolve live worker for fleet inventory reconciliation; will retry later" + ); + schedule_fleet_inventory_retry(retry_after, name, generation, now); + continue; + } + Err(_) => { + tracing::warn!( + worker = %name, + timeout_secs = FLEET_INVENTORY_RECONCILE_LOOKUP_TIMEOUT.as_secs(), + retry_after_secs = FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF.as_secs(), + "timed out resolving live worker for fleet inventory reconciliation; will retry later" + ); + schedule_fleet_inventory_retry(retry_after, name, generation, now); + continue; + } + }; + if agent.name != name.as_str() { + tracing::warn!( + worker = %name, + resolved_name = %agent.name, + retry_after_secs = FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF.as_secs(), + "refusing to reconcile fleet inventory with a mismatched Relaycast identity; will retry later" + ); + schedule_fleet_inventory_retry(retry_after, name, generation, now); + continue; + } + + if let Some(bound_agent_id) = fleet_delivery_book.active_agent_id(name.as_str()) { + if bound_agent_id != agent.id { + tracing::warn!( + worker = %name, + bound_agent_id, + resolved_agent_id = %agent.id, + retry_after_secs = FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF.as_secs(), + "refusing to replace a live worker's authoritative fleet identity; will retry later" + ); + schedule_fleet_inventory_retry(retry_after, name, generation, now); + continue; + } + } else { + fleet_delivery_book.bind_authoritative_identity(agent.name.clone(), agent.id.clone()); + } + + retry_after.remove(&name); + fleet_inventory.insert( + name, + InventoryAgent { + agent_id: agent.id, + name: agent.name, + invocation_id: None, + session_ref, + }, + ); + repaired += 1; + } + + if repaired > 0 { + publish_fleet_inventory_snapshot(fleet_control_tx, fleet_inventory).await; + } + repaired +} + /// Resolve an opaque agent token to the authoritative identity required by /// `inventory.sync` and delivery bookkeeping. /// @@ -2115,6 +2272,19 @@ pub(super) fn fleet_initial_session_ref(spec: &AgentSpec) -> Option { mod tests { use super::*; use crate::protocol::PtyHarnessConfig; + use httpmock::{Method::GET, Method::POST, MockServer}; + + fn live_fleet_worker( + name: &str, + session_ref: Option<&str>, + generation: u128, + ) -> LiveFleetInventoryCandidate { + LiveFleetInventoryCandidate { + name: WorkerName::from(name), + session_ref: session_ref.map(ToOwned::to_owned), + generation: Uuid::from_u128(generation), + } + } #[cfg(unix)] #[tokio::test] @@ -3072,6 +3242,419 @@ mod tests { } } + #[tokio::test] + async fn reconciliation_restores_a_live_worker_missing_from_inventory_without_reregistering() { + let server = MockServer::start(); + let lookup = server.mock(|when, then| { + when.method(GET) + .path("/v1/agents/live-worker") + .header("authorization", "Bearer rk_live_test"); + then.status(200).json_body(serde_json::json!({ + "ok": true, + "data": { + "id": "agent-live-id", + "name": "live-worker", + "type": "agent", + "status": "offline", + "persona": null, + "metadata": {} + } + })); + }); + // A reconciliation must only look the already-running worker up. Any + // registration request can rotate its token and would recreate #1545. + let registration = server.mock(|when, then| { + when.method(POST).path("/v1/agents"); + then.status(500).json_body(serde_json::json!({ + "ok": false, + "error": { "code": "must_not_register", "message": "must not register" } + })); + }); + let relaycast_http = + RelaycastHttpClient::new(Some(server.base_url()), "rk_live_test", "broker", "claude"); + let (tx, mut rx) = mpsc::channel(2); + let mut inventory = HashMap::new(); + let mut delivery_book = FleetDeliveryBook::default(); + let mut retry_after = HashMap::new(); + + let repaired = reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + vec![live_fleet_worker("live-worker", Some("session-live"), 101)], + Instant::now(), + ) + .await; + + assert_eq!(repaired, 1, "the live orphan must be restored"); + assert_eq!( + inventory.get(&WorkerName::from("live-worker")), + Some(&InventoryAgent { + agent_id: "agent-live-id".to_string(), + name: "live-worker".to_string(), + invocation_id: None, + session_ref: Some("session-live".to_string()), + }) + ); + assert_eq!( + delivery_book.active_agent_id("live-worker"), + Some("agent-live-id") + ); + match rx.recv().await { + Some(FleetControlCommand::UpdateInventory(agents)) => { + assert_eq!(agents.len(), 1); + assert_eq!(agents[0].name, "live-worker"); + assert_eq!(agents[0].agent_id, "agent-live-id"); + } + other => panic!("expected repaired inventory snapshot, got {other:?}"), + } + lookup.assert_hits(1); + registration.assert_hits(0); + } + + #[tokio::test] + async fn reconciliation_rejects_a_name_mismatch_instead_of_binding_the_wrong_identity() { + let server = MockServer::start(); + let lookup = server.mock(|when, then| { + when.method(GET) + .path("/v1/agents/live-worker") + .header("authorization", "Bearer rk_live_test"); + then.status(200).json_body(serde_json::json!({ + "ok": true, + "data": { + "id": "agent-other-id", + "name": "other-worker", + "type": "agent", + "status": "offline", + "persona": null, + "metadata": {} + } + })); + }); + let relaycast_http = + RelaycastHttpClient::new(Some(server.base_url()), "rk_live_test", "broker", "claude"); + let (tx, mut rx) = mpsc::channel(2); + let mut inventory = HashMap::new(); + let mut delivery_book = FleetDeliveryBook::default(); + let mut retry_after = HashMap::new(); + + let repaired = reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + vec![live_fleet_worker("live-worker", None, 102)], + Instant::now(), + ) + .await; + + assert_eq!(repaired, 0, "a mismatched name must not be reconciled"); + assert!(inventory.is_empty()); + assert_eq!(delivery_book.active_agent_id("live-worker"), None); + assert!( + rx.try_recv().is_err(), + "a rejected identity must not publish" + ); + lookup.assert_hits(1); + } + + #[tokio::test] + async fn reconciliation_does_not_replace_an_existing_authoritative_identity() { + let server = MockServer::start(); + let lookup = server.mock(|when, then| { + when.method(GET) + .path("/v1/agents/live-worker") + .header("authorization", "Bearer rk_live_test"); + then.status(200).json_body(serde_json::json!({ + "ok": true, + "data": { + "id": "agent-reused-name-id", + "name": "live-worker", + "type": "agent", + "status": "offline", + "persona": null, + "metadata": {} + } + })); + }); + let relaycast_http = + RelaycastHttpClient::new(Some(server.base_url()), "rk_live_test", "broker", "claude"); + let (tx, mut rx) = mpsc::channel(1); + let mut inventory = HashMap::new(); + let mut delivery_book = FleetDeliveryBook::default(); + delivery_book.bind_authoritative_identity("live-worker", "agent-live-id"); + let mut retry_after = HashMap::new(); + let now = Instant::now(); + + let repaired = reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + vec![live_fleet_worker("live-worker", None, 103)], + now, + ) + .await; + + assert_eq!(repaired, 0, "a reused name must not replace live identity"); + assert!(inventory.is_empty()); + assert_eq!( + delivery_book.active_agent_id("live-worker"), + Some("agent-live-id") + ); + assert_eq!( + retry_after.get(&WorkerName::from("live-worker")), + Some(&FleetInventoryRetry { + generation: Uuid::from_u128(103), + retry_after: now + FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF, + }) + ); + assert!( + rx.try_recv().is_err(), + "a rejected replacement must not publish" + ); + lookup.assert_hits(1); + } + + #[tokio::test] + async fn reconciliation_defers_workers_past_the_per_tick_batch() { + let server = MockServer::start(); + let worker_names: Vec<_> = (0..=FLEET_INVENTORY_RECONCILE_BATCH_SIZE) + .map(|index| format!("live-worker-{index}")) + .collect(); + let lookups: Vec<_> = worker_names + .iter() + .map(|worker_name| { + let path = format!("/v1/agents/{worker_name}"); + let resolved_name = worker_name.clone(); + let agent_id = format!("agent-{worker_name}"); + server.mock(move |when, then| { + when.method(GET) + .path(path) + .header("authorization", "Bearer rk_live_test"); + then.status(200).json_body(serde_json::json!({ + "ok": true, + "data": { + "id": agent_id, + "name": resolved_name, + "type": "agent", + "status": "offline", + "persona": null, + "metadata": {} + } + })); + }) + }) + .collect(); + let relaycast_http = + RelaycastHttpClient::new(Some(server.base_url()), "rk_live_test", "broker", "claude"); + let (tx, mut rx) = mpsc::channel(4); + let mut inventory = HashMap::new(); + let mut delivery_book = FleetDeliveryBook::default(); + let mut retry_after = HashMap::new(); + let live_workers = || { + worker_names + .iter() + .enumerate() + .map(|(index, name)| live_fleet_worker(name, None, index as u128 + 1)) + .collect() + }; + let now = Instant::now(); + + let first_repaired = reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + live_workers(), + now, + ) + .await; + + assert_eq!(first_repaired, FLEET_INVENTORY_RECONCILE_BATCH_SIZE); + assert_eq!(inventory.len(), FLEET_INVENTORY_RECONCILE_BATCH_SIZE); + for lookup in lookups.iter().take(FLEET_INVENTORY_RECONCILE_BATCH_SIZE) { + lookup.assert_hits(1); + } + lookups[FLEET_INVENTORY_RECONCILE_BATCH_SIZE].assert_hits(0); + let _ = rx.recv().await.expect("expected first batch snapshot"); + + let second_repaired = reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + live_workers(), + now, + ) + .await; + + assert_eq!(second_repaired, 1); + assert_eq!(inventory.len(), FLEET_INVENTORY_RECONCILE_BATCH_SIZE + 1); + lookups[FLEET_INVENTORY_RECONCILE_BATCH_SIZE].assert_hits(1); + let _ = rx.recv().await.expect("expected second batch snapshot"); + } + + #[tokio::test] + async fn reconciliation_backs_off_failed_identity_lookups() { + let server = MockServer::start(); + let lookup = server.mock(|when, then| { + when.method(GET) + .path("/v1/agents/missing-worker") + .header("authorization", "Bearer rk_live_test"); + then.status(404).json_body(serde_json::json!({ + "ok": false, + "error": { "code": "agent_not_found", "message": "not found" } + })); + }); + let relaycast_http = + RelaycastHttpClient::new(Some(server.base_url()), "rk_live_test", "broker", "claude"); + let (tx, _rx) = mpsc::channel(1); + let mut inventory = HashMap::new(); + let mut delivery_book = FleetDeliveryBook::default(); + let mut retry_after = HashMap::new(); + let live_workers = |generation| vec![live_fleet_worker("missing-worker", None, generation)]; + let now = Instant::now(); + + assert_eq!( + reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + live_workers(104), + now, + ) + .await, + 0 + ); + assert_eq!(lookup.hits(), 1); + assert_eq!( + retry_after.get(&WorkerName::from("missing-worker")), + Some(&FleetInventoryRetry { + generation: Uuid::from_u128(104), + retry_after: now + FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF, + }) + ); + + assert_eq!( + reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + live_workers(104), + now + Duration::from_secs(1), + ) + .await, + 0 + ); + assert_eq!(lookup.hits(), 1, "the backoff must suppress the next tick"); + + assert_eq!( + reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + live_workers(105), + now + Duration::from_secs(2), + ) + .await, + 0 + ); + assert_eq!( + lookup.hits(), + 2, + "a restarted same-name worker bypasses the old retry deadline" + ); + assert_eq!( + retry_after.get(&WorkerName::from("missing-worker")), + Some(&FleetInventoryRetry { + generation: Uuid::from_u128(105), + retry_after: now + + Duration::from_secs(2) + + FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF, + }) + ); + + assert_eq!( + reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + live_workers(105), + now + Duration::from_secs(2) + FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF, + ) + .await, + 0 + ); + assert_eq!(lookup.hits(), 3, "the worker is retried after backoff"); + } + + #[tokio::test] + async fn reconciliation_times_out_a_slow_lookup_and_schedules_backoff() { + let server = MockServer::start(); + let lookup = server.mock(|when, then| { + when.method(GET) + .path("/v1/agents/slow-worker") + .header("authorization", "Bearer rk_live_test"); + then.status(200) + .delay(Duration::from_secs(4)) + .json_body(serde_json::json!({ + "ok": true, + "data": { + "id": "agent-slow-id", + "name": "slow-worker", + "type": "agent", + "status": "offline", + "persona": null, + "metadata": {} + } + })); + }); + let relaycast_http = + RelaycastHttpClient::new(Some(server.base_url()), "rk_live_test", "broker", "claude"); + let (tx, _rx) = mpsc::channel(1); + let mut inventory = HashMap::new(); + let mut delivery_book = FleetDeliveryBook::default(); + let mut retry_after = HashMap::new(); + let now = Instant::now(); + + let repaired = reconcile_fleet_inventory_with_live_workers( + &tx, + &relaycast_http, + &mut delivery_book, + &mut inventory, + &mut retry_after, + vec![live_fleet_worker("slow-worker", None, 106)], + now, + ) + .await; + + assert_eq!(repaired, 0); + assert!(inventory.is_empty()); + assert_eq!(lookup.hits(), 1); + assert_eq!( + retry_after.get(&WorkerName::from("slow-worker")), + Some(&FleetInventoryRetry { + generation: Uuid::from_u128(106), + retry_after: now + FLEET_INVENTORY_RECONCILE_FAILURE_BACKOFF, + }) + ); + } + #[tokio::test] async fn fleet_inventory_snapshot_waits_for_backpressure_instead_of_dropping() { let (tx, mut rx) = mpsc::channel::(1); diff --git a/crates/broker/src/runtime/init.rs b/crates/broker/src/runtime/init.rs index b3aabc292..613534ace 100644 --- a/crates/broker/src/runtime/init.rs +++ b/crates/broker/src/runtime/init.rs @@ -697,6 +697,7 @@ pub(crate) async fn run_init(cmd: InitCommand, telemetry: TelemetryClient) -> Re // field); 0 means unlimited, matching the register manifest. fleet_max_agents: node_max_agents().unwrap_or(0), fleet_inventory: HashMap::new(), + fleet_inventory_reconcile_retry_after: HashMap::new(), sdk_out_tx, worker_event_rx, worker_events_open: true, diff --git a/crates/broker/src/runtime/maintenance.rs b/crates/broker/src/runtime/maintenance.rs index b514dcab4..a7bd25c5d 100644 --- a/crates/broker/src/runtime/maintenance.rs +++ b/crates/broker/src/runtime/maintenance.rs @@ -14,6 +14,7 @@ impl BrokerRuntime { let workers = &mut self.workers; let fleet_control_tx = &self.fleet_control_tx; let fleet_inventory = &mut self.fleet_inventory; + let fleet_inventory_reconcile_retry_after = &mut self.fleet_inventory_reconcile_retry_after; let fleet_delivery_book = &mut self.fleet_delivery_book; let fleet_max_agents = self.fleet_max_agents; // The broker provider's capacity handlers are live whenever it is @@ -696,6 +697,22 @@ impl BrokerRuntime { } } + let live_fleet_workers = workers.live_fleet_inventory_candidates(); + if super::fleet::reconcile_fleet_inventory_with_live_workers( + fleet_control_tx, + relaycast_http, + fleet_delivery_book, + fleet_inventory, + fleet_inventory_reconcile_retry_after, + live_fleet_workers, + now, + ) + .await + > 0 + { + fleet_load_changed = true; + } + // Publish the fleet load snapshot once, after both reaping and restart // handling, so the broadcast count reflects the final post-restart live // worker set rather than a same-tick post-reap intermediate. diff --git a/crates/broker/src/worker.rs b/crates/broker/src/worker.rs index 52e292ca8..3ace26833 100644 --- a/crates/broker/src/worker.rs +++ b/crates/broker/src/worker.rs @@ -238,6 +238,15 @@ pub(crate) enum WorkerEvent { }, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub(crate) struct LiveFleetInventoryCandidate { + pub(crate) name: WorkerName, + pub(crate) session_ref: Option, + /// The process generation distinguishes a restarted same-name worker from the + /// process whose failed Relaycast lookup is currently being backed off. + pub(crate) generation: Uuid, +} + pub(crate) struct WorkerRegistry { pub(crate) workers: HashMap, event_tx: mpsc::Sender, @@ -502,6 +511,34 @@ impl WorkerRegistry { } } + /// Return the live broker-owned workers that must remain present in the + /// Relaycast reconnect inventory. The parent marker is set by the two + /// production spawn surfaces (Dashboard and Relaycast); workers without it + /// are local-only and must not cause a fleet identity lookup. + pub(crate) fn live_fleet_inventory_candidates(&self) -> Vec { + self.workers + .iter() + .filter_map(|(name, handle)| { + if handle.parent.is_none() || !self.is_worker_live(name) { + return None; + } + let session_ref = handle.spec.session_id.clone().or_else(|| { + handle + .spec + .harness_config + .as_ref() + .and_then(ResolvedHarnessConfig::session_id) + .map(ToOwned::to_owned) + }); + Some(LiveFleetInventoryCandidate { + name: name.clone(), + session_ref, + generation: handle.generation, + }) + }) + .collect() + } + pub(crate) fn worker_pid(&self, name: &str) -> Option { self.workers.get(name).and_then(|h| h.child.id()) }