Skip to main content

synd_runtime/daemon/
control.rs

1use std::time::{Duration, Instant};
2
3#[cfg(unix)]
4use rustix::process::Signal;
5use synd_protocol::daemon::DaemonStatusResponse;
6use tokio::time::sleep;
7
8#[cfg(unix)]
9use crate::daemon::{
10    DaemonClaim, DaemonClaimLockAcquirer, SignalTarget, remove_stale_claim,
11    wait_until_claim_released,
12};
13use crate::{
14    DaemonState, DaemonStatus, Error, PlacementSummary, Result, Runtime, ShutdownResult,
15    connection::{RuntimeEndpointConnectionStatus, RuntimeEndpointConnector},
16    placement::PlacementSpec,
17    startup::{StartupLock, StartupLockAcquirer, StartupLockAcquisition},
18};
19
20const SHUTDOWN_WAIT_INTERVAL: Duration = Duration::from_millis(50);
21
22#[derive(Debug, Clone, Copy)]
23pub struct Control<'a> {
24    runtime: &'a Runtime,
25}
26
27impl<'a> Control<'a> {
28    pub(crate) fn new(runtime: &'a Runtime) -> Self {
29        Self { runtime }
30    }
31
32    pub fn runtime(&self) -> &Runtime {
33        self.runtime
34    }
35
36    pub async fn inspect(&self) -> Result<DaemonStatus> {
37        let context = self.resolve_context().await?;
38        let decision = DaemonStatusDecision::from(context);
39
40        match decision {
41            DaemonStatusDecision::RequestStatus { placement } => {
42                let summary = PlacementSummary::from_placement(&placement);
43                match DaemonControlClient::new(self.runtime, &placement)?
44                    .status()
45                    .await
46                {
47                    Ok(status) => Ok(DaemonStatus::running(summary, status.sessions().clone())),
48                    Err(error) if daemon_status_endpoint_missing(&error) => {
49                        Ok(DaemonStatus::new(DaemonState::Running, summary))
50                    }
51                    Err(error) => Err(error.into()),
52                }
53            }
54            DaemonStatusDecision::AlreadyStopped { placement } => Ok(DaemonStatus::new(
55                DaemonState::NotRunning,
56                PlacementSummary::from_placement(&placement),
57            )),
58            DaemonStatusDecision::FailUnavailableEndpoint { placement } => {
59                Err(Error::EndpointUnavailable {
60                    context: "daemon endpoint",
61                    endpoint: placement.endpoint().path().to_path_buf(),
62                })
63            }
64            #[cfg(not(unix))]
65            DaemonStatusDecision::FailUnsupportedTransport { .. } => {
66                Err(Error::UnsupportedTransport {
67                    context: "daemon control transport",
68                })
69            }
70        }
71    }
72
73    pub async fn shutdown(&self) -> Result<ShutdownResult> {
74        let context = self.resolve_context().await?;
75        let decision = DaemonShutdownDecision::from(context);
76
77        match decision {
78            DaemonShutdownDecision::RequestShutdown { placement } => {
79                let summary = PlacementSummary::from_placement(&placement);
80                DaemonControlClient::new(self.runtime, &placement)?
81                    .shutdown()
82                    .await?;
83                DaemonShutdownWaiter::new(
84                    placement,
85                    self.runtime.config().session().acquire_timeout(),
86                )
87                .wait()
88                .await?;
89
90                Ok(ShutdownResult::new(DaemonStatus::new(
91                    DaemonState::NotRunning,
92                    summary,
93                )))
94            }
95            DaemonShutdownDecision::AlreadyStopped { placement } => {
96                Ok(ShutdownResult::new(DaemonStatus::new(
97                    DaemonState::NotRunning,
98                    PlacementSummary::from_placement(&placement),
99                )))
100            }
101            DaemonShutdownDecision::FailUnavailableEndpoint { placement } => {
102                Err(Error::EndpointUnavailable {
103                    context: "daemon endpoint",
104                    endpoint: placement.endpoint().path().to_path_buf(),
105                })
106            }
107            #[cfg(not(unix))]
108            DaemonShutdownDecision::FailUnsupportedTransport { .. } => {
109                Err(Error::UnsupportedTransport {
110                    context: "daemon control transport",
111                })
112            }
113        }
114    }
115
116    pub async fn force_shutdown(&self) -> Result<ShutdownResult> {
117        #[cfg(unix)]
118        {
119            return self.force_shutdown_unix().await;
120        }
121
122        #[cfg(not(unix))]
123        {
124            Err(Error::UnsupportedTransport {
125                context: "daemon control transport",
126            })
127        }
128    }
129
130    #[cfg(unix)]
131    async fn force_shutdown_unix(&self) -> Result<ShutdownResult> {
132        let placement = self.runtime.placement().clone();
133        let summary = PlacementSummary::from_placement(&placement);
134        let _startup_lock = acquire_startup_lock(&placement)?;
135
136        let Some(claim) = DaemonClaim::read(placement.daemon_claim_path())? else {
137            return self.shutdown_without_claim(placement, summary).await;
138        };
139
140        let claim_lock = DaemonClaimLockAcquirer::new(placement.daemon_claim_lock_path());
141        if !claim_lock.is_held()? {
142            return self
143                .shutdown_with_stale_claim(placement, summary, "daemon claim lock is not held")
144                .await;
145        }
146
147        let target = SignalTarget::validate(&placement, &claim)?;
148        if !target.send(Signal::TERM)? {
149            remove_stale_claim(placement.daemon_claim_path())?;
150            return Ok(ShutdownResult::new(DaemonStatus::new(
151                DaemonState::NotRunning,
152                summary,
153            )));
154        }
155
156        let timeout = self.runtime.config().session().acquire_timeout();
157        if !wait_until_claim_released(placement.daemon_claim_lock_path(), timeout).await? {
158            let target = SignalTarget::validate(&placement, &claim)?;
159            let _ = target.send(Signal::KILL)?;
160            if !wait_until_claim_released(placement.daemon_claim_lock_path(), timeout).await? {
161                return Err(Error::ForceShutdownTimeout { pid: claim.pid() });
162            }
163        }
164
165        remove_stale_claim(placement.daemon_claim_path())?;
166
167        Ok(ShutdownResult::new(DaemonStatus::new(
168            DaemonState::NotRunning,
169            summary,
170        )))
171    }
172
173    #[cfg(unix)]
174    async fn shutdown_without_claim(
175        &self,
176        placement: PlacementSpec,
177        summary: PlacementSummary,
178    ) -> Result<ShutdownResult> {
179        match RuntimeEndpointConnector::new(placement.endpoint())
180            .try_connect()
181            .await
182        {
183            RuntimeEndpointConnectionStatus::Missing | RuntimeEndpointConnectionStatus::Stale => {
184                Ok(ShutdownResult::new(DaemonStatus::new(
185                    DaemonState::NotRunning,
186                    summary,
187                )))
188            }
189            RuntimeEndpointConnectionStatus::Connected
190            | RuntimeEndpointConnectionStatus::Unavailable => Err(Error::ForceShutdownRefused {
191                reason: format!(
192                    "daemon claim is missing at {}; cannot prove endpoint owner",
193                    placement.daemon_claim_path().path().display()
194                ),
195            }),
196            #[cfg(not(unix))]
197            RuntimeEndpointConnectionStatus::UnsupportedTransport => {
198                Err(Error::UnsupportedTransport {
199                    context: "daemon control transport",
200                })
201            }
202        }
203    }
204
205    #[cfg(unix)]
206    async fn shutdown_with_stale_claim(
207        &self,
208        placement: PlacementSpec,
209        summary: PlacementSummary,
210        reason: &'static str,
211    ) -> Result<ShutdownResult> {
212        match RuntimeEndpointConnector::new(placement.endpoint())
213            .try_connect()
214            .await
215        {
216            RuntimeEndpointConnectionStatus::Missing | RuntimeEndpointConnectionStatus::Stale => {
217                remove_stale_claim(placement.daemon_claim_path())?;
218                Ok(ShutdownResult::new(DaemonStatus::new(
219                    DaemonState::NotRunning,
220                    summary,
221                )))
222            }
223            RuntimeEndpointConnectionStatus::Connected
224            | RuntimeEndpointConnectionStatus::Unavailable => Err(Error::ForceShutdownRefused {
225                reason: format!(
226                    "{reason}; refusing to use stale daemon claim while endpoint {} is not stopped",
227                    placement.endpoint().path().display()
228                ),
229            }),
230            #[cfg(not(unix))]
231            RuntimeEndpointConnectionStatus::UnsupportedTransport => {
232                Err(Error::UnsupportedTransport {
233                    context: "daemon control transport",
234                })
235            }
236        }
237    }
238
239    async fn resolve_context(&self) -> Result<DaemonControlContext> {
240        let placement = self.runtime.placement().clone();
241        let endpoint_connection = RuntimeEndpointConnector::new(placement.endpoint())
242            .try_connect()
243            .await;
244
245        Ok(DaemonControlContext {
246            placement,
247            endpoint_connection,
248        })
249    }
250}
251
252#[cfg(unix)]
253fn acquire_startup_lock(placement: &PlacementSpec) -> Result<StartupLock> {
254    match StartupLockAcquirer::new(placement.startup_lock_path()).try_acquire()? {
255        StartupLockAcquisition::Acquired(lock) => Ok(lock),
256        StartupLockAcquisition::AlreadyHeld => Err(Error::ForceShutdownRefused {
257            reason: format!(
258                "startup lock is already held at {}",
259                placement.startup_lock_path().path().display()
260            ),
261        }),
262        #[cfg(not(unix))]
263        StartupLockAcquisition::UnsupportedTransport => Err(Error::UnsupportedTransport {
264            context: "daemon startup lock",
265        }),
266    }
267}
268
269/// Facts collected before selecting a daemon control action.
270#[derive(Debug, Clone)]
271struct DaemonControlContext {
272    placement: PlacementSpec,
273    endpoint_connection: RuntimeEndpointConnectionStatus,
274}
275
276/// Branch selected for daemon status inspection.
277#[derive(Debug, Clone, PartialEq, Eq)]
278enum DaemonStatusDecision {
279    RequestStatus {
280        placement: PlacementSpec,
281    },
282    AlreadyStopped {
283        placement: PlacementSpec,
284    },
285    FailUnavailableEndpoint {
286        placement: PlacementSpec,
287    },
288    #[cfg(not(unix))]
289    FailUnsupportedTransport {
290        placement: PlacementSpec,
291    },
292}
293
294impl From<DaemonControlContext> for DaemonStatusDecision {
295    fn from(context: DaemonControlContext) -> Self {
296        match context.endpoint_connection {
297            RuntimeEndpointConnectionStatus::Connected => Self::RequestStatus {
298                placement: context.placement,
299            },
300            RuntimeEndpointConnectionStatus::Missing | RuntimeEndpointConnectionStatus::Stale => {
301                Self::AlreadyStopped {
302                    placement: context.placement,
303                }
304            }
305            RuntimeEndpointConnectionStatus::Unavailable => Self::FailUnavailableEndpoint {
306                placement: context.placement,
307            },
308            #[cfg(not(unix))]
309            RuntimeEndpointConnectionStatus::UnsupportedTransport => {
310                Self::FailUnsupportedTransport {
311                    placement: context.placement,
312                }
313            }
314        }
315    }
316}
317
318/// Branch selected for daemon shutdown.
319#[derive(Debug, Clone, PartialEq, Eq)]
320enum DaemonShutdownDecision {
321    RequestShutdown {
322        placement: PlacementSpec,
323    },
324    AlreadyStopped {
325        placement: PlacementSpec,
326    },
327    FailUnavailableEndpoint {
328        placement: PlacementSpec,
329    },
330    #[cfg(not(unix))]
331    FailUnsupportedTransport {
332        placement: PlacementSpec,
333    },
334}
335
336impl From<DaemonControlContext> for DaemonShutdownDecision {
337    fn from(context: DaemonControlContext) -> Self {
338        match context.endpoint_connection {
339            RuntimeEndpointConnectionStatus::Connected => Self::RequestShutdown {
340                placement: context.placement,
341            },
342            RuntimeEndpointConnectionStatus::Missing | RuntimeEndpointConnectionStatus::Stale => {
343                Self::AlreadyStopped {
344                    placement: context.placement,
345                }
346            }
347            RuntimeEndpointConnectionStatus::Unavailable => Self::FailUnavailableEndpoint {
348                placement: context.placement,
349            },
350            #[cfg(not(unix))]
351            RuntimeEndpointConnectionStatus::UnsupportedTransport => {
352                Self::FailUnsupportedTransport {
353                    placement: context.placement,
354                }
355            }
356        }
357    }
358}
359
360/// Sends daemon control requests over the resolved runtime endpoint.
361struct DaemonControlClient {
362    client: synd_client::Client,
363}
364
365impl DaemonControlClient {
366    #[cfg(unix)]
367    fn new(runtime: &Runtime, placement: &PlacementSpec) -> Result<Self> {
368        let client = synd_client::Client::new_unix(
369            placement.endpoint().path(),
370            synd_client::ClientOptions::new(
371                runtime.config().client().request_timeout(),
372                runtime.config().client().user_agent(),
373            ),
374        )?;
375
376        Ok(Self { client })
377    }
378
379    #[cfg(not(unix))]
380    fn new(runtime: &Runtime, placement: &PlacementSpec) -> Result<Self> {
381        let _ = (runtime, placement);
382
383        Err(Error::UnsupportedTransport {
384            context: "daemon control transport",
385        })
386    }
387
388    async fn shutdown(&self) -> Result<()> {
389        self.client.shutdown_daemon().await?;
390        Ok(())
391    }
392
393    async fn status(&self) -> std::result::Result<DaemonStatusResponse, synd_client::SyndApiError> {
394        self.client.daemon_status().await
395    }
396}
397
398fn daemon_status_endpoint_missing(error: &synd_client::SyndApiError) -> bool {
399    matches!(
400        error,
401        synd_client::SyndApiError::HttpStatus {
402            status,
403            url: Some(url),
404        } if status.as_u16() == 404 && url.path() == synd_protocol::daemon::STATUS_PATH
405    )
406}
407
408/// Waits until the daemon endpoint stops accepting connections.
409struct DaemonShutdownWaiter {
410    placement: PlacementSpec,
411    timeout: Duration,
412}
413
414impl DaemonShutdownWaiter {
415    fn new(placement: PlacementSpec, timeout: Duration) -> Self {
416        Self { placement, timeout }
417    }
418
419    async fn wait(&self) -> Result<()> {
420        let deadline = Instant::now() + self.timeout;
421
422        loop {
423            let endpoint_connection = RuntimeEndpointConnector::new(self.placement.endpoint())
424                .try_connect()
425                .await;
426            match endpoint_connection {
427                RuntimeEndpointConnectionStatus::Missing
428                | RuntimeEndpointConnectionStatus::Stale => {
429                    return Ok(());
430                }
431                RuntimeEndpointConnectionStatus::Connected => {}
432                RuntimeEndpointConnectionStatus::Unavailable => {
433                    return Err(Error::EndpointUnavailable {
434                        context: "daemon endpoint",
435                        endpoint: self.placement.endpoint().path().to_path_buf(),
436                    });
437                }
438                #[cfg(not(unix))]
439                RuntimeEndpointConnectionStatus::UnsupportedTransport => {
440                    return Err(Error::UnsupportedTransport {
441                        context: "daemon control transport",
442                    });
443                }
444            }
445
446            let now = Instant::now();
447            if now >= deadline {
448                return Err(Error::EndpointStopTimeout {
449                    endpoint: self.placement.endpoint().path().to_path_buf(),
450                });
451            }
452
453            sleep(SHUTDOWN_WAIT_INTERVAL.min(deadline - now)).await;
454        }
455    }
456}
457
458#[cfg(test)]
459mod tests {
460    use crate::{
461        RuntimeDatabase,
462        connection::RuntimeEndpointConnectionStatus,
463        instance::RuntimeInstance,
464        placement::{PlacementRoot, PlacementSpec},
465    };
466
467    use super::{DaemonControlContext, DaemonShutdownDecision, DaemonStatusDecision};
468
469    mod status {
470        use super::*;
471
472        #[test]
473        fn from_endpoint_connection() {
474            let cases = [
475                (RuntimeEndpointConnectionStatus::Connected, "request_status"),
476                (RuntimeEndpointConnectionStatus::Missing, "already_stopped"),
477                (RuntimeEndpointConnectionStatus::Stale, "already_stopped"),
478                (
479                    RuntimeEndpointConnectionStatus::Unavailable,
480                    "fail_unavailable_endpoint",
481                ),
482            ];
483
484            for (endpoint_connection, expected_path) in cases {
485                let decision = DaemonStatusDecision::from(context_with(endpoint_connection));
486
487                assert_eq!(status_decision_path(&decision), expected_path);
488            }
489
490            #[cfg(not(unix))]
491            {
492                let decision = DaemonStatusDecision::from(context_with(
493                    RuntimeEndpointConnectionStatus::UnsupportedTransport,
494                ));
495
496                assert_eq!(
497                    status_decision_path(&decision),
498                    "fail_unsupported_transport"
499                );
500            }
501        }
502    }
503
504    mod shutdown_decision {
505        use super::*;
506
507        #[test]
508        fn from_endpoint_connection() {
509            let cases = [
510                (
511                    RuntimeEndpointConnectionStatus::Connected,
512                    "request_shutdown",
513                ),
514                (RuntimeEndpointConnectionStatus::Missing, "already_stopped"),
515                (RuntimeEndpointConnectionStatus::Stale, "already_stopped"),
516                (
517                    RuntimeEndpointConnectionStatus::Unavailable,
518                    "fail_unavailable_endpoint",
519                ),
520            ];
521
522            for (endpoint_connection, expected_path) in cases {
523                let decision = DaemonShutdownDecision::from(context_with(endpoint_connection));
524
525                assert_eq!(shutdown_decision_path(&decision), expected_path);
526            }
527
528            #[cfg(not(unix))]
529            {
530                let decision = DaemonShutdownDecision::from(context_with(
531                    RuntimeEndpointConnectionStatus::UnsupportedTransport,
532                ));
533
534                assert_eq!(
535                    shutdown_decision_path(&decision),
536                    "fail_unsupported_transport"
537                );
538            }
539        }
540    }
541
542    fn context_with(endpoint_connection: RuntimeEndpointConnectionStatus) -> DaemonControlContext {
543        DaemonControlContext {
544            placement: placement(),
545            endpoint_connection,
546        }
547    }
548
549    fn placement() -> PlacementSpec {
550        let tmp = tempfile::tempdir().unwrap();
551        let instance =
552            RuntimeInstance::from_database(&RuntimeDatabase::sqlite(tmp.path().join("synd.db")))
553                .unwrap();
554
555        PlacementSpec::from_instance(PlacementRoot::from(tmp.path().join("runtime")), instance)
556    }
557
558    fn shutdown_decision_path(decision: &DaemonShutdownDecision) -> &'static str {
559        match decision {
560            DaemonShutdownDecision::RequestShutdown { .. } => "request_shutdown",
561            DaemonShutdownDecision::AlreadyStopped { .. } => "already_stopped",
562            DaemonShutdownDecision::FailUnavailableEndpoint { .. } => "fail_unavailable_endpoint",
563            #[cfg(not(unix))]
564            DaemonShutdownDecision::FailUnsupportedTransport { .. } => "fail_unsupported_transport",
565        }
566    }
567
568    fn status_decision_path(decision: &DaemonStatusDecision) -> &'static str {
569        match decision {
570            DaemonStatusDecision::RequestStatus { .. } => "request_status",
571            DaemonStatusDecision::AlreadyStopped { .. } => "already_stopped",
572            DaemonStatusDecision::FailUnavailableEndpoint { .. } => "fail_unavailable_endpoint",
573            #[cfg(not(unix))]
574            DaemonStatusDecision::FailUnsupportedTransport { .. } => "fail_unsupported_transport",
575        }
576    }
577}