1pub mod inventory;
31pub mod remote;
32pub mod serve;
33
34use std::path::PathBuf;
35use std::sync::Arc;
36
37use car_fleet::{FleetComposite, InstanceInventory, InstanceRef, InventoryProvider, WorkerProfile};
38use serde::{Deserialize, Serialize};
39use serde_json::Value;
40
41use crate::handler::JsonRpcMessage;
42use crate::session::{ClientSession, ServerState};
43
44pub use remote::RemoteWorktreeAgent;
45
46const DEFAULT_MAX_PARALLEL: u32 = 2;
48
49const DEFAULT_LOCAL_PARALLEL: u32 = 2;
51
52pub(super) const DEFAULT_MAX_SUBTASK_SECS: u64 = 1800;
60
61const INVENTORY_TIMEOUT: std::time::Duration = car_fleet::DEFAULT_INVENTORY_TIMEOUT;
64
65#[derive(Debug, Clone, Serialize, Deserialize)]
73pub struct FleetWorkerConfig {
74 #[serde(default)]
77 pub accepts_work: bool,
78 #[serde(default)]
81 pub repos: Vec<PathBuf>,
82 #[serde(default = "default_max_parallel")]
84 pub max_parallel: u32,
85 #[serde(default = "default_local_parallel")]
88 pub local_parallel: u32,
89 #[serde(default = "default_dispatches_per_hour")]
96 pub dispatches_per_hour: u32,
97 #[serde(default = "default_max_subtask_secs")]
104 pub max_subtask_secs: u64,
105 #[serde(default)]
116 pub fetch_missing_base: bool,
117 #[serde(default = "default_fetch_remote")]
119 pub fetch_remote: String,
120 #[serde(default, skip_serializing_if = "Option::is_none")]
129 pub allowed_tools: Option<Vec<String>>,
130}
131
132fn default_max_parallel() -> u32 {
133 DEFAULT_MAX_PARALLEL
134}
135
136fn default_local_parallel() -> u32 {
137 DEFAULT_LOCAL_PARALLEL
138}
139
140fn default_dispatches_per_hour() -> u32 {
141 car_fleet::DEFAULT_DISPATCHES_PER_WINDOW
142}
143
144fn default_max_subtask_secs() -> u64 {
145 DEFAULT_MAX_SUBTASK_SECS
146}
147
148fn default_fetch_remote() -> String {
149 "origin".to_string()
150}
151
152impl Default for FleetWorkerConfig {
153 fn default() -> Self {
154 Self {
155 accepts_work: false,
156 repos: Vec::new(),
157 max_parallel: DEFAULT_MAX_PARALLEL,
158 local_parallel: DEFAULT_LOCAL_PARALLEL,
159 dispatches_per_hour: car_fleet::DEFAULT_DISPATCHES_PER_WINDOW,
160 max_subtask_secs: DEFAULT_MAX_SUBTASK_SECS,
161 fetch_missing_base: false,
162 fetch_remote: default_fetch_remote(),
163 allowed_tools: None,
164 }
165 }
166}
167
168impl FleetWorkerConfig {
169 fn path() -> Option<PathBuf> {
170 car_home::root().map(|r| r.join("fleet-worker.json"))
171 }
172
173 pub fn load() -> Self {
179 let Some(path) = Self::path() else {
180 return Self::default();
181 };
182 match std::fs::read_to_string(&path) {
183 Ok(text) => serde_json::from_str(&text).unwrap_or_else(|e| {
184 tracing::warn!(path = %path.display(), error = %e, "unreadable fleet worker config; declining work");
185 Self::default()
186 }),
187 Err(_) => Self::default(),
188 }
189 }
190
191 pub fn save(&self) -> Result<(), String> {
192 let path = Self::path().ok_or("cannot resolve the CAR state root")?;
193 if let Some(parent) = path.parent() {
194 std::fs::create_dir_all(parent)
195 .map_err(|e| format!("create {}: {e}", parent.display()))?;
196 }
197 let json = serde_json::to_string_pretty(self).map_err(|e| e.to_string())?;
198 let tmp = path.with_extension("json.tmp");
199 std::fs::write(&tmp, json).map_err(|e| format!("write {}: {e}", tmp.display()))?;
200 std::fs::rename(&tmp, &path).map_err(|e| format!("rename into {}: {e}", path.display()))
201 }
202}
203
204fn worktree_base() -> PathBuf {
208 car_home::root_or_relative()
209 .join("fleet-worker")
210 .join("worktrees")
211}
212
213async fn detected_adapters() -> Vec<car_external_agents::ExternalAgentSpec> {
220 use tokio::sync::Mutex;
221 static CACHE: std::sync::OnceLock<
222 Mutex<
223 Option<(
224 std::time::Instant,
225 Vec<car_external_agents::ExternalAgentSpec>,
226 )>,
227 >,
228 > = std::sync::OnceLock::new();
229 const TTL: std::time::Duration = std::time::Duration::from_secs(60);
230
231 let cache = CACHE.get_or_init(|| Mutex::new(None));
232 let mut guard = cache.lock().await;
233 if let Some((at, specs)) = guard.as_ref() {
234 if at.elapsed() < TTL {
235 return specs.clone();
236 }
237 }
238 let specs = car_external_agents::detect_runnable().await;
239 *guard = Some((std::time::Instant::now(), specs.clone()));
240 specs
241}
242
243pub async fn worker_profile() -> WorkerProfile {
245 let config = FleetWorkerConfig::load();
246 let adapters: Vec<String> = detected_adapters()
247 .await
248 .into_iter()
249 .map(|s| s.id)
250 .collect();
251 let repo_root_commits = config
255 .repos
256 .iter()
257 .filter_map(|p| car_fleet::root_commit(p).ok())
258 .collect();
259 WorkerProfile {
260 accepts_work: config.accepts_work,
261 adapters,
262 max_parallel: config.max_parallel,
263 repo_root_commits,
264 }
265}
266
267async fn remote_providers(
276 state: &ServerState,
277) -> (Vec<Arc<dyn InventoryProvider>>, Vec<InstanceInventory>) {
278 let identity = {
279 state
280 .peer_identity
281 .lock()
282 .unwrap_or_else(|e| e.into_inner())
283 .clone()
284 };
285
286 let mut providers: Vec<Arc<dyn InventoryProvider>> = Vec::new();
287 let mut seen: std::collections::HashSet<String> = std::collections::HashSet::new();
288
289 for peer in crate::peers::snapshot_parslee(state).await {
292 if let car_peers::PeerAddress::A2a { base_url } = &peer.address {
293 if seen.insert(base_url.clone()) {
294 providers.push(Arc::new(inventory::PeerInventoryProvider::new(
295 InstanceRef::remote(peer.name.clone(), "parslee", base_url.clone()),
296 identity.clone(),
297 )));
298 }
299 }
300 }
301
302 if let Ok(registry) = car_a2a::peers::PeerRegistry::user_default() {
304 for entry in registry.list() {
305 if seen.insert(entry.url.clone()) {
306 let name = entry.label.clone().unwrap_or_else(|| entry.slug.clone());
307 providers.push(Arc::new(inventory::PeerInventoryProvider::new(
308 InstanceRef::remote(name, "registry", entry.url.clone()),
309 identity.clone(),
310 )));
311 }
312 }
313 }
314
315 let visible_only = crate::peers::snapshot_lan(state)
316 .into_iter()
317 .filter_map(|peer| {
318 let car_peers::PeerAddress::A2a { base_url } = &peer.address else {
319 return None;
320 };
321 if seen.contains(base_url) {
322 return None;
323 }
324 Some(InstanceInventory::unreachable(
325 InstanceRef::remote(peer.name.clone(), "lan", base_url.clone()),
326 "discovered on the local network but not a trusted peer — anyone can advertise \
327 any name, so promote it with `a2a.peers.add` before CAR will contact it",
328 car_fleet::now_ms(),
329 ))
330 })
331 .collect();
332
333 (providers, visible_only)
334}
335
336pub async fn composite(
338 state: &ServerState,
339 session: Option<&ClientSession>,
340 include_remote: bool,
341 timeout: std::time::Duration,
342) -> FleetComposite {
343 let local = match session {
344 Some(s) => inventory::local_inventory(state, Some(&s.runtime), Some(&s.memgine)).await,
345 None => inventory::local_inventory(state, None, None).await,
346 };
347 let self_name = local.instance.name.clone();
348
349 let mut all = vec![local];
350 if include_remote {
351 let (providers, visible_only) = remote_providers(state).await;
352 all.extend(car_fleet::gather(&providers, timeout).await);
353 all.extend(visible_only);
354 }
355 car_fleet::compose(self_name, all)
356}
357
358pub async fn handle_fleet_inventory(
361 state: &ServerState,
362 session: &ClientSession,
363) -> Result<Value, String> {
364 let inv =
365 inventory::local_inventory(state, Some(&session.runtime), Some(&session.memgine)).await;
366 serde_json::to_value(inv).map_err(|e| e.to_string())
367}
368
369pub async fn handle_fleet_composite(
375 msg: &JsonRpcMessage,
376 state: &ServerState,
377 session: &ClientSession,
378) -> Result<Value, String> {
379 let include_remote = msg
380 .params
381 .get("include_remote")
382 .and_then(|v| v.as_bool())
383 .unwrap_or(true);
384 let timeout = msg
385 .params
386 .get("timeout_ms")
387 .and_then(|v| v.as_u64())
388 .map(std::time::Duration::from_millis)
389 .unwrap_or(INVENTORY_TIMEOUT);
390 let composite = composite(state, Some(session), include_remote, timeout).await;
391 serde_json::to_value(composite).map_err(|e| e.to_string())
392}
393
394pub async fn handle_fleet_worker_get() -> Result<Value, String> {
396 let config = FleetWorkerConfig::load();
397 let profile = worker_profile().await;
398 Ok(serde_json::json!({
399 "config": config,
400 "profile": profile,
401 }))
402}
403
404pub async fn handle_fleet_worker_set(
419 msg: &JsonRpcMessage,
420 session: &ClientSession,
421) -> Result<Value, String> {
422 let bound_agent = session.agent_id.lock().await.clone();
423 if let Some(agent) = bound_agent {
424 if !session.is_host.load(std::sync::atomic::Ordering::Acquire) {
425 return Err(format!(
426 "`fleet.worker.set` is operator-only: `{agent}` cannot enroll this machine to \
427 run peers' coding subtasks against local checkouts"
428 ));
429 }
430 }
431
432 let mut config = FleetWorkerConfig::load();
433 if let Some(v) = msg.params.get("accepts_work").and_then(|v| v.as_bool()) {
434 config.accepts_work = v;
435 }
436 if let Some(list) = msg.params.get("repos").and_then(|v| v.as_array()) {
437 let mut repos = Vec::new();
438 for entry in list {
439 let path = PathBuf::from(entry.as_str().ok_or("`repos` entries must be strings")?);
440 car_fleet::root_commit(&path)
445 .map_err(|e| format!("`{}` is not a git repository: {e}", path.display()))?;
446 repos.push(path);
447 }
448 config.repos = repos;
449 }
450 if let Some(n) = msg.params.get("max_parallel").and_then(|v| v.as_u64()) {
451 config.max_parallel = n.min(64) as u32;
452 }
453 if let Some(n) = msg.params.get("local_parallel").and_then(|v| v.as_u64()) {
454 config.local_parallel = n.min(64) as u32;
455 }
456 if let Some(n) = msg
457 .params
458 .get("dispatches_per_hour")
459 .and_then(|v| v.as_u64())
460 {
461 config.dispatches_per_hour = n.min(10_000) as u32;
462 }
463 if let Some(n) = msg.params.get("max_subtask_secs").and_then(|v| v.as_u64()) {
464 config.max_subtask_secs = n.clamp(60, 24 * 3600);
468 }
469 if let Some(v) = msg
470 .params
471 .get("fetch_missing_base")
472 .and_then(|v| v.as_bool())
473 {
474 config.fetch_missing_base = v;
475 }
476 if let Some(remote) = msg.params.get("fetch_remote").and_then(|v| v.as_str()) {
477 if remote.is_empty()
481 || remote.starts_with('-')
482 || !remote
483 .chars()
484 .all(|c| c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.' | '/'))
485 {
486 return Err(format!("`{remote}` is not a usable git remote name"));
487 }
488 config.fetch_remote = remote.to_string();
489 }
490 if let Some(list) = msg.params.get("allowed_tools").and_then(|v| v.as_array()) {
491 config.allowed_tools = Some(
492 list.iter()
493 .filter_map(|v| v.as_str().map(String::from))
494 .collect(),
495 );
496 }
497 config.save()?;
498 Ok(serde_json::json!({
499 "config": config,
500 "profile": worker_profile().await,
501 }))
502}
503
504pub struct DaemonFleetResponder {
509 state: std::sync::Weak<ServerState>,
510 runtime: Arc<car_engine::Runtime>,
514}
515
516impl DaemonFleetResponder {
517 pub fn new(state: std::sync::Weak<ServerState>, runtime: Arc<car_engine::Runtime>) -> Self {
518 Self { state, runtime }
519 }
520}
521
522#[async_trait::async_trait]
523impl car_a2a::FleetResponder for DaemonFleetResponder {
524 async fn inventory(&self) -> Result<Value, String> {
525 let state = self
526 .state
527 .upgrade()
528 .ok_or_else(|| "daemon is shutting down".to_string())?;
529 let inv = inventory::local_inventory(&state, Some(&self.runtime), None).await;
532 serde_json::to_value(inv).map_err(|e| e.to_string())
533 }
534
535 async fn run_subtask(&self, dispatch: Value, caller: Option<&str>) -> Result<Value, String> {
536 let Some(caller) = caller else {
541 return Err(
542 "refusing an unattributed fleet dispatch: this surface accepts work only from a CAR peer whose signature identifies it"
543 .to_string(),
544 );
545 };
546 let dispatch: car_fleet::SubtaskDispatch =
547 serde_json::from_value(dispatch).map_err(|e| format!("invalid dispatch: {e}"))?;
548 let outcome = serve::run_dispatch(dispatch, caller).await?;
549 serde_json::to_value(outcome).map_err(|e| e.to_string())
550 }
551}
552
553#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
561#[serde(rename_all = "snake_case")]
562pub enum PoolExclusion {
563 NotEnrolled,
565 RepositoryNotServed,
567 NotRequested,
569 Unreachable,
571 NoAddress,
573}
574
575impl PoolExclusion {
576 pub fn as_str(self) -> &'static str {
577 match self {
578 PoolExclusion::NotEnrolled => "not_enrolled",
579 PoolExclusion::RepositoryNotServed => "repository_not_served",
580 PoolExclusion::NotRequested => "not_requested",
581 PoolExclusion::Unreachable => "unreachable",
582 PoolExclusion::NoAddress => "no_address",
583 }
584 }
585}
586
587#[derive(Debug, Clone, Default, Serialize, Deserialize)]
589pub struct PoolPlan {
590 pub remote_workers: Vec<String>,
593 pub excluded: Vec<(String, String)>,
595}
596
597impl PoolPlan {
598 pub fn local_only(&self) -> bool {
600 self.remote_workers.is_empty()
601 }
602
603 pub fn degraded_reason(&self) -> Option<String> {
605 if !self.local_only() {
606 return None;
607 }
608 if self.excluded.is_empty() {
609 return Some(
610 "no other CAR instance is reachable, so this ran on one machine".to_string(),
611 );
612 }
613 let mut counts: std::collections::BTreeMap<&str, usize> = std::collections::BTreeMap::new();
614 for (_, reason) in &self.excluded {
615 *counts.entry(reason.as_str()).or_default() += 1;
616 }
617 let detail = counts
618 .into_iter()
619 .map(|(reason, n)| format!("{n} {reason}"))
620 .collect::<Vec<_>>()
621 .join(", ");
622 Some(format!(
623 "no peer could take a subtask, so this ran on one machine ({detail})"
624 ))
625 }
626}
627
628pub async fn build_pool(
639 state: &ServerState,
640 repo_root: &std::path::Path,
641 run_id: &str,
642 adapter: &str,
643 only: Option<&[String]>,
644) -> Result<(car_multi::FleetPool, PoolPlan), String> {
645 let fingerprint = car_fleet::read_fingerprint(repo_root).map_err(|e| {
646 format!(
647 "cannot identify the repository at {}: {e}",
648 repo_root.display()
649 )
650 })?;
651 let config = FleetWorkerConfig::load();
652
653 let local: Arc<dyn car_multi::WorktreeAgent> = Arc::new(
654 car_external_agents::ForemanExternalAgent::new(adapter.to_string()),
655 );
656 let mut workers = vec![car_multi::FleetWorker::local(
657 car_a2a::lan::host_label(),
658 local,
659 config.local_parallel.max(1) as usize,
660 )];
661
662 let composite = composite(state, None, true, INVENTORY_TIMEOUT).await;
663 let identity = {
664 state
665 .peer_identity
666 .lock()
667 .unwrap_or_else(|e| e.into_inner())
668 .clone()
669 };
670
671 let mut plan = PoolPlan::default();
674 let eligible: std::collections::HashSet<&str> = composite
675 .workers_for(&fingerprint.root_commit)
676 .into_iter()
677 .map(|c| c.instance.as_str())
678 .collect();
679 for inv in &composite.instances {
680 if inv.instance.kind == car_fleet::InstanceKind::Local {
681 continue;
682 }
683 let name = inv.instance.name.clone();
684 let reason = if !inv.reachable() {
685 PoolExclusion::Unreachable
686 } else if !eligible.contains(name.as_str()) {
687 match &inv.worker {
690 Some(w) if w.accepts_work => PoolExclusion::RepositoryNotServed,
691 _ => PoolExclusion::NotEnrolled,
692 }
693 } else if only.is_some_and(|only| !only.iter().any(|n| n == &name)) {
694 PoolExclusion::NotRequested
695 } else if inv.instance.base_url.is_none() {
696 PoolExclusion::NoAddress
697 } else {
698 continue;
699 };
700 plan.excluded.push((name, reason.as_str().to_string()));
701 }
702
703 for candidate in composite.workers_for(&fingerprint.root_commit) {
704 if candidate.kind == car_fleet::InstanceKind::Local {
705 continue;
706 }
707 if let Some(only) = only {
708 if !only.iter().any(|n| n == &candidate.instance) {
709 continue;
710 }
711 }
712 let Some(base_url) = composite
714 .instances
715 .iter()
716 .find(|i| i.instance.name == candidate.instance)
717 .and_then(|i| i.instance.base_url.clone())
718 else {
719 continue;
720 };
721 let wanted = candidate
725 .adapters
726 .iter()
727 .any(|a| a == adapter)
728 .then(|| adapter.to_string());
729 let agent = RemoteWorktreeAgent::new(
730 candidate.instance.clone(),
731 base_url,
732 identity.clone(),
733 fingerprint.clone(),
734 run_id,
735 )
736 .with_adapter(wanted);
737 plan.remote_workers.push(candidate.instance.clone());
738 workers.push(car_multi::FleetWorker::remote(
739 candidate.instance.clone(),
740 Arc::new(agent),
741 candidate.max_parallel.max(1) as usize,
742 ));
743 }
744
745 if let Some(reason) = plan.degraded_reason() {
746 tracing::warn!(%reason, "distributed foreman run has no peer workers");
747 }
748 Ok((car_multi::FleetPool::new(workers), plan))
749}
750
751pub fn placements_value(placements: &[car_multi::Placement]) -> Value {
761 serde_json::to_value(placements).expect("placement ledger is plain data")
765}
766
767pub fn placements_json(pool: &car_multi::FleetPool) -> Value {
769 placements_value(&pool.placements())
770}
771
772#[cfg(test)]
773mod tests {
774 use super::*;
775
776 #[test]
777 fn a_distributed_run_with_no_peers_says_why_rather_than_going_quiet() {
778 let plan = PoolPlan {
779 remote_workers: Vec::new(),
780 excluded: vec![
781 ("studio".into(), "not_enrolled".into()),
782 ("laptop".into(), "not_enrolled".into()),
783 ("ci-box".into(), "repository_not_served".into()),
784 ],
785 };
786 assert!(plan.local_only());
787 let reason = plan.degraded_reason().expect("degraded");
788 assert!(reason.contains("2 not_enrolled"), "{reason}");
789 assert!(reason.contains("1 repository_not_served"), "{reason}");
790 }
791
792 #[test]
793 fn a_pool_with_peers_reports_no_degradation() {
794 let plan = PoolPlan {
795 remote_workers: vec!["studio".into()],
796 excluded: vec![("laptop".into(), "not_enrolled".into())],
797 };
798 assert!(!plan.local_only());
799 assert!(plan.degraded_reason().is_none());
800 }
801
802 #[test]
803 fn no_reachable_peers_at_all_is_its_own_message() {
804 let plan = PoolPlan::default();
805 let reason = plan.degraded_reason().expect("degraded");
806 assert!(
807 reason.contains("no other CAR instance is reachable"),
808 "{reason}"
809 );
810 }
811
812 #[test]
813 fn the_shipped_default_declines_work() {
814 let c = FleetWorkerConfig::default();
817 assert!(!c.accepts_work);
818 assert!(c.repos.is_empty());
819 }
820
821 #[test]
822 fn a_laptop_does_not_fetch_unless_told_to() {
823 let c = FleetWorkerConfig::default();
826 assert!(!c.fetch_missing_base);
827 assert_eq!(c.fetch_remote, "origin");
828 }
829
830 #[test]
831 fn a_config_written_before_runner_mode_existed_still_declines() {
832 let parsed: FleetWorkerConfig =
833 serde_json::from_str("{\"accepts_work\": true}").expect("parses");
834 assert!(
835 !parsed.fetch_missing_base,
836 "absent must not read as enabled"
837 );
838 assert_eq!(parsed.fetch_remote, "origin");
839 }
840
841 #[test]
842 fn the_sender_does_not_choose_this_machines_limits() {
843 let parsed: FleetWorkerConfig =
847 serde_json::from_str("{\"accepts_work\": true, \"repos\": []}").expect("parses");
848 assert_eq!(
849 parsed.dispatches_per_hour,
850 car_fleet::DEFAULT_DISPATCHES_PER_WINDOW
851 );
852 assert_eq!(parsed.max_subtask_secs, DEFAULT_MAX_SUBTASK_SECS);
853 assert!(parsed.allowed_tools.is_none());
854 }
855
856 #[test]
857 fn an_unparseable_config_declines_rather_than_half_accepting() {
858 let parsed: FleetWorkerConfig =
859 serde_json::from_str("{\"accepts_work\": true}").expect("partial config parses");
860 assert!(parsed.accepts_work);
861 assert_eq!(parsed.max_parallel, DEFAULT_MAX_PARALLEL, "limits default");
862 assert!(
863 parsed.repos.is_empty(),
864 "and with no repos it can still serve nothing"
865 );
866 }
867
868 #[test]
869 fn worker_worktrees_live_outside_every_served_repository() {
870 let base = worktree_base();
871 assert!(
872 base.ends_with("fleet-worker/worktrees") || base.ends_with("fleet-worker\\worktrees")
873 );
874 }
875}