use std::marker::PhantomData;
use std::sync::Arc;
use serde::Serialize;
use zeroize::Zeroizing;
use crate::backend::DurableBackendEnum;
use crate::backend::local::now_unix_millis;
use crate::error::DurableError;
use crate::ids::{ExecutionId, PromiseId};
pub(crate) const RESOLVER_TOKEN_LEN: usize = 32;
const RESOLVER_CONTEXT: &str = "zeph-durable v1 promise resolver-token 2026";
#[derive(Debug, Clone)]
pub(crate) struct PromiseRecord {
pub(crate) execution_id: ExecutionId,
pub(crate) resolver_token_hash: [u8; 32],
pub(crate) resolved: bool,
pub(crate) payload: Option<Vec<u8>>,
}
pub(crate) fn resolver_token_hash(
promise_id: PromiseId,
execution_id: ExecutionId,
token: &[u8; RESOLVER_TOKEN_LEN],
) -> blake3::Hash {
let mut input = [0u8; 16 + 16 + RESOLVER_TOKEN_LEN];
input[..16].copy_from_slice(promise_id.as_uuid().as_bytes());
input[16..32].copy_from_slice(execution_id.as_bytes());
input[32..].copy_from_slice(token);
blake3::Hash::from(blake3::derive_key(RESOLVER_CONTEXT, &input))
}
#[derive(Debug)]
pub struct DurablePromise<T> {
id: PromiseId,
resolver_token: Option<Zeroizing<[u8; RESOLVER_TOKEN_LEN]>>,
_t: PhantomData<fn() -> T>,
}
impl<T> DurablePromise<T> {
pub(crate) fn fresh(
id: PromiseId,
resolver_token: Zeroizing<[u8; RESOLVER_TOKEN_LEN]>,
) -> Self {
Self {
id,
resolver_token: Some(resolver_token),
_t: PhantomData,
}
}
pub(crate) fn resumed(id: PromiseId) -> Self {
Self {
id,
resolver_token: None,
_t: PhantomData,
}
}
#[must_use]
pub fn id(&self) -> PromiseId {
self.id
}
#[must_use]
pub fn resolver_token(&self) -> Option<&[u8; RESOLVER_TOKEN_LEN]> {
self.resolver_token.as_deref()
}
#[must_use]
pub fn is_resumed(&self) -> bool {
self.resolver_token.is_none()
}
}
#[derive(Clone, Debug)]
pub struct DurableHandle {
backend: Arc<DurableBackendEnum>,
}
impl DurableHandle {
#[must_use]
pub fn new(backend: Arc<DurableBackendEnum>) -> Self {
Self { backend }
}
#[tracing::instrument(
name = "durable.promise.resolve",
skip(self, resolver_token, value),
fields(promise_id = %id.as_uuid())
)]
pub async fn resolve<T: Serialize>(
&self,
id: PromiseId,
resolver_token: &[u8; RESOLVER_TOKEN_LEN],
value: T,
) -> Result<(), DurableError> {
let record = self
.backend
.promise_state(id)
.await?
.ok_or(DurableError::UnknownPromise)?;
let presented = resolver_token_hash(id, record.execution_id, resolver_token);
if presented != blake3::Hash::from(record.resolver_token_hash) {
return Err(DurableError::PromiseRejected);
}
if record.resolved {
return Ok(());
}
let payload = serde_json::to_vec(&value).map_err(|_| DurableError::Serialize {
step: "promise.resolve",
})?;
self.backend
.resolve_promise(id, record.execution_id, &payload, now_unix_millis())
.await?;
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ids::StepId;
#[test]
fn resolver_hash_binds_promise_and_execution() {
let exec = ExecutionId::new();
let promise = PromiseId::derive(exec, StepId::new(0));
let token = [7u8; RESOLVER_TOKEN_LEN];
let base = resolver_token_hash(promise, exec, &token);
assert_eq!(
base,
resolver_token_hash(promise, exec, &token),
"deterministic for fixed inputs"
);
let other_promise = PromiseId::derive(exec, StepId::new(1));
assert_ne!(base, resolver_token_hash(other_promise, exec, &token));
assert_ne!(
base,
resolver_token_hash(promise, ExecutionId::new(), &token)
);
assert_ne!(
base,
resolver_token_hash(promise, exec, &[8u8; RESOLVER_TOKEN_LEN])
);
}
#[test]
fn fresh_promise_carries_token_resumed_does_not() {
let fresh: DurablePromise<u32> =
DurablePromise::fresh(PromiseId::new(), Zeroizing::new([1u8; RESOLVER_TOKEN_LEN]));
assert!(fresh.resolver_token().is_some());
assert!(!fresh.is_resumed());
let resumed: DurablePromise<u32> = DurablePromise::resumed(PromiseId::new());
assert!(resumed.resolver_token().is_none());
assert!(resumed.is_resumed());
}
}