//! Local implementation of PeerHost using Store directly. //! //! Enables standalone peer session persistence, idempotent request admission, //! correlated results, cancellation, and lifecycle management within klbr-runtime. use crate::protocol::*; use crate::store::{self, PeerAdmitOutcome, PeerRequestRecord, PeerSessionRecord, Store}; use anyhow::{bail, ensure}; use serde_json::{json, Value}; #[derive(Clone, Default)] pub struct LocalPeerHost; impl PeerHost for LocalPeerHost { fn open_peer<'a>( &'a self, store: &'a Store, caller_session_id: &'a SessionId, req: PeerOpenReq, ) -> HostFuture<'a> { Box::pin(async move { let peer_id_str = req .peer_id .unwrap_or_else(|| format!("peer-{}", uuid::Uuid::new_v4())); let peer_sid = SessionId::try_from(peer_id_str.clone()) .map_err(|e| anyhow::anyhow!("invalid peer_id: {e}"))?; let mode = req.mode.as_deref().unwrap_or("independent"); let now = store::now_ms()?; if mode == "resume" { let existing = store.session_state(&peer_sid).await?; ensure!( existing.is_some(), "cannot resume non-existent peer session '{}'", peer_id_str ); } else { store.ensure_session(&peer_sid).await?; } store.ensure_session(caller_session_id).await?; let peer_rec = PeerSessionRecord { peer_id: peer_id_str.clone(), owner_session_id: caller_session_id.to_string(), init_mode: mode.to_string(), profile: req.profile.clone(), workspace: req.workspace.clone(), kernel_target: req.kernel_target.as_ref().map(|v| v.to_string()), budget: req.budget.as_ref().map(|v| v.to_string()), branch_parent_revision: req.branch_parent_revision.clone(), created_ms: now, }; store.record_peer_session(peer_rec).await?; Ok(json!({ "peer_id": peer_id_str, "status": "opened", })) }) } fn get_peer<'a>( &'a self, store: &'a Store, _caller_session_id: &'a SessionId, req: PeerGetReq, ) -> HostFuture<'a> { Box::pin(async move { let peer_sid = SessionId::try_from(req.peer_id.clone()) .map_err(|e| anyhow::anyhow!("invalid peer_id: {e}"))?; let meta = store.get_peer_session(&req.peer_id).await?; let state = store .session_state(&peer_sid) .await? .map(|(st, _)| st) .unwrap_or_else(|| "ready".into()); Ok(json!({ "peer_id": req.peer_id, "status": state, "profile": meta.as_ref().and_then(|m| m.profile.clone()), "workspace": meta.as_ref().and_then(|m| m.workspace.clone()), "init_mode": meta.as_ref().map(|m| m.init_mode.clone()), })) }) } fn list_peers<'a>( &'a self, store: &'a Store, caller_session_id: &'a SessionId, ) -> HostFuture<'a> { Box::pin(async move { let peers = store.list_peer_sessions(caller_session_id.as_str()).await?; let peers_json: Vec = peers .into_iter() .map(|p| { json!({ "peer_id": p.peer_id, "owner_session_id": p.owner_session_id, "profile": p.profile, "workspace": p.workspace, "init_mode": p.init_mode, "created_ms": p.created_ms, }) }) .collect(); Ok(json!({ "peers": peers_json, })) }) } fn send_peer_message<'a>( &'a self, store: &'a Store, caller_session_id: &'a SessionId, req: PeerSendReq, ) -> HostFuture<'a> { Box::pin(async move { let peer_sid = SessionId::try_from(req.peer_id.clone()) .map_err(|e| anyhow::anyhow!("invalid peer_id: {e}"))?; store.ensure_session(&peer_sid).await?; store.ensure_session(caller_session_id).await?; let now = store::now_ms()?; let record = PeerRequestRecord { request_id: req.request_id.clone(), sender_session_id: caller_session_id.to_string(), peer_id: req.peer_id.clone(), text: req.content.clone(), evidence_refs: req.evidence.clone(), status: "admitted".into(), admitted_at_ms: now, completed_at_ms: None, result: None, error: None, }; let outcome = store.admit_peer_request(record).await?; let (admitted_rec, duplicate) = match outcome { PeerAdmitOutcome::Admitted(rec) => (rec, false), PeerAdmitOutcome::AlreadyAdmitted(rec) => (rec, true), }; // Admission receipt returned immediately (admission != completion) // Non-generative acknowledgment ensures no echo loops or unbounded inference Ok(json!({ "request_id": admitted_rec.request_id, "peer_id": admitted_rec.peer_id, "status": "admitted", "admitted_at_ms": admitted_rec.admitted_at_ms, "duplicate": duplicate, })) }) } fn get_peer_result<'a>( &'a self, store: &'a Store, _caller_session_id: &'a SessionId, req: PeerResultReq, ) -> HostFuture<'a> { Box::pin(async move { let req_id = &req.request_id; let deadline = req .wait_ms .map(|ms| tokio::time::Instant::now() + std::time::Duration::from_millis(ms)); loop { let rec = store.get_peer_request(req_id).await?; match rec { None => bail!("peer_request_not_found: request '{}' not found", req_id), Some(r) => { if r.status == "completed" || r.status == "cancelled" || r.status == "failed" { return Ok(json!({ "request_id": r.request_id, "peer_id": r.peer_id, "status": r.status, "result": r.result, "error": r.error, "completed_at_ms": r.completed_at_ms, })); } if deadline.is_none_or(|deadline| tokio::time::Instant::now() >= deadline) { return Ok(json!({ "request_id": r.request_id, "peer_id": r.peer_id, "status": "pending", })); } tokio::time::sleep(std::time::Duration::from_millis(20)).await; } } } }) } fn submit_peer_result<'a>( &'a self, store: &'a Store, _caller_session_id: &'a SessionId, req: PeerSubmitResultReq, ) -> HostFuture<'a> { Box::pin(async move { let now = store::now_ms()?; let updated = store .complete_peer_request(&req.request_id, req.result.clone(), now) .await?; if !updated { let rec = store.get_peer_request(&req.request_id).await?; if let Some(r) = rec { if r.status == "completed" { return Ok(json!({ "request_id": req.request_id, "status": "already_completed", "result": r.result, })); } if r.status == "cancelled" { bail!( "request_cancelled: cannot submit result for cancelled request '{}'", req.request_id ); } } bail!("request_not_found: request '{}' not found", req.request_id); } Ok(json!({ "request_id": req.request_id, "status": "recorded", })) }) } fn cancel_peer_request<'a>( &'a self, store: &'a Store, _caller_session_id: &'a SessionId, req: PeerCancelReq, ) -> HostFuture<'a> { Box::pin(async move { let cancelled = store.cancel_peer_request(&req.request_id).await?; Ok(json!({ "request_id": req.request_id, "status": "cancelled", "cancelled": cancelled, })) }) } fn stop_peer<'a>( &'a self, store: &'a Store, _caller_session_id: &'a SessionId, req: PeerStopReq, ) -> HostFuture<'a> { Box::pin(async move { let peer_sid = SessionId::try_from(req.peer_id.clone()) .map_err(|e| anyhow::anyhow!("invalid peer_id: {e}"))?; let drain = req.drain.unwrap_or(false); let state_str = if drain { "draining" } else { "stopped" }; store.set_session_state(&peer_sid, state_str).await?; Ok(json!({ "peer_id": req.peer_id, "status": "stopped", })) }) } }