use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use crate::{DurableControlReceipt, ManagedRunTarget, RunAdmissionLease};
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum DurableRunControlEffect {
Steer {
text: String,
},
Interrupt {
#[serde(default, skip_serializing_if = "Option::is_none")]
reason: Option<String>,
},
}
impl DurableRunControlEffect {
#[must_use]
pub const fn operation(&self) -> &'static str {
match self {
Self::Steer { .. } => "steer",
Self::Interrupt { .. } => "interrupt",
}
}
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum DurableRunControlStatus {
Pending,
Delivered,
Consumed,
Reconciled,
}
impl DurableRunControlStatus {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Pending => "pending",
Self::Delivered => "delivered",
Self::Consumed => "consumed",
Self::Reconciled => "reconciled",
}
}
#[must_use]
pub const fn can_advance_to(self, next: Self) -> bool {
self as u8 == next as u8
|| matches!(
(self, next),
(Self::Pending, Self::Delivered | Self::Reconciled)
| (Self::Delivered, Self::Consumed | Self::Reconciled)
| (Self::Consumed, Self::Reconciled)
)
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct AdmitRunControl {
pub lease: RunAdmissionLease,
pub authority_binding: String,
pub operation_id: String,
pub receipt_id: String,
pub idempotency_key: String,
pub command_fingerprint: String,
pub effect: DurableRunControlEffect,
pub created_at: DateTime<Utc>,
}
impl AdmitRunControl {
#[must_use]
pub fn into_intent(self) -> DurableRunControlIntent {
let receipt = DurableControlReceipt {
receipt_id: self.receipt_id,
target: self.lease.target.clone(),
operation_id: self.operation_id.clone(),
operation: self.effect.operation().to_string(),
idempotency_key: self.idempotency_key,
command_fingerprint: self.command_fingerprint,
fencing_generation: self.lease.fencing_generation,
state: DurableRunControlStatus::Pending.as_str().to_string(),
created_at: self.created_at,
};
DurableRunControlIntent {
operation_id: self.operation_id,
target: self.lease.target,
authority_binding: self.authority_binding,
admission_id: self.lease.admission_id,
host_instance_id: self.lease.host_instance_id,
fencing_generation: self.lease.fencing_generation,
idempotency_key: receipt.idempotency_key.clone(),
command_fingerprint: receipt.command_fingerprint.clone(),
receipt,
effect: self.effect,
status: DurableRunControlStatus::Pending,
created_at: self.created_at,
delivered_at: None,
consumed_at: None,
reconciled_at: None,
}
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
pub struct DurableRunControlIntent {
pub operation_id: String,
pub target: ManagedRunTarget,
pub authority_binding: String,
pub admission_id: String,
pub host_instance_id: String,
pub fencing_generation: u64,
pub idempotency_key: String,
pub command_fingerprint: String,
pub receipt: DurableControlReceipt,
pub effect: DurableRunControlEffect,
pub status: DurableRunControlStatus,
pub created_at: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub delivered_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub consumed_at: Option<DateTime<Utc>>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub reconciled_at: Option<DateTime<Utc>>,
}
impl starweaver_core::VersionedRecord for DurableRunControlIntent {
const SCHEMA: &'static str = "starweaver.session.durable_run_control_intent";
}
impl DurableRunControlIntent {
#[must_use]
pub fn matches_admission(&self, request: &AdmitRunControl) -> bool {
self.target == request.lease.target
&& self.authority_binding == request.authority_binding
&& self.admission_id == request.lease.admission_id
&& self.host_instance_id == request.lease.host_instance_id
&& self.fencing_generation == request.lease.fencing_generation
&& self.operation_id == request.operation_id
&& self.receipt.receipt_id == request.receipt_id
&& self.idempotency_key == request.idempotency_key
&& self.command_fingerprint == request.command_fingerprint
&& self.effect == request.effect
}
pub fn advance(
&mut self,
next: DurableRunControlStatus,
occurred_at: DateTime<Utc>,
) -> Result<(), &'static str> {
if !self.status.can_advance_to(next) {
return Err("invalid durable run control state transition");
}
if self.status == next {
return Ok(());
}
self.status = next;
self.receipt.state = next.as_str().to_string();
match next {
DurableRunControlStatus::Pending => {}
DurableRunControlStatus::Delivered => self.delivered_at = Some(occurred_at),
DurableRunControlStatus::Consumed => {
self.delivered_at.get_or_insert(occurred_at);
self.consumed_at = Some(occurred_at);
}
DurableRunControlStatus::Reconciled => self.reconciled_at = Some(occurred_at),
}
Ok(())
}
}
#[must_use]
pub fn deterministic_run_control_operation_id(
operation: &str,
authority_binding: &str,
target: &ManagedRunTarget,
idempotency_key: &str,
command_fingerprint: &str,
) -> String {
let mut digest = Sha256::new();
for component in [
"starweaver.session.run_control.operation.v1",
operation,
authority_binding,
target.namespace_id.as_str(),
target.session_id.as_str(),
target.run_id.as_str(),
idempotency_key,
command_fingerprint,
] {
digest.update(component.len().to_be_bytes());
digest.update(component.as_bytes());
}
format!("run_control_{:x}", digest.finalize())
}
#[must_use]
pub fn deterministic_run_control_receipt_id(operation_id: &str) -> String {
let mut digest = Sha256::new();
digest.update(b"starweaver.session.run_control.receipt.v1\0");
digest.update(operation_id.as_bytes());
format!("control_{:x}", digest.finalize())
}
#[cfg(test)]
mod tests {
use super::*;
use starweaver_core::{RunId, SessionId};
#[test]
fn operation_identity_binds_authority_target_key_and_fingerprint() {
let target = ManagedRunTarget::new(
"local",
SessionId::from_string("session-a"),
RunId::from_string("run-a"),
);
let first = deterministic_run_control_operation_id(
"steer",
"authority-a",
&target,
"key-a",
"sha256:a",
);
assert_eq!(
first,
deterministic_run_control_operation_id(
"steer",
"authority-a",
&target,
"key-a",
"sha256:a"
)
);
assert_ne!(
first,
deterministic_run_control_operation_id(
"steer",
"authority-b",
&target,
"key-a",
"sha256:a"
)
);
assert_ne!(
first,
deterministic_run_control_operation_id(
"interrupt",
"authority-a",
&target,
"key-a",
"sha256:a"
)
);
}
#[test]
fn control_states_are_monotonic() {
assert!(
DurableRunControlStatus::Pending.can_advance_to(DurableRunControlStatus::Delivered)
);
assert!(
DurableRunControlStatus::Delivered.can_advance_to(DurableRunControlStatus::Consumed)
);
assert!(
DurableRunControlStatus::Consumed.can_advance_to(DurableRunControlStatus::Reconciled)
);
assert!(
!DurableRunControlStatus::Consumed.can_advance_to(DurableRunControlStatus::Pending)
);
}
}