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#[derive(Debug, Clone)]
271struct DaemonControlContext {
272 placement: PlacementSpec,
273 endpoint_connection: RuntimeEndpointConnectionStatus,
274}
275
276#[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#[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
360struct 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
408struct 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}