use std::future::Future;
use std::time::Duration;
use tokio::time::Instant;
use crate::bounded_git::GitDeadline;
use crate::client::completion::{
certify, classify, is_settled_success_claim, CertificationStage, Disposition, Verdict,
};
use crate::client::envelope::{Operation, Outcome, ResultEnvelope};
use crate::client::session::{observe, Connection, Observation};
use crate::client::transport::Wake;
use crate::web::remote_control_api::dto::{BlockerKind, ChangeBlocker, OwnerExecutionContract};
const POLL_INTERVAL: Duration = Duration::from_millis(250);
const RETRY_INTERVAL: Duration = Duration::from_millis(50);
const GIT_INVOCATION_BUDGET: Duration = Duration::from_secs(30);
const UNCERTIFIED_CLAIM_ROUNDS: usize = 2;
async fn within<F: Future>(deadline: Option<Instant>, future: F) -> Option<F::Output> {
match deadline {
Some(deadline) => tokio::time::timeout_at(deadline, future).await.ok(),
None => Some(future.await),
}
}
fn reached(deadline: Option<Instant>) -> bool {
deadline.is_some_and(|deadline| Instant::now() >= deadline)
}
fn git_deadline(deadline: Option<Instant>) -> GitDeadline {
match deadline {
Some(deadline) => GitDeadline::Operation(deadline),
None => GitDeadline::PerChild(GIT_INVOCATION_BUDGET),
}
}
pub async fn run(
connection: &Connection,
change_id: &str,
timeout: Option<Duration>,
) -> ResultEnvelope {
let mut diagnostics = Diagnostics::new(timeout);
let deadline = timeout.map(|timeout| Instant::now() + timeout);
let initial = match within(deadline, observe(connection, Some(change_id))).await {
Some(Ok(initial)) => initial,
Some(Err(error)) => return error.into_envelope(Operation::Wait).with_change(change_id),
None => {
return diagnostics.expire(change_id, None, TimeoutStage::InitialObservation, None);
}
};
diagnostics.record(&initial, change_id);
let instance_id = initial.instance_id.clone();
let Some(contract) = initial.contract.contract.clone() else {
return ResultEnvelope::new(Operation::Wait, Outcome::UnsupportedTerminalMode)
.with_change(change_id)
.with_instance(Some(instance_id))
.with_message(
"the owner published no execution contract, so this client cannot tell what would \
prove the change finished. Waiting without a terminal mode could only end in a \
timeout",
);
};
let Some(repo_root) = connection.repo_root().map(|root| root.to_path_buf()) else {
return ResultEnvelope::new(Operation::Wait, Outcome::NotInRepository)
.with_change(change_id)
.with_instance(Some(instance_id))
.with_message(
"completion is certified from repository evidence, so `wait` must run inside the \
owner's Git repository",
);
};
let mut observation = initial;
let mut uncertified_rounds = 0usize;
let mut first_observation = true;
loop {
match evaluate(
&observation,
&instance_id,
change_id,
&repo_root,
&contract,
deadline,
first_observation,
)
.await
{
Step::Settled(envelope) => return *envelope,
Step::Expired { stage, detail } => {
return diagnostics.expire(change_id, Some(&instance_id), stage, detail)
}
Step::UncertifiedClaim { status, detail } => {
uncertified_rounds += 1;
if uncertified_rounds >= UNCERTIFIED_CLAIM_ROUNDS {
return requires_action_envelope(
change_id,
&instance_id,
&status,
Some(detail),
None,
format!(
"'{change_id}' is at settled status '{status}', but repository \
evidence still does not prove completion"
),
);
}
if reached(deadline) {
return diagnostics.expire(
change_id,
Some(&instance_id),
TimeoutStage::ObservingOwner,
Some(detail),
);
}
}
Step::KeepObserving { detail } => {
uncertified_rounds = 0;
if reached(deadline) {
return diagnostics.expire(
change_id,
Some(&instance_id),
TimeoutStage::ObservingOwner,
detail,
);
}
}
}
first_observation = false;
let budget = match deadline {
Some(deadline) => {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return diagnostics.expire(
change_id,
Some(&instance_id),
TimeoutStage::ObservingOwner,
None,
);
}
POLL_INTERVAL.min(remaining)
}
None => POLL_INTERVAL,
};
let wake = connection
.client()
.wake_on_activity(observation.event_sequence, &instance_id, budget)
.await;
debug_assert!(matches!(wake, Wake::Activity | Wake::Gap | Wake::Idle));
observation = loop {
let Some(result) = within(deadline, observe(connection, Some(change_id))).await else {
return diagnostics.expire(
change_id,
Some(&instance_id),
TimeoutStage::ObservingOwner,
None,
);
};
match result {
Ok(next) => break next,
Err(error) if error.is_transient() => {
if reached(deadline) {
return diagnostics.expire(
change_id,
Some(&instance_id),
TimeoutStage::ObservingOwner,
Some(error.message().to_string()),
);
}
let backoff = match deadline {
Some(deadline) => {
RETRY_INTERVAL.min(deadline.saturating_duration_since(Instant::now()))
}
None => RETRY_INTERVAL,
};
tokio::time::sleep(backoff).await;
}
Err(error) => {
match certify(change_id, &repo_root, &contract, git_deadline(deadline)).await {
Verdict::Completed { evidence } => {
return completed_envelope(change_id, &instance_id, &contract, evidence)
}
Verdict::DeadlineExpired { stage } => {
return diagnostics.expire(
change_id,
Some(&instance_id),
stage.into(),
None,
)
}
_ => {}
}
return error.into_envelope(Operation::Wait).with_change(change_id);
}
}
};
diagnostics.record(&observation, change_id);
}
}
enum Step {
Settled(Box<ResultEnvelope>),
KeepObserving { detail: Option<String> },
UncertifiedClaim { status: String, detail: String },
Expired {
stage: TimeoutStage,
detail: Option<String>,
},
}
impl Step {
fn settled(envelope: ResultEnvelope) -> Self {
Self::Settled(Box::new(envelope))
}
}
async fn evaluate(
observation: &Observation,
instance_id: &str,
change_id: &str,
repo_root: &std::path::Path,
contract: &OwnerExecutionContract,
deadline: Option<Instant>,
first_observation: bool,
) -> Step {
if observation.instance_id != instance_id {
let verdict = certify(change_id, repo_root, contract, git_deadline(deadline)).await;
return match verdict {
Verdict::Completed { evidence } => Step::settled(completed_envelope(
change_id,
instance_id,
contract,
evidence,
)),
Verdict::DeadlineExpired { stage } => Step::Expired {
stage: stage.into(),
detail: None,
},
_ => Step::settled(
ResultEnvelope::new(Operation::Wait, Outcome::OwnerRestarted)
.with_change(change_id)
.with_instance(Some(observation.instance_id.clone()))
.with_message(
"the socket began serving a different owner incarnation and current \
repository evidence does not prove completion on its own",
)
.with_detail(serde_json::json!({
"expected_instance_id": instance_id,
"observed_instance_id": observation.instance_id,
})),
),
};
}
if let Some(process_error) = observation.state.snapshot.process_error.clone() {
return Step::settled(
ResultEnvelope::new(Operation::Wait, Outcome::ProcessFailed)
.with_change(change_id)
.with_instance(Some(instance_id.to_string()))
.with_message(format!("the owner reported a fatal error: {process_error}")),
);
}
let change = observation.change(change_id);
if first_observation && change.is_none() {
return match certify(change_id, repo_root, contract, git_deadline(deadline)).await {
Verdict::Completed { evidence } => Step::settled(completed_envelope(
change_id,
instance_id,
contract,
evidence,
)),
Verdict::DeadlineExpired { stage } => Step::Expired {
stage: stage.into(),
detail: None,
},
Verdict::NotCompleted { detail }
| Verdict::Broken { detail }
| Verdict::Unsupported { detail } => {
Step::settled(unknown_change_envelope(change_id, instance_id, detail))
}
};
}
let status = change.map(|change| change.display_status.as_str());
let blocker = change.and_then(|change| change.blocker.as_ref());
let error_detail = change.and_then(|change| change.error_detail.clone());
match classify(status, blocker.map(|blocker| blocker.kind)) {
Disposition::Rejected => {
let detail = error_detail.unwrap_or_else(|| "the change was rejected".to_string());
return Step::settled(
ResultEnvelope::new(Operation::Wait, Outcome::ChangeRejected)
.with_change(change_id)
.with_instance(Some(instance_id.to_string()))
.with_message(detail),
);
}
Disposition::RequiresAction => {
let status = status.unwrap_or_default();
let external = blocker.filter(|blocker| blocker.kind == BlockerKind::External);
let message = match external {
Some(_) => format!(
"'{change_id}' is blocked on an external prerequisite the owner cannot clear \
by itself, so it will not advance without a new operator action"
),
None => format!(
"'{change_id}' is at status '{status}', which the owner cannot advance without \
a new operator action"
),
};
return Step::settled(requires_action_envelope(
change_id,
instance_id,
status,
error_detail.or_else(|| external.and_then(|blocker| blocker.detail.clone())),
external,
message,
));
}
Disposition::KeepObserving => {
return Step::KeepObserving {
detail: status.map(|status| {
format!("'{change_id}' is at status '{status}' with no terminal evidence yet")
}),
}
}
Disposition::Certify => {}
}
match certify(change_id, repo_root, contract, git_deadline(deadline)).await {
Verdict::Completed { evidence } => Step::settled(completed_envelope(
change_id,
instance_id,
contract,
evidence,
)),
Verdict::NotCompleted { detail } if is_settled_success_claim(status) => {
Step::UncertifiedClaim {
status: status.unwrap_or_default().to_string(),
detail,
}
}
Verdict::NotCompleted { detail } => Step::KeepObserving {
detail: Some(detail),
},
Verdict::Broken { detail } => Step::settled(
ResultEnvelope::new(Operation::Wait, Outcome::EvidenceError)
.with_change(change_id)
.with_instance(Some(instance_id.to_string()))
.with_message(format!(
"repository completion evidence for '{change_id}' is unusable: {detail}"
))
.with_detail(contract_detail(contract)),
),
Verdict::Unsupported { detail } => Step::settled(
ResultEnvelope::new(Operation::Wait, Outcome::UnsupportedTerminalMode)
.with_change(change_id)
.with_instance(Some(instance_id.to_string()))
.with_message(detail)
.with_detail(contract_detail(contract)),
),
Verdict::DeadlineExpired { stage } => Step::Expired {
stage: stage.into(),
detail: None,
},
}
}
fn completed_envelope(
change_id: &str,
instance_id: &str,
contract: &OwnerExecutionContract,
evidence: String,
) -> ResultEnvelope {
let mut detail = contract_detail(contract);
if let Some(object) = detail.as_object_mut() {
object.insert(
"evidence".to_string(),
serde_json::Value::String(evidence.clone()),
);
}
ResultEnvelope::new(Operation::Wait, Outcome::Completed)
.with_change(change_id)
.with_instance(Some(instance_id.to_string()))
.with_message(evidence)
.with_detail(detail)
}
fn unknown_change_envelope(
change_id: &str,
instance_id: &str,
evidence_detail: String,
) -> ResultEnvelope {
ResultEnvelope::new(Operation::Wait, Outcome::ChangeNotFound)
.with_change(change_id)
.with_instance(Some(instance_id.to_string()))
.with_message(format!(
"the owner does not track a proposal named '{change_id}', and repository evidence does \
not prove one finished"
))
.with_detail(serde_json::json!({
"commands_submitted": 0,
"evidence_detail": evidence_detail,
}))
}
fn requires_action_envelope(
change_id: &str,
instance_id: &str,
observed_status: &str,
error_detail: Option<String>,
blocker: Option<&ChangeBlocker>,
message: String,
) -> ResultEnvelope {
let mut detail = serde_json::json!({
"observed_status": observed_status,
"commands_submitted": 0,
});
if let (Some(object), Some(error_detail)) = (detail.as_object_mut(), error_detail) {
object.insert(
"error_detail".to_string(),
serde_json::Value::String(error_detail),
);
}
if let (Some(object), Some(blocker)) = (detail.as_object_mut(), blocker) {
if let Ok(blocker) = serde_json::to_value(blocker) {
object.insert("blocker".to_string(), blocker);
}
}
ResultEnvelope::new(Operation::Wait, Outcome::ChangeRequiresAction)
.with_change(change_id)
.with_instance(Some(instance_id.to_string()))
.with_message(message)
.with_detail(detail)
}
fn contract_detail(contract: &OwnerExecutionContract) -> serde_json::Value {
serde_json::json!({
"terminal_mode": contract.terminal_mode,
"base_branch": contract.base_branch,
"remote": contract.remote,
"pushed_branch": contract.pushed_branch,
"commands_submitted": 0,
})
}
fn budget_phrase(timeout: Option<Duration>) -> String {
match timeout {
Some(timeout) => format!(" within {}ms", timeout.as_millis()),
None => String::new(),
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum TimeoutStage {
InitialObservation,
ObservingOwner,
RepositoryCertification,
RemoteVerification,
}
impl TimeoutStage {
fn as_str(self) -> &'static str {
match self {
Self::InitialObservation => "initial_observation",
Self::ObservingOwner => "observing_owner",
Self::RepositoryCertification => "repository_certification",
Self::RemoteVerification => "remote_verification",
}
}
}
impl From<CertificationStage> for TimeoutStage {
fn from(stage: CertificationStage) -> Self {
match stage {
CertificationStage::Repository => Self::RepositoryCertification,
CertificationStage::Remote => Self::RemoteVerification,
}
}
}
struct Diagnostics {
timeout: Option<Duration>,
started: Instant,
last_observation: Option<serde_json::Value>,
}
impl Diagnostics {
fn new(timeout: Option<Duration>) -> Self {
Self {
timeout,
started: Instant::now(),
last_observation: None,
}
}
fn record(&mut self, observation: &Observation, change_id: &str) {
let change = observation.change(change_id);
let execution = observation
.execution
.changes
.iter()
.find(|status| status.id == change_id);
self.last_observation = Some(serde_json::json!({
"observed_at": observation.execution.observed_at,
"state_revision": observation.state_revision,
"event_sequence": observation.event_sequence,
"change": change,
"execution": execution,
}));
}
fn expire(
&self,
change_id: &str,
instance_id: Option<&str>,
stage: TimeoutStage,
detail: Option<String>,
) -> ResultEnvelope {
let mut message = match stage {
TimeoutStage::InitialObservation => format!(
"the owner did not answer the first observation{}",
budget_phrase(self.timeout)
),
_ => format!(
"no verified terminal outcome for '{change_id}'{}",
budget_phrase(self.timeout)
),
};
if let Some(detail) = detail {
message.push_str(&format!("; {detail}"));
}
let mut body = serde_json::json!({
"commands_submitted": 0,
"timeout_stage": stage.as_str(),
"wait_elapsed_ms": self.started.elapsed().as_millis() as u64,
"last_observation": self.last_observation,
});
if let (Some(object), Some(timeout)) = (body.as_object_mut(), self.timeout) {
object.insert(
"timeout_ms".to_string(),
serde_json::json!(timeout.as_millis() as u64),
);
}
ResultEnvelope::new(Operation::Wait, Outcome::Timeout)
.with_change(change_id)
.with_instance(instance_id.map(str::to_string))
.with_message(message)
.with_detail(body)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::web::remote_control_api::dto::{
AttentionState, ChangeActivity, ChangeExecutionState, ChangeExecutionStatus,
ChangeResource, ChangeTiming, ExecutionPhase, LatestLogProjection, ParallelEligibility,
QueueIntent, TerminalMode,
};
fn change(id: &str, display_status: &str) -> ChangeResource {
ChangeResource {
id: id.to_string(),
display_status: display_status.to_string(),
progress_status: "in_progress".to_string(),
completed_tasks: 3,
total_tasks: 7,
progress_percent: 42.0,
dependencies: Vec::new(),
iteration_number: Some(2),
execution_marked: true,
queue_intent: QueueIntent::Queued,
attention: AttentionState::None,
blocker: None,
error_detail: None,
actions: crate::web::remote_control_api::projection::change_actions_for_test(
"select",
display_status,
None,
),
parallel: ParallelEligibility::default(),
timing: ChangeTiming::default(),
latest_activity: None,
worktree: None,
}
}
fn execution(id: &str, phase: ExecutionPhase) -> ChangeExecutionStatus {
ChangeExecutionStatus {
id: id.to_string(),
execution_id: Some(format!("x-{id}")),
execution_state: ChangeExecutionState::Active,
current_phase: phase,
last_completed_phase: Some(ExecutionPhase::Apply),
iteration: Some(2),
phase_started_at: Some("2026-01-01T00:00:00Z".to_string()),
last_completed_at: Some("2026-01-01T00:00:01Z".to_string()),
run_started_at: Some("2026-01-01T00:00:00Z".to_string()),
run_completed_at: None,
latest_activity: Some(ChangeActivity {
event_type: "phase_started".to_string(),
timestamp: "2026-01-01T00:00:00Z".to_string(),
detail: Some("acceptance started".to_string()),
}),
latest_log: Some(LatestLogProjection {
message: "running acceptance".to_string(),
level: crate::events::LogLevel::Info,
operation: Some("acceptance".to_string()),
iteration: Some(2),
created_at: "2026-01-01T00:00:02Z".to_string(),
}),
}
}
fn observation(
changes: Vec<ChangeResource>,
executions: Vec<ChangeExecutionStatus>,
revision: u64,
) -> Observation {
let mut observation = crate::client::session::observation_for_test(changes);
observation.state_revision = revision;
observation.event_sequence = revision * 10;
observation.execution.changes = executions;
observation.execution.observed_at = format!("2026-01-01T00:00:0{revision}Z");
observation
}
fn contract() -> OwnerExecutionContract {
OwnerExecutionContract {
base_branch: "main".to_string(),
terminal_mode: TerminalMode::Merged,
remote: None,
pushed_branch: None,
}
}
fn external_blocker() -> ChangeBlocker {
ChangeBlocker {
status: "blocked".to_string(),
kind: BlockerKind::External,
category: Some("pending_verification".to_string()),
detail: Some("waiting on the signing certificate".to_string()),
unblock_condition: Some("the certificate is issued".to_string()),
prerequisite_owner: Some("release".to_string()),
origin: Some("apply".to_string()),
resumable: true,
dependencies: Vec::new(),
}
}
#[test]
fn a_completion_envelope_carries_its_evidence_and_no_command_count() {
let envelope = completed_envelope("alpha", "i-1", &contract(), "proof".to_string());
assert_eq!(envelope.outcome, Outcome::Completed);
assert!(envelope.ok);
assert_eq!(envelope.detail["evidence"], "proof");
assert_eq!(envelope.detail["terminal_mode"], "merged");
assert_eq!(envelope.detail["commands_submitted"], 0);
}
#[test]
fn a_timeout_reports_the_budget_and_the_missing_evidence() {
let diagnostics = Diagnostics::new(Some(Duration::from_secs(90)));
let envelope = diagnostics.expire(
"alpha",
Some("i-1"),
TimeoutStage::ObservingOwner,
Some("no archive entry".to_string()),
);
assert_eq!(envelope.outcome, Outcome::Timeout);
assert!(!envelope.ok);
let message = envelope.message.unwrap();
assert!(message.contains("within 90000ms"), "{message}");
assert!(message.contains("no archive entry"), "{message}");
assert_eq!(envelope.detail["commands_submitted"], 0);
assert_eq!(envelope.detail["timeout_ms"], 90_000);
assert_eq!(envelope.detail["timeout_stage"], "observing_owner");
assert!(envelope.detail["wait_elapsed_ms"].is_u64());
}
#[test]
fn every_timeout_stage_has_its_own_stable_wire_spelling() {
assert_eq!(
TimeoutStage::InitialObservation.as_str(),
"initial_observation"
);
assert_eq!(TimeoutStage::ObservingOwner.as_str(), "observing_owner");
assert_eq!(
TimeoutStage::RepositoryCertification.as_str(),
"repository_certification"
);
assert_eq!(
TimeoutStage::RemoteVerification.as_str(),
"remote_verification"
);
assert_eq!(
TimeoutStage::from(CertificationStage::Repository),
TimeoutStage::RepositoryCertification
);
assert_eq!(
TimeoutStage::from(CertificationStage::Remote),
TimeoutStage::RemoteVerification
);
}
#[test]
fn a_timeout_before_the_first_observation_invents_no_owner_or_state() {
let diagnostics = Diagnostics::new(Some(Duration::from_millis(500)));
let envelope = diagnostics.expire("alpha", None, TimeoutStage::InitialObservation, None);
assert_eq!(envelope.outcome, Outcome::Timeout);
assert_eq!(envelope.instance_id, None);
assert_eq!(envelope.detail["timeout_stage"], "initial_observation");
assert!(envelope.detail["last_observation"].is_null());
assert_eq!(envelope.detail["timeout_ms"], 500);
assert_eq!(envelope.detail["commands_submitted"], 0);
}
#[test]
fn a_timeout_after_an_observation_reports_the_target_row_and_its_execution() {
let mut diagnostics = Diagnostics::new(Some(Duration::from_secs(30)));
diagnostics.record(
&observation(
vec![change("alpha", "accepting")],
vec![execution("alpha", ExecutionPhase::Acceptance)],
4,
),
"alpha",
);
let envelope = diagnostics.expire("alpha", Some("i-1"), TimeoutStage::ObservingOwner, None);
let last = &envelope.detail["last_observation"];
assert_eq!(last["state_revision"], 4);
assert_eq!(last["event_sequence"], 40);
assert_eq!(last["observed_at"], "2026-01-01T00:00:04Z");
assert_eq!(last["change"]["id"], "alpha");
assert_eq!(last["change"]["display_status"], "accepting");
assert_eq!(last["change"]["completed_tasks"], 3);
assert_eq!(last["change"]["total_tasks"], 7);
assert_eq!(last["execution"]["execution_id"], "x-alpha");
assert_eq!(last["execution"]["current_phase"], "acceptance");
assert_eq!(last["execution"]["last_completed_phase"], "apply");
assert_eq!(last["execution"]["execution_state"], "active");
assert_eq!(last["execution"]["run_started_at"], "2026-01-01T00:00:00Z");
assert_eq!(
last["execution"]["latest_activity"]["event_type"],
"phase_started"
);
assert_eq!(
last["execution"]["latest_log"]["message"],
"running acceptance"
);
}
#[test]
fn a_retained_observation_carries_no_change_but_the_requested_one() {
let mut diagnostics = Diagnostics::new(Some(Duration::from_secs(30)));
diagnostics.record(
&observation(
vec![change("alpha", "applying"), change("beta", "archiving")],
vec![
execution("alpha", ExecutionPhase::Apply),
execution("beta", ExecutionPhase::Archive),
],
2,
),
"alpha",
);
let envelope = diagnostics.expire("alpha", Some("i-1"), TimeoutStage::ObservingOwner, None);
let rendered = envelope.detail["last_observation"].to_string();
assert!(rendered.contains("alpha"), "{rendered}");
assert!(!rendered.contains("beta"), "{rendered}");
assert!(!rendered.contains("archiving"), "{rendered}");
}
#[test]
fn a_newer_coherent_observation_replaces_the_retained_one() {
let mut diagnostics = Diagnostics::new(Some(Duration::from_secs(30)));
diagnostics.record(
&observation(
vec![change("alpha", "applying")],
vec![execution("alpha", ExecutionPhase::Apply)],
1,
),
"alpha",
);
diagnostics.record(
&observation(
vec![change("alpha", "accepting")],
vec![execution("alpha", ExecutionPhase::Acceptance)],
5,
),
"alpha",
);
let envelope = diagnostics.expire("alpha", Some("i-1"), TimeoutStage::ObservingOwner, None);
let last = &envelope.detail["last_observation"];
assert_eq!(last["state_revision"], 5);
assert_eq!(last["change"]["display_status"], "accepting");
assert_eq!(last["execution"]["current_phase"], "acceptance");
}
#[test]
fn a_retained_observation_reports_an_absent_target_rather_than_a_stale_row() {
let mut diagnostics = Diagnostics::new(Some(Duration::from_secs(30)));
diagnostics.record(
&observation(
vec![change("alpha", "applying")],
vec![execution("alpha", ExecutionPhase::Apply)],
1,
),
"alpha",
);
diagnostics.record(&observation(Vec::new(), Vec::new(), 2), "alpha");
let envelope = diagnostics.expire("alpha", Some("i-1"), TimeoutStage::ObservingOwner, None);
let last = &envelope.detail["last_observation"];
assert_eq!(last["state_revision"], 2);
assert!(last["change"].is_null());
assert!(last["execution"].is_null());
}
#[test]
fn a_certification_timeout_reports_its_stage_over_the_observation_it_started_from() {
let mut diagnostics = Diagnostics::new(Some(Duration::from_secs(30)));
diagnostics.record(
&observation(
vec![change("alpha", "merged")],
vec![execution("alpha", ExecutionPhase::Merge)],
9,
),
"alpha",
);
for (stage, expected) in [
(CertificationStage::Repository, "repository_certification"),
(CertificationStage::Remote, "remote_verification"),
] {
let envelope = diagnostics.expire("alpha", Some("i-1"), stage.into(), None);
assert_eq!(envelope.detail["timeout_stage"], expected);
assert_eq!(envelope.detail["last_observation"]["state_revision"], 9);
assert_eq!(
envelope.detail["last_observation"]["change"]["display_status"],
"merged"
);
assert_eq!(envelope.detail["commands_submitted"], 0);
}
}
#[test]
fn a_manual_action_release_names_the_status_and_submits_nothing() {
let envelope = requires_action_envelope(
"alpha",
"i-1",
"merge wait",
Some("conflict in src/lib.rs".to_string()),
None,
"needs an operator".to_string(),
);
assert_eq!(envelope.outcome, Outcome::ChangeRequiresAction);
assert!(!envelope.ok);
assert_eq!(envelope.exit_code(), 27);
assert_eq!(envelope.change_id.as_deref(), Some("alpha"));
assert_eq!(envelope.detail["observed_status"], "merge wait");
assert_eq!(envelope.detail["error_detail"], "conflict in src/lib.rs");
assert_eq!(envelope.detail["commands_submitted"], 0);
assert!(!envelope.detail.as_object().unwrap().contains_key("blocker"));
}
#[test]
fn an_external_blocker_release_publishes_the_prerequisite_facts() {
let envelope = requires_action_envelope(
"alpha",
"i-1",
"blocked",
Some("waiting on the signing certificate".to_string()),
Some(&external_blocker()),
"blocked externally".to_string(),
);
assert_eq!(envelope.outcome, Outcome::ChangeRequiresAction);
assert_eq!(envelope.exit_code(), 27);
assert_eq!(envelope.detail["observed_status"], "blocked");
assert_eq!(envelope.detail["commands_submitted"], 0);
assert_eq!(
envelope.detail["error_detail"],
"waiting on the signing certificate"
);
assert_eq!(envelope.detail["blocker"]["kind"], "external");
assert_eq!(
envelope.detail["blocker"]["category"],
"pending_verification"
);
assert_eq!(
envelope.detail["blocker"]["unblock_condition"],
"the certificate is issued"
);
assert_eq!(envelope.detail["blocker"]["prerequisite_owner"], "release");
assert_eq!(envelope.detail["blocker"]["resumable"], true);
}
#[test]
fn an_unknown_target_is_refused_with_the_shared_outcome_and_submits_nothing() {
let envelope =
unknown_change_envelope("aaaa", "i-1", "no archive entry on 'main'".to_string());
assert_eq!(envelope.outcome, Outcome::ChangeNotFound);
assert!(!envelope.ok);
assert_eq!(envelope.exit_code(), 9);
assert_eq!(envelope.change_id.as_deref(), Some("aaaa"));
assert_eq!(envelope.instance_id.as_deref(), Some("i-1"));
assert_eq!(envelope.detail["commands_submitted"], 0);
assert_eq!(
envelope.detail["evidence_detail"],
"no archive entry on 'main'"
);
}
#[test]
fn a_manual_action_release_omits_error_detail_when_the_owner_published_none() {
let envelope =
requires_action_envelope("alpha", "i-1", "stopped", None, None, "stopped".to_string());
let detail = envelope.detail.as_object().unwrap();
assert!(!detail.contains_key("error_detail"));
assert_eq!(detail["observed_status"], "stopped");
assert_eq!(detail["commands_submitted"], 0);
}
#[test]
fn a_settled_success_claim_gets_exactly_one_second_look() {
assert_eq!(UNCERTIFIED_CLAIM_ROUNDS, 2);
}
#[test]
fn an_unbounded_wait_never_reaches_an_operation_deadline() {
assert!(!reached(None));
assert!(reached(Some(Instant::now() - Duration::from_secs(1))));
assert!(!reached(Some(Instant::now() + Duration::from_secs(60))));
}
#[test]
fn git_children_are_bounded_per_child_exactly_when_the_operation_is_not() {
assert!(matches!(
git_deadline(None),
GitDeadline::PerChild(GIT_INVOCATION_BUDGET)
));
assert!(git_deadline(Some(Instant::now())).is_operation_deadline());
}
#[tokio::test]
async fn a_step_without_a_deadline_is_run_rather_than_bounded() {
assert_eq!(within(None, async { 7 }).await, Some(7));
assert_eq!(
within(
Some(Instant::now() - Duration::from_secs(1)),
std::future::pending::<i32>()
)
.await,
None
);
}
}