klieo-ops 0.3.0

Operational layer above klieo-core: supervisor, governor, gates, escalation, worklog, handoff.
Documentation
//! `KvHandoff`: default `Handoff` impl persisting envelopes in `KvStore`.
//!
//! Envelopes are stored content-addressed under
//! `ops.handoff.envelopes/<content_hash>`. Delivery pointers land under
//! `ops.handoff.delivery/<agent_id>:<content_hash>`. Verification performs
//! expiry + signature-math checks.

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";

/// Optional audit sink for `KvHandoff`.
struct AuditSink {
    episodic: Arc<dyn EpisodicMemory>,
    run_id: RunId,
}

/// KV-backed `Handoff` implementation.
pub struct KvHandoff {
    kv: Arc<dyn KvStore>,
    audit: Option<AuditSink>,
}

impl KvHandoff {
    /// Construct a `KvHandoff` over the given `KvStore`.
    #[must_use]
    pub fn new(kv: Arc<dyn KvStore>) -> Self {
        Self { kv, audit: None }
    }

    /// Wire an audit sink. When set, every state-changing call emits the
    /// corresponding `OpsEvent` into `episodic` under `run_id`.
    #[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();

        // Build a stub envelope to derive the canonical payload for signing.
        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));
        }
        // Phase B M9 scope: source-identity trust is delegated to the caller.
        // Receiver crates wire their own trust registry. `KvHandoff::verify`
        // is intentionally trust-registry-agnostic at this milestone.
        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(),
        })
    }
}