use super::envelope::{
verify_signature, CanonicalPayload, EnvelopeSignature, HandoffEnvelope, HandoffProof,
MerklePrefixProof,
};
use super::trait_::{Handoff, HandoffError, HandoffState};
use crate::emit::emit_ops_event;
use crate::ops_event::OpsEvent;
use crate::types::AgentId;
use async_trait::async_trait;
use bytes::Bytes;
use chrono::Utc;
use ed25519_dalek::{Signer as _, SigningKey};
use klieo_core::ids::RunId;
use klieo_core::memory::EpisodicMemory;
use klieo_core::KvStore;
use std::sync::Arc;
const BUCKET_ENVELOPES: &str = "ops.handoff.envelopes";
const BUCKET_DELIVERY: &str = "ops.handoff.delivery";
struct AuditSink {
episodic: Arc<dyn EpisodicMemory>,
run_id: RunId,
}
pub struct KvHandoff {
kv: Arc<dyn KvStore>,
audit: Option<AuditSink>,
}
impl KvHandoff {
#[must_use]
pub fn new(kv: Arc<dyn KvStore>) -> Self {
Self { kv, audit: None }
}
#[must_use]
pub fn with_audit(mut self, episodic: Arc<dyn EpisodicMemory>, run_id: RunId) -> Self {
self.audit = Some(AuditSink { episodic, run_id });
self
}
async fn emit(&self, event: OpsEvent) {
if let Some(sink) = &self.audit {
if let Err(err) = emit_ops_event(&*sink.episodic, sink.run_id, event).await {
tracing::warn!(
target: "klieo.ops.handoff.audit",
error = %err,
"audit emit failed"
);
}
}
}
fn content_hash(env: &HandoffEnvelope) -> String {
let canonical_bytes = CanonicalPayload::from_envelope(env).to_bytes();
blake3::hash(&canonical_bytes).to_hex().to_string()
}
}
#[async_trait]
impl Handoff for KvHandoff {
async fn package(
&self,
state: HandoffState,
signer: &SigningKey,
ttl: chrono::Duration,
) -> Result<HandoffEnvelope, HandoffError> {
let now = Utc::now();
let expires = now + ttl;
let issued_at = now.to_rfc3339();
let expires_at = expires.to_rfc3339();
let stub = HandoffEnvelope {
from_run: state.from_run.clone(),
merkle_prefix_proof: MerklePrefixProof {
root: state.merkle_root.clone(),
at_seq: state.at_seq,
},
redacted_state: state.redacted_state.clone(),
source_signature: EnvelopeSignature {
signature_hex: String::new(),
verifying_key_hex: String::new(),
},
source_identity: state.source_identity.clone(),
tenant: state.tenant.clone(),
issued_at: issued_at.clone(),
expires_at: expires_at.clone(),
};
let canonical_bytes = CanonicalPayload::from_envelope(&stub).to_bytes();
let sig = signer.sign(&canonical_bytes);
let vk = signer.verifying_key();
let source_signature = EnvelopeSignature {
signature_hex: hex::encode(sig.to_bytes()),
verifying_key_hex: hex::encode(vk.to_bytes()),
};
let envelope = HandoffEnvelope {
source_signature,
..stub
};
self.emit(OpsEvent::HandoffPackaged {
tenant: envelope.tenant.clone(),
from_run: envelope.from_run.clone(),
merkle_root: envelope.merkle_prefix_proof.root.clone(),
})
.await;
Ok(envelope)
}
async fn deliver(
&self,
envelope: HandoffEnvelope,
to: AgentId,
) -> Result<String, HandoffError> {
let body = serde_json::to_vec(&envelope)
.map_err(|e| HandoffError::Serialisation(format!("envelope serialise: {e}")))?;
let hash = Self::content_hash(&envelope);
self.kv
.put(BUCKET_ENVELOPES, &hash, Bytes::from(body))
.await
.map_err(|e| HandoffError::Storage(e.to_string()))?;
let delivery_key = format!("{}:{}", to.0, hash);
self.kv
.put(
BUCKET_DELIVERY,
&delivery_key,
Bytes::from(envelope.from_run.clone()),
)
.await
.map_err(|e| HandoffError::Storage(e.to_string()))?;
self.emit(OpsEvent::HandoffDelivered {
tenant: envelope.tenant.clone(),
from_run: envelope.from_run.clone(),
to_agent: to.0.clone(),
})
.await;
Ok(envelope.from_run)
}
async fn verify(&self, envelope: &HandoffEnvelope) -> Result<HandoffProof, HandoffError> {
let expires = chrono::DateTime::parse_from_rfc3339(&envelope.expires_at)
.map_err(|e| HandoffError::Internal(format!("parse expires_at: {e}")))?;
if Utc::now() > expires {
let err = HandoffError::Expired {
expired_at: envelope.expires_at.clone(),
};
self.emit(OpsEvent::HandoffVerificationFailed {
tenant: envelope.tenant.clone(),
from_run: Some(envelope.from_run.clone()),
reason: err.to_string(),
})
.await;
return Err(err);
}
if let Err(msg) = verify_signature(envelope) {
self.emit(OpsEvent::HandoffVerificationFailed {
tenant: envelope.tenant.clone(),
from_run: Some(envelope.from_run.clone()),
reason: msg.clone(),
})
.await;
return Err(HandoffError::SignatureInvalid(msg));
}
self.emit(OpsEvent::HandoffVerified {
tenant: envelope.tenant.clone(),
from_run: envelope.from_run.clone(),
})
.await;
Ok(HandoffProof {
from_run: envelope.from_run.clone(),
source_identity: envelope.source_identity.clone(),
tenant: envelope.tenant.clone(),
redacted_state: envelope.redacted_state.clone(),
})
}
}