1use super::{
2 ProcessHandle, ProcessObserver, ProcessStatus, ReadinessAttemptEvidence, ReadinessEvidence,
3 ReadinessFailure, ReadinessFailureKind, ReadinessObserver, SystemProcessRuntime,
4 TargetRegistryMatchEvidence,
5};
6use crate::interrupt;
7use crate::operation_bound::{AttemptBound, OperationBound, OperationTerminalCause, Remaining};
8use crate::plan::{ProcessEndpointPlan, ReadinessPlan, TargetRegistryExpectedTarget};
9use std::thread;
10use std::time::Duration;
11
12const POLL_INTERVAL: Duration = Duration::from_millis(100);
13const MAX_PROBE_INTERVAL: Duration = Duration::from_secs(5);
14const READINESS_START_BOUNDARY: &str = "after_process_spawn_before_readiness_attempt";
15
16#[derive(Debug, thiserror::Error)]
17pub(super) enum ReadinessProbeError {
18 #[error("readiness operation deadline expired")]
19 Deadline,
20 #[error("bounded readiness attempt was unexpectedly unbounded")]
21 UnexpectedUnbounded,
22 #[error("{label} request failed: {source}")]
23 Request {
24 label: String,
25 #[source]
26 source: reqwest::Error,
27 },
28 #[error("{label} returned HTTP {status}")]
29 HttpStatus { label: String, status: u16 },
30 #[error("{label} returned invalid JSON: {source}")]
31 InvalidJson {
32 label: String,
33 #[source]
34 source: serde_json::Error,
35 },
36 #[error("target registry observation mismatch: {details}")]
37 RegistryMismatch { details: String },
38}
39
40pub(super) fn ensure_alive(status: ProcessStatus) -> Result<(), ReadinessFailure> {
41 if !status.queried {
42 return Err(readiness_failure(
43 ReadinessFailureKind::Exited,
44 status
45 .error
46 .unwrap_or_else(|| "failed to query server process group".to_owned()),
47 ));
48 }
49 if !status.alive {
50 return Err(readiness_failure(
51 ReadinessFailureKind::Exited,
52 status
53 .error
54 .unwrap_or_else(|| "server process group exited before readiness".to_owned()),
55 ));
56 }
57 Ok(())
58}
59
60fn readiness_failure(kind: ReadinessFailureKind, message: String) -> ReadinessFailure {
61 ReadinessFailure {
62 kind,
63 message,
64 timing: None,
65 diagnostic_attempts: Vec::new(),
66 }
67}
68
69pub(super) fn timed_readiness_failure(
70 mut failure: ReadinessFailure,
71 bound: &OperationBound,
72 terminal_cause: OperationTerminalCause,
73 diagnostic_attempts: Vec<ReadinessAttemptEvidence>,
74) -> ReadinessFailure {
75 failure.timing = Some(bound.timing(READINESS_START_BOUNDARY, terminal_cause));
76 failure.diagnostic_attempts = diagnostic_attempts;
77 failure
78}
79
80pub(super) fn wait_http_ready<R: ProcessObserver>(
81 runtime: &R,
82 handle: &ProcessHandle,
83 endpoint: &ProcessEndpointPlan,
84 path: &str,
85 attempt_timeout_seconds: u64,
86 bound: &OperationBound,
87 on_probe_failure: &mut dyn FnMut(&str),
88) -> Result<ReadinessEvidence, ReadinessFailure> {
89 let url = format!("http://{}:{}{}", endpoint.host, endpoint.port, path);
93 let mut attempts = 0_u32;
94 let mut diagnostic_attempts = Vec::new();
95 let mut probe_interval = POLL_INTERVAL;
101 loop {
102 ensure_readiness_active(bound, "no readiness probe completed").map_err(|failure| {
103 timed_readiness_failure(
104 failure,
105 bound,
106 OperationTerminalCause::TimedOut,
107 diagnostic_attempts.clone(),
108 )
109 })?;
110 if interrupt::received() {
111 return Err(timed_readiness_failure(
112 readiness_failure(
113 ReadinessFailureKind::Interrupted,
114 "server startup was interrupted".to_owned(),
115 ),
116 bound,
117 OperationTerminalCause::Interrupted,
118 diagnostic_attempts,
119 ));
120 }
121 let status_attempt = readiness_attempt(bound, attempt_timeout_seconds);
122 let status_effective_bound_ms = status_attempt.configured_ms().unwrap_or_default();
123 let status_bound = status_attempt.into_operation_bound();
124 let status = runtime.status_with_bound(handle, &status_bound);
125 diagnostic_attempts = vec![process_status_evidence(&status, status_effective_bound_ms)];
126 if !status.queried && status_bound.is_expired() && !bound.is_expired() {
127 let last_error = status
128 .error
129 .as_deref()
130 .unwrap_or("process status attempt deadline expired");
131 on_probe_failure(last_error);
132 sleep_within_readiness(bound, probe_interval);
133 probe_interval = (probe_interval * 2).min(MAX_PROBE_INTERVAL);
134 continue;
135 }
136 ensure_readiness_active(
137 bound,
138 "the server process status attempt did not complete in time",
139 )
140 .map_err(|failure| {
141 timed_readiness_failure(
142 failure,
143 bound,
144 OperationTerminalCause::TimedOut,
145 diagnostic_attempts.clone(),
146 )
147 })?;
148 if interrupt::received() {
149 return Err(timed_readiness_failure(
150 readiness_failure(
151 ReadinessFailureKind::Interrupted,
152 "server startup was interrupted".to_owned(),
153 ),
154 bound,
155 OperationTerminalCause::Interrupted,
156 diagnostic_attempts,
157 ));
158 }
159 ensure_alive(status).map_err(|failure| {
160 timed_readiness_failure(
161 failure,
162 bound,
163 OperationTerminalCause::Failed,
164 diagnostic_attempts.clone(),
165 )
166 })?;
167 attempts = attempts.saturating_add(1);
168 let attempt = probe_http_attempt(
169 &endpoint.host,
170 endpoint.port,
171 path,
172 bound,
173 attempt_timeout_seconds,
174 );
175 let effective_bound_ms = attempt.effective_bound_ms;
176 let last_error = match attempt.outcome {
177 Ok(()) => {
178 diagnostic_attempts.push(ReadinessAttemptEvidence {
179 operation: "http_readiness".to_owned(),
180 effective_bound_ms,
181 succeeded: true,
182 error: None,
183 });
184 let ready_unix_ms = unix_time_millis().map_err(|failure| {
185 timed_readiness_failure(
186 failure,
187 bound,
188 OperationTerminalCause::Failed,
189 diagnostic_attempts.clone(),
190 )
191 })?;
192 ensure_readiness_active(
193 bound,
194 "the readiness response completed after the deadline",
195 )
196 .map_err(|failure| {
197 timed_readiness_failure(
198 failure,
199 bound,
200 OperationTerminalCause::TimedOut,
201 diagnostic_attempts.clone(),
202 )
203 })?;
204 return Ok(ReadinessEvidence::Http {
205 url,
206 attempts,
207 ready_unix_ms,
208 timing: bound
209 .timing(READINESS_START_BOUNDARY, OperationTerminalCause::Succeeded),
210 diagnostic_attempts,
211 });
212 }
213 Err(error) => {
214 let error = error.to_string();
215 diagnostic_attempts.push(ReadinessAttemptEvidence {
216 operation: "http_readiness".to_owned(),
217 effective_bound_ms,
218 succeeded: false,
219 error: Some(error.clone()),
220 });
221 error
222 }
223 };
224 on_probe_failure(&last_error);
225 if bound.is_expired() {
226 let timeout_seconds = readiness_timeout_seconds(bound);
227 return Err(timed_readiness_failure(
228 readiness_failure(
229 ReadinessFailureKind::Timeout,
230 format!(
231 "server did not become ready within {timeout_seconds} seconds; last probe error: {last_error}"
232 ),
233 ),
234 bound,
235 OperationTerminalCause::TimedOut,
236 diagnostic_attempts,
237 ));
238 }
239 sleep_within_readiness(bound, probe_interval);
240 probe_interval = (probe_interval * 2).min(MAX_PROBE_INTERVAL);
241 }
242}
243
244pub(super) struct HttpTargetRegistryProbe<'a> {
245 pub(super) readiness_path: &'a str,
246 pub(super) registry_path: &'a str,
247 pub(super) targets_field: &'a str,
248 pub(super) target_url_field: &'a str,
249 pub(super) target_role_field: &'a str,
250 pub(super) target_healthy_field: &'a str,
251 pub(super) target_bootstrap_port_field: &'a str,
252 pub(super) expected_targets: &'a [TargetRegistryExpectedTarget],
253}
254
255fn sleep_within_readiness(bound: &OperationBound, cadence: Duration) {
256 match bound.remaining() {
257 Remaining::Finite(remaining) => thread::sleep(cadence.min(remaining)),
258 Remaining::Expired => {}
259 Remaining::Unbounded => thread::sleep(cadence),
260 }
261}
262
263fn process_status_evidence(
264 status: &ProcessStatus,
265 effective_bound_ms: u64,
266) -> ReadinessAttemptEvidence {
267 let succeeded = status.queried && status.alive;
268 ReadinessAttemptEvidence {
269 operation: "process_status".to_owned(),
270 effective_bound_ms,
271 succeeded,
272 error: if succeeded {
273 None
274 } else {
275 Some(
276 status
277 .error
278 .clone()
279 .unwrap_or_else(|| "server process group is not alive".to_owned()),
280 )
281 },
282 }
283}
284
285fn ensure_readiness_active(
286 bound: &OperationBound,
287 last_error: &str,
288) -> Result<(), ReadinessFailure> {
289 if !bound.is_expired() {
290 return Ok(());
291 }
292 Err(readiness_failure(
293 ReadinessFailureKind::Timeout,
294 format!(
295 "server did not become ready within {} seconds; last probe error: {last_error}",
296 readiness_timeout_seconds(bound)
297 ),
298 ))
299}
300
301fn readiness_timeout_seconds(bound: &OperationBound) -> u64 {
302 bound.configured_ms().unwrap_or_default() / 1_000
303}
304
305fn attempt_remaining(attempt: &AttemptBound) -> Result<Duration, ReadinessProbeError> {
306 match attempt.remaining() {
307 Remaining::Finite(remaining) => Ok(remaining),
308 Remaining::Expired => Err(ReadinessProbeError::Deadline),
309 Remaining::Unbounded => Err(ReadinessProbeError::UnexpectedUnbounded),
310 }
311}
312
313pub(super) fn wait_http_target_registry_ready(
314 status: impl Fn(&OperationBound) -> ProcessStatus,
315 endpoint: &ProcessEndpointPlan,
316 probe: HttpTargetRegistryProbe<'_>,
317 attempt_timeout_seconds: u64,
318 bound: &OperationBound,
319 on_probe_failure: &mut dyn FnMut(&str),
320) -> Result<ReadinessEvidence, ReadinessFailure> {
321 let readiness_url = format!(
322 "http://{}:{}{}",
323 endpoint.host, endpoint.port, probe.readiness_path
324 );
325 let registry_url = format!(
326 "http://{}:{}{}",
327 endpoint.host, endpoint.port, probe.registry_path
328 );
329 let mut attempts = 0_u32;
330 let mut diagnostic_attempts = Vec::new();
331 let mut probe_interval = POLL_INTERVAL;
332 loop {
333 ensure_readiness_active(bound, "no readiness probe completed").map_err(|failure| {
334 timed_readiness_failure(
335 failure,
336 bound,
337 OperationTerminalCause::TimedOut,
338 diagnostic_attempts.clone(),
339 )
340 })?;
341 if interrupt::received() {
342 return Err(timed_readiness_failure(
343 readiness_failure(
344 ReadinessFailureKind::Interrupted,
345 "server startup was interrupted".to_owned(),
346 ),
347 bound,
348 OperationTerminalCause::Interrupted,
349 diagnostic_attempts,
350 ));
351 }
352 let status_attempt = readiness_attempt(bound, attempt_timeout_seconds);
353 let status_effective_bound_ms = status_attempt.configured_ms().unwrap_or_default();
354 let status_bound = status_attempt.into_operation_bound();
355 let process_status = status(&status_bound);
356 diagnostic_attempts = vec![process_status_evidence(
357 &process_status,
358 status_effective_bound_ms,
359 )];
360 if !process_status.queried && status_bound.is_expired() && !bound.is_expired() {
361 let last_error = process_status
362 .error
363 .as_deref()
364 .unwrap_or("process status attempt deadline expired");
365 on_probe_failure(last_error);
366 sleep_within_readiness(bound, probe_interval);
367 probe_interval = (probe_interval * 2).min(MAX_PROBE_INTERVAL);
368 continue;
369 }
370 ensure_readiness_active(
371 bound,
372 "the server process status attempt did not complete in time",
373 )
374 .map_err(|failure| {
375 timed_readiness_failure(
376 failure,
377 bound,
378 OperationTerminalCause::TimedOut,
379 diagnostic_attempts.clone(),
380 )
381 })?;
382 if interrupt::received() {
383 return Err(timed_readiness_failure(
384 readiness_failure(
385 ReadinessFailureKind::Interrupted,
386 "server startup was interrupted".to_owned(),
387 ),
388 bound,
389 OperationTerminalCause::Interrupted,
390 diagnostic_attempts,
391 ));
392 }
393 ensure_alive(process_status).map_err(|failure| {
394 timed_readiness_failure(
395 failure,
396 bound,
397 OperationTerminalCause::Failed,
398 diagnostic_attempts.clone(),
399 )
400 })?;
401 attempts = attempts.saturating_add(1);
402 let public_attempt = probe_http_attempt(
403 &endpoint.host,
404 endpoint.port,
405 probe.readiness_path,
406 bound,
407 attempt_timeout_seconds,
408 );
409 let public_effective_bound_ms = public_attempt.effective_bound_ms;
410 let last_error = match public_attempt.outcome {
411 Ok(()) => {
412 let registry_attempt = probe_target_registry_attempt(
413 &endpoint.host,
414 endpoint.port,
415 &probe,
416 bound,
417 attempt_timeout_seconds,
418 );
419 let registry_effective_bound_ms = registry_attempt.effective_bound_ms;
420 match registry_attempt.outcome {
421 Ok(matched_targets) => {
422 diagnostic_attempts.extend([
423 ReadinessAttemptEvidence {
424 operation: "public_http_readiness".to_owned(),
425 effective_bound_ms: public_effective_bound_ms,
426 succeeded: true,
427 error: None,
428 },
429 ReadinessAttemptEvidence {
430 operation: "target_registry".to_owned(),
431 effective_bound_ms: registry_effective_bound_ms,
432 succeeded: true,
433 error: None,
434 },
435 ]);
436 let ready_unix_ms = unix_time_millis().map_err(|failure| {
437 timed_readiness_failure(
438 failure,
439 bound,
440 OperationTerminalCause::Failed,
441 diagnostic_attempts.clone(),
442 )
443 })?;
444 ensure_readiness_active(
445 bound,
446 "the target registry response completed after the deadline",
447 )
448 .map_err(|failure| {
449 timed_readiness_failure(
450 failure,
451 bound,
452 OperationTerminalCause::TimedOut,
453 diagnostic_attempts.clone(),
454 )
455 })?;
456 return Ok(ReadinessEvidence::HttpTargetRegistry {
457 readiness_url,
458 registry_url,
459 attempts,
460 ready_unix_ms,
461 matched_targets,
462 timing: bound.timing(
463 READINESS_START_BOUNDARY,
464 OperationTerminalCause::Succeeded,
465 ),
466 diagnostic_attempts,
467 });
468 }
469 Err(error) => {
470 let error = error.to_string();
471 diagnostic_attempts.extend([
472 ReadinessAttemptEvidence {
473 operation: "public_http_readiness".to_owned(),
474 effective_bound_ms: public_effective_bound_ms,
475 succeeded: true,
476 error: None,
477 },
478 ReadinessAttemptEvidence {
479 operation: "target_registry".to_owned(),
480 effective_bound_ms: registry_effective_bound_ms,
481 succeeded: false,
482 error: Some(error.clone()),
483 },
484 ]);
485 error
486 }
487 }
488 }
489 Err(error) => {
490 let error = error.to_string();
491 diagnostic_attempts.push(ReadinessAttemptEvidence {
492 operation: "public_http_readiness".to_owned(),
493 effective_bound_ms: public_effective_bound_ms,
494 succeeded: false,
495 error: Some(error.clone()),
496 });
497 format!("public readiness probe failed: {error}")
498 }
499 };
500 on_probe_failure(&last_error);
501 if bound.is_expired() {
502 let timeout_seconds = readiness_timeout_seconds(bound);
503 return Err(timed_readiness_failure(
504 readiness_failure(
505 ReadinessFailureKind::Timeout,
506 format!(
507 "server did not become ready within {timeout_seconds} seconds; last probe error: {last_error}"
508 ),
509 ),
510 bound,
511 OperationTerminalCause::TimedOut,
512 diagnostic_attempts,
513 ));
514 }
515 sleep_within_readiness(bound, probe_interval);
516 probe_interval = (probe_interval * 2).min(MAX_PROBE_INTERVAL);
517 }
518}
519
520fn probe_target_registry_attempt(
521 host: &str,
522 port: u16,
523 probe: &HttpTargetRegistryProbe<'_>,
524 bound: &OperationBound,
525 attempt_timeout_seconds: u64,
526) -> ProbeAttempt<Vec<TargetRegistryMatchEvidence>> {
527 let response = probe_http_json_attempt(
528 host,
529 port,
530 probe.registry_path,
531 "target registry",
532 bound,
533 attempt_timeout_seconds,
534 );
535 let effective_bound_ms = response.effective_bound_ms;
536 let outcome = response
537 .outcome
538 .and_then(|response| match_target_registry(&response, probe, bound));
539 ProbeAttempt {
540 effective_bound_ms,
541 outcome,
542 }
543}
544
545pub(super) fn match_target_registry(
546 response: &serde_json::Value,
547 probe: &HttpTargetRegistryProbe<'_>,
548 bound: &OperationBound,
549) -> Result<Vec<TargetRegistryMatchEvidence>, ReadinessProbeError> {
550 readiness_remaining(bound)?;
551 let targets = response
552 .get(probe.targets_field)
553 .and_then(serde_json::Value::as_array)
554 .ok_or_else(|| ReadinessProbeError::RegistryMismatch {
555 details: format!(
556 "target registry response has no array field {:?}",
557 probe.targets_field
558 ),
559 })?;
560 let mut evidence = Vec::with_capacity(probe.expected_targets.len());
561 for expected in probe.expected_targets {
562 readiness_remaining(bound)?;
563 let matches: Vec<&serde_json::Map<String, serde_json::Value>> = targets
564 .iter()
565 .filter_map(serde_json::Value::as_object)
566 .filter(|target| {
567 target
568 .get(probe.target_url_field)
569 .and_then(serde_json::Value::as_str)
570 == Some(expected.url.as_str())
571 && target
572 .get(probe.target_role_field)
573 .and_then(serde_json::Value::as_str)
574 == Some(expected.role.as_str())
575 })
576 .collect();
577 let target = match matches.as_slice() {
578 [] => {
579 return Err(ReadinessProbeError::RegistryMismatch {
580 details: format!(
581 "target registry has no {:?} target at {:?}",
582 expected.role, expected.url
583 ),
584 });
585 }
586 [target] => *target,
587 _ => {
588 return Err(ReadinessProbeError::RegistryMismatch {
589 details: format!(
590 "target registry has multiple {:?} targets at {:?}",
591 expected.role, expected.url
592 ),
593 });
594 }
595 };
596 let healthy = target
597 .get(probe.target_healthy_field)
598 .and_then(serde_json::Value::as_bool)
599 .ok_or_else(|| ReadinessProbeError::RegistryMismatch {
600 details: format!(
601 "target registry entry for {:?} at {:?} has no boolean {:?} field",
602 expected.role, expected.url, probe.target_healthy_field
603 ),
604 })?;
605 if !healthy {
606 return Err(ReadinessProbeError::RegistryMismatch {
607 details: format!(
608 "target registry entry for {:?} at {:?} is not healthy",
609 expected.role, expected.url
610 ),
611 });
612 }
613 let bootstrap_port = match target.get(probe.target_bootstrap_port_field) {
614 None | Some(serde_json::Value::Null) => None,
615 Some(value) => {
616 let port = value.as_u64().and_then(|port| u16::try_from(port).ok());
617 Some(port.ok_or_else(|| ReadinessProbeError::RegistryMismatch {
618 details: format!(
619 "target registry entry for {:?} at {:?} has invalid {:?}",
620 expected.role, expected.url, probe.target_bootstrap_port_field
621 ),
622 })?)
623 }
624 };
625 if let Some(expected_port) = expected.bootstrap_port
626 && bootstrap_port != Some(expected_port)
627 {
628 return Err(ReadinessProbeError::RegistryMismatch {
629 details: format!(
630 "target registry entry for {:?} at {:?} has bootstrap port {bootstrap_port:?}, expected {expected_port}",
631 expected.role, expected.url
632 ),
633 });
634 }
635 evidence.push(TargetRegistryMatchEvidence {
636 url: expected.url.clone(),
637 role: expected.role.clone(),
638 healthy,
639 bootstrap_port,
640 });
641 }
642 readiness_remaining(bound)?;
643 Ok(evidence)
644}
645
646struct ProbeAttempt<T> {
647 effective_bound_ms: u64,
648 outcome: Result<T, ReadinessProbeError>,
649}
650
651fn readiness_attempt(bound: &OperationBound, attempt_timeout_seconds: u64) -> AttemptBound {
652 bound.attempt(Some(Duration::from_secs(attempt_timeout_seconds)))
653}
654
655#[cfg(test)]
656pub(super) fn probe_http(
657 host: &str,
658 port: u16,
659 path: &str,
660 bound: &OperationBound,
661 attempt_timeout_seconds: u64,
662) -> Result<(), ReadinessProbeError> {
663 probe_http_attempt(host, port, path, bound, attempt_timeout_seconds).outcome
664}
665
666fn probe_http_attempt(
667 host: &str,
668 port: u16,
669 path: &str,
670 bound: &OperationBound,
671 attempt_timeout_seconds: u64,
672) -> ProbeAttempt<()> {
673 let attempt = readiness_attempt(bound, attempt_timeout_seconds);
674 let effective_bound_ms = attempt.configured_ms().unwrap_or_default();
675 let outcome = (|| {
676 let url = format!("http://{host}:{port}{path}");
677 let timeout = attempt_remaining(&attempt)?;
678 let client = reqwest::blocking::Client::builder()
679 .timeout(timeout)
680 .connect_timeout(timeout)
681 .redirect(reqwest::redirect::Policy::none())
682 .no_proxy()
683 .build()
684 .map_err(|source| readiness_request_error("readiness", source))?;
685 let mut response = client
686 .get(&url)
687 .send()
688 .map_err(|source| readiness_request_error("readiness", source))?;
689 let status = response.status().as_u16();
690 response
691 .copy_to(&mut std::io::sink())
692 .map_err(|source| readiness_request_error("readiness", source))?;
693 if (200..300).contains(&status) {
694 Ok(())
695 } else {
696 Err(ReadinessProbeError::HttpStatus {
697 label: "readiness".to_owned(),
698 status,
699 })
700 }
701 })();
702 ProbeAttempt {
703 effective_bound_ms,
704 outcome,
705 }
706}
707
708#[cfg(test)]
709pub(super) fn probe_http_json(
710 host: &str,
711 port: u16,
712 path: &str,
713 label: &str,
714 bound: &OperationBound,
715 attempt_timeout_seconds: u64,
716) -> Result<serde_json::Value, ReadinessProbeError> {
717 probe_http_json_attempt(host, port, path, label, bound, attempt_timeout_seconds).outcome
718}
719
720fn probe_http_json_attempt(
721 host: &str,
722 port: u16,
723 path: &str,
724 label: &str,
725 bound: &OperationBound,
726 attempt_timeout_seconds: u64,
727) -> ProbeAttempt<serde_json::Value> {
728 let attempt = readiness_attempt(bound, attempt_timeout_seconds);
729 let effective_bound_ms = attempt.configured_ms().unwrap_or_default();
730 let outcome = (|| {
731 let url = format!("http://{host}:{port}{path}");
732 let timeout = attempt_remaining(&attempt)?;
733 let client = reqwest::blocking::Client::builder()
734 .timeout(timeout)
735 .connect_timeout(timeout)
736 .redirect(reqwest::redirect::Policy::none())
737 .no_proxy()
738 .build()
739 .map_err(|source| readiness_request_error(label, source))?;
740 let response = client
741 .get(&url)
742 .send()
743 .map_err(|source| readiness_request_error(label, source))?;
744 let status = response.status().as_u16();
745 if !(200..300).contains(&status) {
746 return Err(ReadinessProbeError::HttpStatus {
747 label: label.to_owned(),
748 status,
749 });
750 }
751 let body = response
752 .bytes()
753 .map_err(|source| readiness_request_error(label, source))?;
754 let value =
755 serde_json::from_slice(&body).map_err(|source| ReadinessProbeError::InvalidJson {
756 label: label.to_owned(),
757 source,
758 })?;
759 readiness_remaining(bound)?;
760 Ok(value)
761 })();
762 ProbeAttempt {
763 effective_bound_ms,
764 outcome,
765 }
766}
767
768fn readiness_request_error(label: &str, source: reqwest::Error) -> ReadinessProbeError {
769 if source.is_timeout() {
770 ReadinessProbeError::Deadline
771 } else {
772 ReadinessProbeError::Request {
773 label: label.to_owned(),
774 source,
775 }
776 }
777}
778
779fn readiness_remaining(bound: &OperationBound) -> Result<(), ReadinessProbeError> {
780 match bound.remaining() {
781 Remaining::Expired => Err(ReadinessProbeError::Deadline),
782 Remaining::Finite(_) | Remaining::Unbounded => Ok(()),
783 }
784}
785
786pub(super) fn unix_time_millis() -> Result<u64, ReadinessFailure> {
787 std::time::SystemTime::now()
788 .duration_since(std::time::UNIX_EPOCH)
789 .map(crate::operation_bound::duration_millis)
790 .map_err(|error| {
791 readiness_failure(
792 ReadinessFailureKind::Exited,
793 format!("system clock is before Unix epoch: {error}"),
794 )
795 })
796}
797
798pub(super) fn wait_process_alive_ready(
799 status: impl Fn(&OperationBound) -> ProcessStatus,
800 attempt_timeout_seconds: u64,
801 bound: &OperationBound,
802 on_probe_failure: &mut dyn FnMut(&str),
803) -> Result<ReadinessEvidence, ReadinessFailure> {
804 loop {
805 ensure_readiness_active(bound, "no process status attempt completed").map_err(
806 |failure| {
807 timed_readiness_failure(
808 failure,
809 bound,
810 OperationTerminalCause::TimedOut,
811 Vec::new(),
812 )
813 },
814 )?;
815 if interrupt::received() {
816 return Err(timed_readiness_failure(
817 readiness_failure(
818 ReadinessFailureKind::Interrupted,
819 "server startup was interrupted".to_owned(),
820 ),
821 bound,
822 OperationTerminalCause::Interrupted,
823 Vec::new(),
824 ));
825 }
826 let attempt = readiness_attempt(bound, attempt_timeout_seconds);
827 let effective_bound_ms = attempt.configured_ms().unwrap_or_default();
828 let attempt_bound = attempt.into_operation_bound();
829 let process_status = status(&attempt_bound);
830 let diagnostic_attempts =
831 vec![process_status_evidence(&process_status, effective_bound_ms)];
832 if !process_status.queried && attempt_bound.is_expired() && !bound.is_expired() {
833 on_probe_failure(
834 process_status
835 .error
836 .as_deref()
837 .unwrap_or("process status attempt deadline expired"),
838 );
839 sleep_within_readiness(bound, POLL_INTERVAL);
840 continue;
841 }
842 ensure_readiness_active(
843 bound,
844 "the server process status attempt did not complete in time",
845 )
846 .map_err(|failure| {
847 timed_readiness_failure(
848 failure,
849 bound,
850 OperationTerminalCause::TimedOut,
851 diagnostic_attempts.clone(),
852 )
853 })?;
854 ensure_alive(process_status).map_err(|failure| {
855 timed_readiness_failure(
856 failure,
857 bound,
858 OperationTerminalCause::Failed,
859 diagnostic_attempts.clone(),
860 )
861 })?;
862 return Ok(ReadinessEvidence::ProcessAlive {
863 ready_unix_ms: unix_time_millis().map_err(|failure| {
864 timed_readiness_failure(
865 failure,
866 bound,
867 OperationTerminalCause::Failed,
868 diagnostic_attempts.clone(),
869 )
870 })?,
871 timing: bound.timing(READINESS_START_BOUNDARY, OperationTerminalCause::Succeeded),
872 diagnostic_attempts,
873 });
874 }
875}
876
877impl ReadinessObserver for SystemProcessRuntime {
878 fn wait_ready(
879 &self,
880 handle: &ProcessHandle,
881 endpoint: &ProcessEndpointPlan,
882 readiness: &ReadinessPlan,
883 bound: &OperationBound,
884 on_probe_failure: &mut dyn FnMut(&str),
885 ) -> Result<ReadinessEvidence, ReadinessFailure> {
886 match readiness {
887 ReadinessPlan::ProcessAlive {
888 attempt_timeout_seconds,
889 ..
890 } => wait_process_alive_ready(
891 |bound| self.status_with_bound(handle, bound),
892 *attempt_timeout_seconds,
893 bound,
894 on_probe_failure,
895 ),
896 ReadinessPlan::Http {
897 path,
898 attempt_timeout_seconds,
899 ..
900 } => wait_http_ready(
901 self,
902 handle,
903 endpoint,
904 path,
905 *attempt_timeout_seconds,
906 bound,
907 on_probe_failure,
908 ),
909 ReadinessPlan::HttpTargetRegistry {
910 readiness_path,
911 registry_path,
912 targets_field,
913 target_url_field,
914 target_role_field,
915 target_healthy_field,
916 target_bootstrap_port_field,
917 expected_targets,
918 attempt_timeout_seconds,
919 ..
920 } => wait_http_target_registry_ready(
921 |bound| self.status_with_bound(handle, bound),
922 endpoint,
923 HttpTargetRegistryProbe {
924 readiness_path,
925 registry_path,
926 targets_field,
927 target_url_field,
928 target_role_field,
929 target_healthy_field,
930 target_bootstrap_port_field,
931 expected_targets,
932 },
933 *attempt_timeout_seconds,
934 bound,
935 on_probe_failure,
936 ),
937 }
938 }
939}