ive harnessed the harness
Something went wrong. Try again.
10 kB · 275 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276//! 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<Value> = 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", })) }) }}