Skip to main content

inferlab_runtime/server/
readiness.rs

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    // A capture-armed server carries no readiness deadline
90    // ([[RFC-0004:C-WORKLOAD-PROFILING]]); the loop still terminates on
91    // readiness, process-group exit, or interruption.
92    let url = format!("http://{}:{}{}", endpoint.host, endpoint.port, path);
93    let mut attempts = 0_u32;
94    let mut diagnostic_attempts = Vec::new();
95    // The probe cadence backs off from POLL_INTERVAL to a cap: sub-second
96    // detection for ordinary startups without tens of thousands of no-op
97    // probes across a capture-armed unbounded wait. The sleep is clamped to
98    // the remaining deadline so a configured timeout fires within one
99    // interval.
100    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}