Skip to main content

scv_server/
components.rs

1//! Server-owned lifecycle for every long-running integration.
2
3use 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/// Components must observe cancellation and must not detach child tasks.
19/// Return on failure; the supervisor owns retries and bounded shutdown.
20#[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    /// Idempotent start: replacement must first stop and join the old instance.
88    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                // Catch task panics without letting them take down the daemon or skip retries.
108                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            // The runner owns abort/join of its child, so never abort the runner first.
162            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        // A bridge state commit briefly holds this same lock. Retry next refresh
334        // rather than interrupting healthy work for ordinary lock contention.
335        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                // Persist disabled and tool-free first, so failed deletion can
380                // neither resurrect a live account nor hand a later login the grant.
381                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
400/// The longest tool call an owner session in `workspace` may make, from the
401/// configuration its sessions load. Read at each (re)start of the component.
402fn 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
412/// Only the authenticated account owner may receive tools. Credentials without
413/// a known owner ID grant tools to nobody, even when the setting asks for it.
414fn 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}