1use anyhow::{Result, bail};
4use async_trait::async_trait;
5use scv_clawbot::state::{self, Account, AccountSettings};
6use scv_protocol::{ComponentHealth, ComponentState, DaemonCommand, DaemonStatus, RemoteTools};
7use std::{
8 collections::BTreeMap,
9 path::PathBuf,
10 sync::{Arc, Mutex},
11 time::{Duration, Instant, SystemTime, UNIX_EPOCH},
12};
13use tokio::task::JoinHandle;
14use tokio_util::sync::CancellationToken;
15
16const STOP_GRACE: Duration = Duration::from_secs(5);
17
18#[async_trait]
21pub trait Component: Send + Sync + 'static {
22 async fn run(&self, cancellation: CancellationToken, health: HealthReporter) -> Result<()>;
23}
24
25#[derive(Clone)]
26pub struct HealthReporter(Arc<Mutex<ComponentHealth>>);
27
28impl HealthReporter {
29 pub fn contact(&self, connected: bool) {
30 let mut health = self.0.lock().unwrap();
31 if matches!(
32 health.state,
33 ComponentState::Stopping | ComponentState::Stopped
34 ) {
35 return;
36 }
37 health.state = if connected {
38 ComponentState::Connected
39 } else {
40 ComponentState::Disconnected
41 };
42 health.error = (!connected).then(|| "Component contact failed".into());
43 if connected {
44 health.last_success_unix_seconds = Some(
45 SystemTime::now()
46 .duration_since(UNIX_EPOCH)
47 .unwrap_or_default()
48 .as_secs(),
49 );
50 }
51 }
52
53 fn transition(&self, state: ComponentState, error: Option<&str>) {
54 let mut health = self.0.lock().unwrap();
55 health.state = state;
56 health.error = error.map(str::to_owned);
57 }
58
59 fn snapshot(&self) -> ComponentHealth {
60 self.0.lock().unwrap().clone()
61 }
62}
63
64pub struct Supervisor {
65 tasks: BTreeMap<String, RunningComponent>,
66 grace: Duration,
67 initial_backoff: Duration,
68}
69
70struct RunningComponent {
71 cancellation: CancellationToken,
72 task: JoinHandle<()>,
73 health: HealthReporter,
74}
75
76impl Default for Supervisor {
77 fn default() -> Self {
78 Self {
79 tasks: BTreeMap::new(),
80 grace: STOP_GRACE,
81 initial_backoff: Duration::from_secs(1),
82 }
83 }
84}
85
86impl Supervisor {
87 pub fn start(&mut self, component: Arc<dyn Component>, health: ComponentHealth) {
89 if self.tasks.contains_key(&health.id) {
90 return;
91 }
92 let id = health.id.clone();
93 let health = HealthReporter(Arc::new(Mutex::new(health)));
94 let cancellation = CancellationToken::new();
95 let cancel = cancellation.clone();
96 let report = health.clone();
97 let initial_backoff = self.initial_backoff;
98 let grace = self.grace;
99 let task = tokio::spawn(async move {
100 let mut delay = initial_backoff;
101 loop {
102 if cancel.is_cancelled() {
103 break;
104 }
105 report.transition(ComponentState::Starting, None);
106 let started = Instant::now();
107 let instance = component.clone();
109 let child_cancel = cancel.clone();
110 let child_report = report.clone();
111 let mut child =
112 tokio::spawn(async move { instance.run(child_cancel, child_report).await });
113 tokio::select! {
114 biased;
115 _ = cancel.cancelled() => {
116 report.transition(ComponentState::Stopping, None);
117 if tokio::time::timeout(grace, &mut child).await.is_err() {
118 child.abort();
119 let _ = child.await;
120 }
121 break;
122 }
123 _ = &mut child => {}
124 }
125 report.transition(
126 ComponentState::Backoff,
127 Some("Component stopped unexpectedly; retrying"),
128 );
129 report.0.lock().unwrap().restarts += 1;
130 if started.elapsed() >= Duration::from_secs(60) {
131 delay = initial_backoff;
132 }
133 tokio::select! {
134 _ = cancel.cancelled() => break,
135 _ = tokio::time::sleep(delay) => {}
136 }
137 delay = (delay * 2).min(Duration::from_secs(60));
138 }
139 report.transition(ComponentState::Stopped, None);
140 });
141 self.tasks.insert(
142 id,
143 RunningComponent {
144 cancellation,
145 task,
146 health,
147 },
148 );
149 }
150
151 pub fn health(&self) -> Vec<ComponentHealth> {
152 self.tasks
153 .values()
154 .map(|task| task.health.snapshot())
155 .collect()
156 }
157
158 pub async fn stop(&mut self, id: &str) {
159 if let Some(running) = self.tasks.get_mut(id) {
160 running.cancellation.cancel();
161 let _ = (&mut running.task).await;
163 }
164 self.tasks.remove(id);
165 }
166
167 pub async fn shutdown(&mut self) {
168 for task in self.tasks.values() {
169 task.cancellation.cancel();
170 }
171 for id in self.tasks.keys().cloned().collect::<Vec<_>>() {
172 self.stop(&id).await;
173 }
174 }
175}
176
177struct ClawBot {
178 account: String,
179 credentials: Account,
180 workspace: PathBuf,
181 socket: PathBuf,
182 tool_owner: Option<String>,
183}
184
185#[async_trait]
186impl Component for ClawBot {
187 async fn run(&self, cancellation: CancellationToken, health: HealthReporter) -> Result<()> {
188 let tool_owner = self.tool_owner.clone().map(|user_id| {
189 let turn_timeout = scv_clawbot::owner_turn_timeout(max_tool_timeout(&self.workspace));
190 tracing::info!(
191 "ClawBot {} owner turns may run up to {} seconds",
192 self.account,
193 turn_timeout.as_secs()
194 );
195 scv_clawbot::ToolOwner {
196 user_id,
197 turn_timeout,
198 }
199 });
200 scv_clawbot::run_supervised(
201 &self.credentials.token,
202 &self.credentials.base_url,
203 &self.account,
204 &self.workspace,
205 &self.socket,
206 tool_owner.as_ref(),
207 cancellation,
208 Arc::new(move |connected| health.contact(connected)),
209 )
210 .await
211 }
212}
213
214pub(crate) struct Components {
215 supervisor: Supervisor,
216 desired: BTreeMap<String, (Account, AccountSettings)>,
217 inactive: BTreeMap<String, ComponentHealth>,
218 socket: PathBuf,
219 workspace: PathBuf,
220}
221
222impl Components {
223 pub fn new(socket: PathBuf, workspace: PathBuf) -> Self {
224 Self {
225 supervisor: Supervisor::default(),
226 desired: BTreeMap::new(),
227 inactive: BTreeMap::new(),
228 socket,
229 workspace,
230 }
231 }
232
233 pub fn status(&self) -> DaemonStatus {
234 let mut components = self.supervisor.health();
235 components.extend(self.inactive.values().cloned());
236 components.sort_by(|a, b| a.id.cmp(&b.id));
237 DaemonStatus {
238 version: env!("CARGO_PKG_VERSION").into(),
239 pid: std::process::id(),
240 components,
241 }
242 }
243
244 pub async fn reconcile(&mut self) -> Result<()> {
245 let names = match state::account_names() {
246 Ok(names) => names,
247 Err(_) => {
248 self.supervisor.shutdown().await;
249 self.desired.clear();
250 self.inactive.clear();
251 let mut health = initial_health("discovery", None, false);
252 health.id = "clawbot:discovery-error".into();
253 health.state = ComponentState::Failed;
254 health.error = Some(
255 "Account discovery failed; components stopped until configuration is readable"
256 .into(),
257 );
258 self.inactive.insert("discovery-error".into(), health);
259 bail!("Account discovery failed");
260 }
261 };
262 for name in self
263 .desired
264 .keys()
265 .chain(self.inactive.keys())
266 .cloned()
267 .collect::<Vec<_>>()
268 {
269 if !names.contains(&name) {
270 self.supervisor.stop(&format!("clawbot:{name}")).await;
271 self.desired.remove(&name);
272 self.inactive.remove(&name);
273 }
274 }
275 for name in names {
276 let loaded = (|| -> Result<_> {
277 let (account, settings) = state::account_snapshot(&name)?;
278 Ok((
279 account.ok_or_else(|| anyhow::anyhow!("missing account"))?,
280 settings,
281 ))
282 })();
283 let (credentials, settings) = match loaded {
284 Ok(value) => value,
285 Err(error) => {
286 self.account_error(name, error).await;
287 continue;
288 }
289 };
290 if self.desired.get(&name) == Some(&(credentials.clone(), settings.clone())) {
291 continue;
292 }
293 self.supervisor.stop(&format!("clawbot:{name}")).await;
294 self.inactive.remove(&name);
295 let mut health = initial_health(&name, Some(&credentials), settings.enabled);
296 let tool_owner = tool_owner(&credentials, &settings);
297 if tool_owner.is_some() {
298 health.remote_tools = RemoteTools::Owner;
299 }
300 if settings.enabled {
301 let workspace = settings
302 .workspace
303 .clone()
304 .unwrap_or_else(|| self.workspace.clone());
305 if !workspace.is_absolute() || !workspace.is_dir() {
306 health.state = ComponentState::Failed;
307 health.error =
308 Some("Component workspace must be an existing absolute directory".into());
309 self.inactive.insert(name.clone(), health);
310 self.desired.remove(&name);
311 continue;
312 }
313 self.supervisor.start(
314 Arc::new(ClawBot {
315 account: name.clone(),
316 credentials: credentials.clone(),
317 workspace,
318 socket: self.socket.clone(),
319 tool_owner,
320 }),
321 health,
322 );
323 } else {
324 health.state = ComponentState::Disabled;
325 self.inactive.insert(name.clone(), health);
326 }
327 self.desired.insert(name, (credentials, settings));
328 }
329 Ok(())
330 }
331
332 async fn account_error(&mut self, name: String, error: anyhow::Error) {
333 if error
336 .downcast_ref::<std::io::Error>()
337 .is_some_and(|error| error.kind() == std::io::ErrorKind::WouldBlock)
338 {
339 return;
340 }
341 self.supervisor.stop(&format!("clawbot:{name}")).await;
342 self.desired.remove(&name);
343 let mut health = initial_health(&name, None, true);
344 health.state = ComponentState::Failed;
345 health.error = Some("Invalid or inaccessible account/settings".into());
346 self.inactive.insert(name, health);
347 }
348
349 pub async fn control(&mut self, command: DaemonCommand) -> Result<DaemonStatus> {
350 match command {
351 DaemonCommand::Status => return Ok(self.status()),
352 DaemonCommand::Reload => {}
353 DaemonCommand::ClawbotSet {
354 account,
355 enabled,
356 workspace,
357 remote_tools,
358 } => {
359 state::validate_name(&account)?;
360 if state::account(&account)?.is_none() {
361 bail!("Account is not logged in");
362 }
363 let mut settings = state::settings(&account)?;
364 settings.enabled = enabled;
365 if let Some(path) = workspace {
366 let path = PathBuf::from(path);
367 if !path.is_absolute() || !path.is_dir() {
368 bail!("Invalid component workspace");
369 }
370 settings.workspace = Some(std::fs::canonicalize(path)?);
371 }
372 if let Some(mode) = remote_tools {
373 settings.remote_tools = mode;
374 }
375 state::save_settings(&account, &settings)?;
376 }
377 DaemonCommand::ClawbotLogout { account } => {
378 state::validate_name(&account)?;
379 let mut settings = state::settings(&account)?;
382 settings.enabled = false;
383 settings.remote_tools = RemoteTools::None;
384 state::save_settings(&account, &settings)?;
385 self.supervisor.stop(&format!("clawbot:{account}")).await;
386 self.desired.remove(&account);
387 self.inactive.remove(&account);
388 state::remove(&account)?;
389 }
390 }
391 self.reconcile().await?;
392 Ok(self.status())
393 }
394
395 pub async fn shutdown(&mut self) {
396 self.supervisor.shutdown().await;
397 }
398}
399
400fn max_tool_timeout(workspace: &std::path::Path) -> std::time::Duration {
403 let seconds = crate::Config::load(workspace, crate::ConfigOverrides::default())
404 .map(|config| config.tools.max_timeout_seconds)
405 .unwrap_or_else(|error| {
406 tracing::warn!("ClawBot uses the default tool timeout ceiling: {error:#}");
407 crate::config::ToolConfig::default().max_timeout_seconds
408 });
409 std::time::Duration::from_secs(seconds)
410}
411
412fn tool_owner(credentials: &Account, settings: &AccountSettings) -> Option<String> {
415 (settings.remote_tools == RemoteTools::Owner)
416 .then(|| credentials.user_id.clone())
417 .flatten()
418 .filter(|owner| !owner.is_empty())
419}
420
421fn initial_health(account: &str, credentials: Option<&Account>, enabled: bool) -> ComponentHealth {
422 ComponentHealth {
423 id: format!("clawbot:{account}"),
424 account: account.into(),
425 bot_id: credentials.and_then(|a| a.bot_id.clone()),
426 user_id: credentials.and_then(|a| a.user_id.clone()),
427 enabled,
428 state: ComponentState::Starting,
429 last_success_unix_seconds: None,
430 error: None,
431 restarts: 0,
432 remote_tools: RemoteTools::None,
433 }
434}
435
436#[cfg(test)]
437mod tests {
438 use super::*;
439 use std::sync::atomic::{AtomicUsize, Ordering};
440
441 struct Fake {
442 starts: Arc<AtomicUsize>,
443 stops: Arc<AtomicUsize>,
444 fail_first: bool,
445 }
446 #[async_trait]
447 impl Component for Fake {
448 async fn run(&self, cancellation: CancellationToken, health: HealthReporter) -> Result<()> {
449 let attempt = self.starts.fetch_add(1, Ordering::SeqCst);
450 if self.fail_first && attempt == 0 {
451 bail!("secret error must never enter status");
452 }
453 health.contact(true);
454 cancellation.cancelled().await;
455 self.stops.fetch_add(1, Ordering::SeqCst);
456 Ok(())
457 }
458 }
459
460 #[tokio::test]
461 async fn starts_once_recovers_reports_contact_and_joins_before_restoration() {
462 let starts = Arc::new(AtomicUsize::new(0));
463 let stops = Arc::new(AtomicUsize::new(0));
464 let fake = Arc::new(Fake {
465 starts: starts.clone(),
466 stops: stops.clone(),
467 fail_first: true,
468 });
469 let mut supervisor = Supervisor {
470 initial_backoff: Duration::from_millis(10),
471 ..Supervisor::default()
472 };
473 supervisor.start(fake.clone(), initial_health("test", None, true));
474 supervisor.start(fake.clone(), initial_health("test", None, true));
475 tokio::time::timeout(Duration::from_secs(2), async {
476 loop {
477 if supervisor.health()[0].state == ComponentState::Connected {
478 break;
479 }
480 tokio::time::sleep(Duration::from_millis(1)).await;
481 }
482 })
483 .await
484 .unwrap();
485 let health = &supervisor.health()[0];
486 assert_eq!(starts.load(Ordering::SeqCst), 2);
487 assert_eq!(health.restarts, 1);
488 assert!(health.last_success_unix_seconds.is_some());
489 assert!(health.error.is_none());
490 supervisor.shutdown().await;
491 assert_eq!(stops.load(Ordering::SeqCst), 1);
492 supervisor.start(fake, initial_health("test", None, true));
493 tokio::time::sleep(Duration::from_millis(20)).await;
494 supervisor.shutdown().await;
495 assert_eq!(starts.load(Ordering::SeqCst), 3);
496 assert_eq!(stops.load(Ordering::SeqCst), 2);
497 }
498
499 #[test]
500 fn credentials_are_not_connection_evidence() {
501 let health = initial_health("saved", None, true);
502 assert_eq!(health.state, ComponentState::Starting);
503 assert_eq!(health.last_success_unix_seconds, None);
504 }
505
506 #[tokio::test]
507 async fn busy_account_snapshot_preserves_live_work_but_invalid_settings_stop_it() {
508 let starts = Arc::new(AtomicUsize::new(0));
509 let stops = Arc::new(AtomicUsize::new(0));
510 let mut components = Components::new(PathBuf::from("/unused.sock"), PathBuf::from("/"));
511 components.supervisor.start(
512 Arc::new(Fake {
513 starts: starts.clone(),
514 stops: stops.clone(),
515 fail_first: false,
516 }),
517 initial_health("test", None, true),
518 );
519 tokio::time::timeout(Duration::from_secs(1), async {
520 while starts.load(Ordering::SeqCst) == 0 {
521 tokio::task::yield_now().await;
522 }
523 })
524 .await
525 .unwrap();
526 components
527 .account_error(
528 "test".into(),
529 std::io::Error::from(std::io::ErrorKind::WouldBlock).into(),
530 )
531 .await;
532 assert_eq!(
533 components.status().components[0].state,
534 ComponentState::Connected
535 );
536 assert_eq!(starts.load(Ordering::SeqCst), 1);
537 assert_eq!(stops.load(Ordering::SeqCst), 0);
538 components
539 .account_error("test".into(), anyhow::anyhow!("invalid settings"))
540 .await;
541 assert_eq!(stops.load(Ordering::SeqCst), 1);
542 assert_eq!(
543 components.status().components[0].state,
544 ComponentState::Failed
545 );
546 }
547
548 struct Stubborn;
549 #[async_trait]
550 impl Component for Stubborn {
551 async fn run(&self, _: CancellationToken, _: HealthReporter) -> Result<()> {
552 std::future::pending().await
553 }
554 }
555
556 #[tokio::test]
557 async fn bounded_stop_aborts_uncooperative_component_and_cancels_backoff() {
558 let mut supervisor = Supervisor {
559 grace: Duration::from_millis(20),
560 ..Supervisor::default()
561 };
562 supervisor.start(Arc::new(Stubborn), initial_health("stubborn", None, true));
563 tokio::task::yield_now().await;
564 tokio::time::timeout(Duration::from_secs(1), supervisor.shutdown())
565 .await
566 .unwrap();
567 assert!(supervisor.health().is_empty());
568 let fake = Arc::new(Fake {
569 starts: Arc::new(AtomicUsize::new(0)),
570 stops: Arc::new(AtomicUsize::new(0)),
571 fail_first: true,
572 });
573 supervisor.start(fake, initial_health("backoff", None, true));
574 tokio::time::sleep(Duration::from_millis(10)).await;
575 assert_eq!(supervisor.health()[0].state, ComponentState::Backoff);
576 assert_eq!(
577 supervisor.health()[0].error.as_deref(),
578 Some("Component stopped unexpectedly; retrying")
579 );
580 tokio::time::timeout(Duration::from_millis(100), supervisor.shutdown())
581 .await
582 .unwrap();
583 }
584
585 #[test]
586 fn remote_tools_require_owner_mode_and_known_owner() {
587 let account = |user_id: Option<&str>| Account {
588 token: "token".into(),
589 base_url: "https://example.invalid".into(),
590 bot_id: Some("bot".into()),
591 user_id: user_id.map(Into::into),
592 };
593 let owner = AccountSettings {
594 remote_tools: RemoteTools::Owner,
595 ..Default::default()
596 };
597 assert_eq!(
598 tool_owner(&account(Some("owner@im.wechat")), &owner).as_deref(),
599 Some("owner@im.wechat")
600 );
601 assert_eq!(tool_owner(&account(None), &owner), None);
602 assert_eq!(tool_owner(&account(Some("")), &owner), None);
603 assert_eq!(
604 tool_owner(
605 &account(Some("owner@im.wechat")),
606 &AccountSettings::default()
607 ),
608 None
609 );
610 }
611}