refactor: Podのメインループのリファクタリング
This commit is contained in:
+122
-31
@@ -3,15 +3,19 @@ use std::sync::Arc;
|
||||
|
||||
use llm_worker::WorkerError;
|
||||
use llm_worker::llm_client::client::LlmClient;
|
||||
use llm_worker::llm_client::types::{Item, Role};
|
||||
use session_store::Store;
|
||||
use tokio::sync::{broadcast, mpsc, oneshot};
|
||||
|
||||
use llm_worker::Item;
|
||||
use session_store::LogEntry;
|
||||
use session_store::session_log;
|
||||
|
||||
use crate::ipc::alerter::Alerter;
|
||||
use crate::ipc::notify_buffer::NotifyBuffer;
|
||||
use crate::ipc::server::SocketServer;
|
||||
use crate::pod::{Pod, PodError, PodRunResult};
|
||||
use crate::pod::{LogCommand, LogDrainHandle, Pod, PodError, PodRunResult};
|
||||
use crate::runtime::dir::RuntimeDir;
|
||||
use crate::session_log_sink::SessionLogSink;
|
||||
use crate::shared_state::PodSharedState;
|
||||
use crate::spawn::comm_tools::{
|
||||
list_pods_tool, read_pod_output_tool, send_to_pod_tool, stop_pod_tool,
|
||||
@@ -22,16 +26,6 @@ use protocol::{
|
||||
AlertLevel, AlertSource, ErrorCode, Event, Method, PodStatus, RunResult, Segment, TurnResult,
|
||||
};
|
||||
|
||||
fn is_system_message_item(item: &Item) -> bool {
|
||||
matches!(
|
||||
item,
|
||||
Item::Message {
|
||||
role: Role::System,
|
||||
..
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// PodHandle — client-facing, Clone-able
|
||||
// ---------------------------------------------------------------------------
|
||||
@@ -43,6 +37,10 @@ pub struct PodHandle {
|
||||
pub shared_state: Arc<PodSharedState>,
|
||||
pub runtime_dir: Arc<RuntimeDir>,
|
||||
pub alerter: Alerter,
|
||||
/// Session-log mirror + broadcast handle. The IPC server snapshots
|
||||
/// it on every new connection (Event::Snapshot) and forwards
|
||||
/// subsequent commits (Event::Entry) on the receiver.
|
||||
pub sink: SessionLogSink,
|
||||
}
|
||||
|
||||
impl PodHandle {
|
||||
@@ -86,11 +84,11 @@ async fn finish_controller_run<C, St>(
|
||||
C: LlmClient + Clone + 'static,
|
||||
St: Store + Clone + 'static,
|
||||
{
|
||||
let items = pod.worker().history().to_vec();
|
||||
shared_state.update_history(items);
|
||||
shared_state.set_user_segments(pod.user_segments().to_vec());
|
||||
// history / user_segments are no longer mirrored on PodSharedState —
|
||||
// clients reconstruct them from `Event::Snapshot` + live
|
||||
// `Event::Entry` deliveries driven by the session-log sink. We
|
||||
// only flip the status and kick post-run memory jobs here.
|
||||
set_controller_status(shared_state, runtime_dir, event_tx, new_status).await;
|
||||
let _ = runtime_dir.write_history(shared_state).await;
|
||||
pod.spawn_post_run_memory_jobs();
|
||||
}
|
||||
|
||||
@@ -167,8 +165,23 @@ impl PodController {
|
||||
}])
|
||||
.map_err(std::io::Error::other)?;
|
||||
|
||||
// === 1.5. Per-item history-commit drain task ===
|
||||
//
|
||||
// Worker callbacks fire `on_history_append` for each assistant
|
||||
// item / tool result / hook-injected item that lands in
|
||||
// history. The drain task picks them up off an unbounded mpsc
|
||||
// and commits each as a typed `LogEntry` through the sink,
|
||||
// serialised against the same `session_head` lock the Pod uses
|
||||
// for its own commits. This gives mid-turn snapshot visibility:
|
||||
// a late-attaching client sees in-flight tool calls + completed
|
||||
// assistant blocks without waiting for the turn-end persist.
|
||||
let (log_cmd_tx, log_cmd_rx) = mpsc::unbounded_channel::<LogCommand>();
|
||||
let drain_ctx = pod.log_drain_handle();
|
||||
let _drain_task = tokio::spawn(run_log_drain(log_cmd_rx, drain_ctx));
|
||||
pod.attach_log_cmd_tx(log_cmd_tx.clone());
|
||||
|
||||
// === 2. Worker event bridge wiring ===
|
||||
wire_event_bridges_on_worker(&mut pod, &event_tx, &alerter);
|
||||
wire_event_bridges_on_worker(&mut pod, &event_tx, &alerter, log_cmd_tx);
|
||||
|
||||
// === 3. Tool registration (builtin / memory / spawn-orchestration) ===
|
||||
let fs_for_view = register_pod_tools(
|
||||
@@ -193,8 +206,6 @@ impl PodController {
|
||||
manifest_toml.clone(),
|
||||
greeting,
|
||||
));
|
||||
shared_state.update_history(pod.worker().history().to_vec());
|
||||
shared_state.set_user_segments(pod.user_segments().to_vec());
|
||||
shared_state.set_fs_view(crate::fs_view::PodFsView::new(fs_for_view));
|
||||
shared_state.set_workflows(
|
||||
pod.workflow_completions()
|
||||
@@ -210,7 +221,6 @@ impl PodController {
|
||||
);
|
||||
runtime_dir.write_manifest(&manifest_toml).await?;
|
||||
runtime_dir.write_status(&shared_state).await?;
|
||||
runtime_dir.write_history(&shared_state).await?;
|
||||
|
||||
let handle = PodHandle {
|
||||
method_tx,
|
||||
@@ -218,6 +228,7 @@ impl PodController {
|
||||
shared_state: shared_state.clone(),
|
||||
runtime_dir: runtime_dir.clone(),
|
||||
alerter: alerter.clone(),
|
||||
sink: pod.sink(),
|
||||
};
|
||||
|
||||
let socket_server = SocketServer::start(&handle).await?;
|
||||
@@ -251,16 +262,30 @@ impl PodController {
|
||||
/// Wire the per-event broadcast bridges on the Pod's Worker. Each callback
|
||||
/// re-publishes a worker-level signal as a `protocol::Event` on `event_tx`
|
||||
/// so subscribers (TUI, socket clients) get a single typed stream.
|
||||
///
|
||||
/// Also wires `on_history_append` into the per-item drain channel so
|
||||
/// every history append observed by the worker becomes a typed
|
||||
/// `LogEntry` commit (via the drain task).
|
||||
fn wire_event_bridges_on_worker<C, St>(
|
||||
pod: &mut Pod<C, St>,
|
||||
event_tx: &broadcast::Sender<Event>,
|
||||
alerter: &Alerter,
|
||||
log_cmd_tx: mpsc::UnboundedSender<LogCommand>,
|
||||
) where
|
||||
C: LlmClient + Clone + 'static,
|
||||
St: Store + Clone + 'static,
|
||||
{
|
||||
let worker = pod.worker_mut();
|
||||
|
||||
// Per-history-append → drain channel. Sends are infallible-by-design
|
||||
// here (UnboundedSender never blocks); a closed receiver just means
|
||||
// the controller is shutting down, in which case dropping the item
|
||||
// is acceptable.
|
||||
let drain_tx = log_cmd_tx.clone();
|
||||
worker.on_history_append(move |item| {
|
||||
let _ = drain_tx.send(LogCommand::Item(item.clone()));
|
||||
});
|
||||
|
||||
let tx = event_tx.clone();
|
||||
worker.on_turn_start(move |turn| {
|
||||
let _ = tx.send(Event::TurnStart { turn });
|
||||
@@ -365,13 +390,80 @@ fn wire_event_bridges_on_worker<C, St>(
|
||||
alerter_for_worker.alert(AlertLevel::Warn, AlertSource::Worker, message.to_owned());
|
||||
});
|
||||
|
||||
let tx = event_tx.clone();
|
||||
worker.on_history_append(move |item| {
|
||||
if is_system_message_item(item) {
|
||||
let value = serde_json::to_value(item).expect("Item is Serialize");
|
||||
let _ = tx.send(Event::SystemMessage { item: value });
|
||||
// History-append broadcasts (previously `Event::SystemMessage`)
|
||||
// have been removed: every persistent history item is now committed
|
||||
// through the session-log sink as a typed `LogEntry`, and clients
|
||||
// see it via `Event::Snapshot` + live `Event::Entry`. The
|
||||
// per-item commit channel is wired at the top of this function.
|
||||
}
|
||||
|
||||
/// Drain task: consumes `LogCommand::Item` and `LogCommand::Flush`
|
||||
/// off the channel and commits each item as a typed `LogEntry` through
|
||||
/// the supplied store + sink. Lives as long as the controller; exits
|
||||
/// when the sender is dropped (controller shutdown).
|
||||
async fn run_log_drain<St>(
|
||||
mut rx: mpsc::UnboundedReceiver<LogCommand>,
|
||||
ctx: LogDrainHandle<St>,
|
||||
) where
|
||||
St: session_store::Store + Clone + Send + 'static,
|
||||
{
|
||||
while let Some(cmd) = rx.recv().await {
|
||||
match cmd {
|
||||
LogCommand::Item(item) => {
|
||||
let Some(entry) = classify_history_item(item) else {
|
||||
continue;
|
||||
};
|
||||
let mut head = ctx.session_head.lock().await;
|
||||
match session_store::append_entry_with_hash(
|
||||
&ctx.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
entry.clone(),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(_) => {
|
||||
// Publish under the same critical section view
|
||||
// a `subscribe_with_snapshot` would observe.
|
||||
ctx.sink.publish(entry);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, "drain: append_entry failed; entry dropped");
|
||||
}
|
||||
}
|
||||
}
|
||||
LogCommand::Flush(ack) => {
|
||||
let _ = ack.send(());
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
/// Map a single worker-history `Item` to its corresponding `LogEntry`
|
||||
/// classification. `None` is the skip signal for `user_message` items —
|
||||
/// those are committed via `LogEntry::UserInput` by `Pod::run` at
|
||||
/// submit time and would otherwise produce a duplicate entry here.
|
||||
fn classify_history_item(item: Item) -> Option<LogEntry> {
|
||||
let ts = session_log::now_millis();
|
||||
if item.is_user_message() {
|
||||
return None;
|
||||
}
|
||||
if item.is_tool_result() {
|
||||
return Some(LogEntry::ToolResults {
|
||||
ts,
|
||||
items: vec![session_store::LoggedItem::from(&item)],
|
||||
});
|
||||
}
|
||||
if item.is_assistant_message() || item.is_tool_call() || item.is_reasoning() {
|
||||
return Some(LogEntry::AssistantItems {
|
||||
ts,
|
||||
items: vec![session_store::LoggedItem::from(&item)],
|
||||
});
|
||||
}
|
||||
Some(LogEntry::HookInjectedItems {
|
||||
ts,
|
||||
items: vec![session_store::LoggedItem::from(&item)],
|
||||
})
|
||||
}
|
||||
|
||||
/// Register the builtin file-manipulation tools, optional memory tools,
|
||||
@@ -656,10 +748,9 @@ async fn controller_loop<C, St>(
|
||||
break;
|
||||
}
|
||||
|
||||
// GetHistory / ListCompletions are handled at the socket
|
||||
// layer (direct response). If they reach the controller,
|
||||
// ignore them.
|
||||
Method::GetHistory | Method::ListCompletions { .. } => {}
|
||||
// ListCompletions is handled at the socket layer (direct
|
||||
// response). If it reaches the controller, ignore it.
|
||||
Method::ListCompletions { .. } => {}
|
||||
|
||||
Method::PodEvent(event) => {
|
||||
// Echo the received event to all subscribers so every
|
||||
@@ -820,7 +911,7 @@ where
|
||||
// drain it at its next pre_llm_request.
|
||||
notify_buffer.push(message);
|
||||
}
|
||||
Some(Method::GetHistory | Method::ListCompletions { .. }) => {}
|
||||
Some(Method::ListCompletions { .. }) => {}
|
||||
Some(Method::PodEvent(event)) => {
|
||||
let _ = event_tx.send(Event::PodEvent(event.clone()));
|
||||
// mpsc is consume-once, so we cannot defer this
|
||||
|
||||
@@ -62,6 +62,13 @@ async fn handle_connection(stream: tokio::net::UnixStream, handle: PodHandle) {
|
||||
let mut reader = JsonLineReader::new(reader);
|
||||
let mut writer = JsonLineWriter::new(writer);
|
||||
|
||||
// Atomically subscribe to the session-log mirror first. The
|
||||
// returned (snapshot, rx) pair partitions the entry timeline:
|
||||
// entries committed before this call appear in `entries`, every
|
||||
// entry after lands on `entry_rx`. Doing this before the alert
|
||||
// snapshot keeps both ordering pairs internally consistent.
|
||||
let (entries_snapshot, mut entry_rx) = handle.sink.subscribe_with_snapshot();
|
||||
|
||||
// Atomically subscribe and snapshot buffered alerts so that
|
||||
// warnings emitted before this client connected are replayed
|
||||
// exactly once — they appear in the snapshot, and any alert
|
||||
@@ -73,8 +80,41 @@ async fn handle_connection(stream: tokio::net::UnixStream, handle: PodHandle) {
|
||||
}
|
||||
}
|
||||
|
||||
// Send the typed snapshot up front so late attachers can
|
||||
// reconstruct view state without an extra round trip.
|
||||
let snapshot_event = Event::Snapshot {
|
||||
entries: entries_snapshot
|
||||
.into_iter()
|
||||
.map(|e| serde_json::to_value(&e).expect("LogEntry is Serialize"))
|
||||
.collect(),
|
||||
greeting: handle.shared_state.greeting.clone(),
|
||||
status: handle.shared_state.get_status(),
|
||||
};
|
||||
if writer.write(&snapshot_event).await.is_err() {
|
||||
return;
|
||||
}
|
||||
|
||||
loop {
|
||||
tokio::select! {
|
||||
// Live session-log entries → this client as Event::Entry.
|
||||
entry = entry_rx.recv() => {
|
||||
match entry {
|
||||
Ok(entry) => {
|
||||
let value = serde_json::to_value(&entry)
|
||||
.expect("LogEntry is Serialize");
|
||||
if writer.write(&Event::Entry { entry: value }).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {
|
||||
// Slow client fell behind the broadcast buffer.
|
||||
// Drop the connection so the next reconnect
|
||||
// re-seeds the prefix via subscribe_with_snapshot.
|
||||
break;
|
||||
}
|
||||
Err(tokio::sync::broadcast::error::RecvError::Closed) => break,
|
||||
}
|
||||
}
|
||||
// Broadcast events → this client
|
||||
event = rx.recv() => {
|
||||
match event {
|
||||
@@ -129,57 +169,6 @@ async fn handle_connection(stream: tokio::net::UnixStream, handle: PodHandle) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Ok(Some(Method::GetHistory)) => {
|
||||
let items = handle.shared_state.history();
|
||||
let segments_per_user = handle.shared_state.user_segments();
|
||||
// Embed `segments` on user-message JSON values so
|
||||
// the TUI can re-render typed atoms on restore.
|
||||
// Alignment: segments are recorded only for
|
||||
// submissions made during the live session, never
|
||||
// for seed history loaded via `SessionStart.history`
|
||||
// (post-compaction). The seed user_messages always
|
||||
// come first in worker history, so the last
|
||||
// `segments_per_user.len()` user_messages are the
|
||||
// ones that map 1:1 to the segments list.
|
||||
let total_user_msgs =
|
||||
items.iter().filter(|i| i.is_user_message()).count();
|
||||
let skip = total_user_msgs.saturating_sub(segments_per_user.len());
|
||||
let mut user_idx = 0usize;
|
||||
let values = items
|
||||
.iter()
|
||||
.map(|item| {
|
||||
let mut value =
|
||||
serde_json::to_value(item).expect("Item is Serialize");
|
||||
if item.is_user_message() {
|
||||
if user_idx >= skip {
|
||||
let seg_idx = user_idx - skip;
|
||||
if let Some(obj) = value.as_object_mut() {
|
||||
let segs = serde_json::to_value(
|
||||
&segments_per_user[seg_idx],
|
||||
)
|
||||
.expect("Segment is Serialize");
|
||||
obj.insert("segments".into(), segs);
|
||||
}
|
||||
}
|
||||
user_idx += 1;
|
||||
}
|
||||
value
|
||||
})
|
||||
.collect();
|
||||
let greeting = handle.shared_state.greeting.clone();
|
||||
let status = handle.shared_state.get_status();
|
||||
if writer
|
||||
.write(&Event::History {
|
||||
items: values,
|
||||
greeting,
|
||||
status,
|
||||
})
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
break;
|
||||
}
|
||||
}
|
||||
Ok(Some(method)) => {
|
||||
let _ = handle.send(method).await;
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ pub mod hook;
|
||||
pub mod ipc;
|
||||
pub mod prompt;
|
||||
pub mod runtime;
|
||||
pub mod session_log_sink;
|
||||
pub mod shared_state;
|
||||
pub mod spawn;
|
||||
pub mod workflow;
|
||||
@@ -30,4 +31,5 @@ pub use prompt::system::{SystemPromptContext, SystemPromptError, SystemPromptTem
|
||||
pub use protocol::{ErrorCode, Event, Method, PodStatus, TurnResult};
|
||||
pub use provider::{ProviderError, build_client};
|
||||
pub use runtime::dir::RuntimeDir;
|
||||
pub use session_log_sink::{SessionLogSink, SessionLogWriter};
|
||||
pub use shared_state::PodSharedState;
|
||||
|
||||
+341
-162
@@ -7,10 +7,40 @@ use llm_worker::Item;
|
||||
use llm_worker::llm_client::RequestConfig;
|
||||
use llm_worker::llm_client::client::LlmClient;
|
||||
use llm_worker::state::Mutable;
|
||||
use llm_worker::{Role, ToolOutputLimits, UsageRecord, Worker, WorkerError, WorkerResult};
|
||||
use session_store::{EntryHash, PodScopeSnapshot, SessionId, SessionStartState, Store, StoreError};
|
||||
use llm_worker::{ToolOutputLimits, UsageRecord, Worker, WorkerError, WorkerResult};
|
||||
use session_store::{
|
||||
EntryHash, HashedEntry, LogEntry, PodScopeSnapshot, SessionId, Store, StoreError, session_log,
|
||||
to_logged,
|
||||
};
|
||||
use tracing::{info, warn};
|
||||
|
||||
use crate::session_log_sink::SessionLogSink;
|
||||
|
||||
/// Command sent to the per-Pod history-drain task.
|
||||
///
|
||||
/// `Item` carries one worker-history append observed via
|
||||
/// `Worker::on_history_append`; the drain classifies it into a
|
||||
/// `LogEntry::AssistantItems` / `LogEntry::ToolResults` /
|
||||
/// `LogEntry::HookInjectedItems` and commits it through the sink.
|
||||
/// `Flush(ack)` is the barrier used by `persist_turn` to ensure every
|
||||
/// in-flight item is committed before the trailing `TurnEnd` entry
|
||||
/// lands.
|
||||
#[derive(Debug)]
|
||||
pub enum LogCommand {
|
||||
Item(Item),
|
||||
Flush(tokio::sync::oneshot::Sender<()>),
|
||||
}
|
||||
|
||||
/// State shared between Pod and the controller-spawned history-drain
|
||||
/// task: store + session-head lock + broadcast sink. All three are
|
||||
/// `Clone`able (the latter two as `Arc` clones, the store per its
|
||||
/// `Clone` impl) so handing a copy to the drain task is cheap.
|
||||
pub struct LogDrainHandle<St> {
|
||||
pub store: St,
|
||||
pub session_head: Arc<AsyncMutex<SessionHead>>,
|
||||
pub sink: SessionLogSink,
|
||||
}
|
||||
|
||||
use manifest::{
|
||||
Permission, PodManifest, PodManifestConfig, ResolveError, Scope, ScopeConfig, ScopeError,
|
||||
ScopeRule, SharedScope, WorkerManifest,
|
||||
@@ -38,9 +68,9 @@ use protocol::{AlertLevel, AlertSource, Event, Segment};
|
||||
use tokio::sync::broadcast;
|
||||
use tokio::task::JoinHandle;
|
||||
|
||||
struct SessionHead {
|
||||
session_id: SessionId,
|
||||
head_hash: Option<EntryHash>,
|
||||
pub struct SessionHead {
|
||||
pub session_id: SessionId,
|
||||
pub head_hash: Option<EntryHash>,
|
||||
}
|
||||
|
||||
/// Pre-LLM-request hook that records `history.len()` at send time into a
|
||||
@@ -190,9 +220,22 @@ pub struct Pod<C: LlmClient, St: Store> {
|
||||
/// the K-th `Item::user_message` in `worker.history()` (modulo seed
|
||||
/// history loaded via `SessionStart.history`, whose original segments
|
||||
/// are not preserved). Populated from log on `restore_from_manifest`,
|
||||
/// appended after `save_user_input` on each `run`. Mirrored to
|
||||
/// `PodSharedState` by the controller for `Event::History` use.
|
||||
/// appended after `save_user_input` on each `run`. Pre-`Event::Snapshot`
|
||||
/// this fed `PodSharedState.user_segments`; the new wire format
|
||||
/// carries typed atoms via `LogEntry::UserInput { segments }` so
|
||||
/// this remains purely an in-memory tracker for compact alignment.
|
||||
user_segments: Vec<Vec<Segment>>,
|
||||
/// Pod-side session-log mirror + broadcast sink. Populated alongside
|
||||
/// every successful `session_store::append_entry` write so connected
|
||||
/// clients see a `(snapshot, live)` stream consistent with what's
|
||||
/// on disk.
|
||||
sink: SessionLogSink,
|
||||
/// Sender into the controller-spawned history-drain task.
|
||||
/// `None` when no controller has wired one (tests, low-level Pod
|
||||
/// usage). The drain task is the source of mid-turn `AssistantItems`
|
||||
/// / `ToolResults` / `HookInjectedItems` commits, fed by the
|
||||
/// `Worker::on_history_append` callback.
|
||||
log_cmd_tx: Option<tokio::sync::mpsc::UnboundedSender<LogCommand>>,
|
||||
}
|
||||
|
||||
impl<C: LlmClient + 'static, St: Store + 'static> Pod<C, St> {
|
||||
@@ -249,6 +292,22 @@ impl<C: LlmClient + Clone + 'static, St: Store + Clone + 'static> Pod<C, St> {
|
||||
extract_pointer: self.extract_pointer.clone(),
|
||||
memory_task: None,
|
||||
user_segments: self.user_segments.clone(),
|
||||
// The memory-task clone never appends to the session log
|
||||
// (it only reads `worker.history()`), so a fresh sink is
|
||||
// fine — nothing observes its broadcast.
|
||||
sink: SessionLogSink::new(),
|
||||
log_cmd_tx: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Build a `LogDrainHandle` carrying everything the controller's
|
||||
/// drain task needs: store handle, the shared session-head lock,
|
||||
/// and the broadcast sink. All three are cheap clones.
|
||||
pub fn log_drain_handle(&self) -> LogDrainHandle<St> {
|
||||
LogDrainHandle {
|
||||
store: self.store.clone(),
|
||||
session_head: self.session_head.clone(),
|
||||
sink: self.sink.clone(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -332,6 +391,8 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
extract_pointer: Arc::new(Mutex::new(None)),
|
||||
memory_task: None,
|
||||
user_segments: Vec::new(),
|
||||
sink: SessionLogSink::new(),
|
||||
log_cmd_tx: None,
|
||||
};
|
||||
pod.apply_permissions_from_manifest();
|
||||
pod.apply_prune_from_manifest();
|
||||
@@ -426,8 +487,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
/// can restore the narrowed scope instead of reclaiming delegated
|
||||
/// writes.
|
||||
pub async fn persist_scope_snapshot(&mut self) -> Result<(), StoreError> {
|
||||
let mut head = self.session_head.lock().await;
|
||||
if head.head_hash.is_none() {
|
||||
if self.session_head.lock().await.head_hash.is_none() {
|
||||
return Ok(());
|
||||
}
|
||||
let snapshot = {
|
||||
@@ -437,8 +497,50 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
deny: scope.deny_rules(),
|
||||
}
|
||||
};
|
||||
session_store::save_pod_scope(&self.store, head.session_id, &mut head.head_hash, &snapshot)
|
||||
.await
|
||||
let payload = serde_json::to_value(&snapshot).expect("PodScopeSnapshot is Serialize");
|
||||
self.commit_entry(LogEntry::Extension {
|
||||
ts: session_log::now_millis(),
|
||||
domain: session_store::POD_SCOPE_EXTENSION_DOMAIN.into(),
|
||||
payload,
|
||||
})
|
||||
.await
|
||||
.map(|_| ())
|
||||
}
|
||||
|
||||
/// Append `entry` to the session log AND publish it through the
|
||||
/// broadcast sink. Holds the session-head async lock across the
|
||||
/// disk write and the sink publish so subscribers see a gap-free
|
||||
/// `(snapshot, live)` stream consistent with what's on disk.
|
||||
pub(crate) async fn commit_entry(
|
||||
&self,
|
||||
entry: LogEntry,
|
||||
) -> Result<EntryHash, StoreError> {
|
||||
let mut head = self.session_head.lock().await;
|
||||
let hash = session_store::append_entry_with_hash(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
entry.clone(),
|
||||
)
|
||||
.await?;
|
||||
self.sink.publish(entry);
|
||||
Ok(hash)
|
||||
}
|
||||
|
||||
/// Cloneable sink handle. Exposed to the controller so the IPC
|
||||
/// layer can `subscribe_with_snapshot` and stream entries to
|
||||
/// clients without consulting any other state.
|
||||
pub fn sink(&self) -> SessionLogSink {
|
||||
self.sink.clone()
|
||||
}
|
||||
|
||||
/// Wire a history-drain task. The controller calls this once per
|
||||
/// Pod after the drain task is spawned; the matching mpsc receiver
|
||||
/// drives per-item commits of assistant items / tool results /
|
||||
/// hook-injected items committed by the worker via
|
||||
/// `Worker::on_history_append`.
|
||||
pub fn attach_log_cmd_tx(&mut self, tx: tokio::sync::mpsc::UnboundedSender<LogCommand>) {
|
||||
self.log_cmd_tx = Some(tx);
|
||||
}
|
||||
|
||||
/// Cloneable callback handed to dynamic-scope tools. It cannot append
|
||||
@@ -459,13 +561,12 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
.expect("pending_scope_snapshot poisoned")
|
||||
.take();
|
||||
if let Some(snapshot) = snapshot {
|
||||
let mut head = self.session_head.lock().await;
|
||||
session_store::save_pod_scope(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
&snapshot,
|
||||
)
|
||||
let payload = serde_json::to_value(&snapshot).expect("PodScopeSnapshot is Serialize");
|
||||
self.commit_entry(LogEntry::Extension {
|
||||
ts: session_log::now_millis(),
|
||||
domain: session_store::POD_SCOPE_EXTENSION_DOMAIN.into(),
|
||||
payload,
|
||||
})
|
||||
.await?;
|
||||
}
|
||||
Ok(())
|
||||
@@ -629,15 +730,13 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
/// (the entry is dropped) and a `Warn` alert + `tracing::warn!` are
|
||||
/// emitted so the failure isn't completely silent.
|
||||
async fn try_record_metric(&mut self, metric: &session_metrics::Metric) {
|
||||
let mut head = self.session_head.lock().await;
|
||||
if let Err(err) = session_metrics::record_metric(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
metric,
|
||||
)
|
||||
.await
|
||||
{
|
||||
let payload = serde_json::to_value(metric).expect("Metric is Serialize");
|
||||
let entry = LogEntry::Extension {
|
||||
ts: session_log::now_millis(),
|
||||
domain: session_metrics::DOMAIN.into(),
|
||||
payload,
|
||||
};
|
||||
if let Err(err) = self.commit_entry(entry).await {
|
||||
warn!(name = %metric.name, error = %err, "failed to record session metric; dropping");
|
||||
self.alert(
|
||||
AlertLevel::Warn,
|
||||
@@ -656,20 +755,6 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
}
|
||||
}
|
||||
|
||||
fn broadcast_system_message_item(&self, item: &Item) {
|
||||
if !matches!(
|
||||
item,
|
||||
Item::Message {
|
||||
role: Role::System,
|
||||
..
|
||||
}
|
||||
) {
|
||||
return;
|
||||
}
|
||||
let value = serde_json::to_value(item).expect("Item is Serialize");
|
||||
self.send_event(Event::SystemMessage { item: value });
|
||||
}
|
||||
|
||||
/// Push a `Method::Notify` (or rendered `Method::PodEvent`) entry
|
||||
/// onto the pending buffer.
|
||||
///
|
||||
@@ -975,17 +1060,12 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
// Persist the user input as typed segments before the worker
|
||||
// pushes its flattened copy into history. save_delta deliberately
|
||||
// skips the resulting `is_user_message()` item to avoid double-write.
|
||||
{
|
||||
let mut head = self.session_head.lock().await;
|
||||
self.session_id = head.session_id;
|
||||
session_store::save_user_input(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
input.clone(),
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
self.session_id = self.session_head.lock().await.session_id;
|
||||
self.commit_entry(LogEntry::UserInput {
|
||||
ts: session_log::now_millis(),
|
||||
segments: input.clone(),
|
||||
})
|
||||
.await?;
|
||||
self.user_segments.push(input.clone());
|
||||
|
||||
// Resolve `@<path>` refs, `#<slug>` Knowledge refs, and `/<slug>`
|
||||
@@ -1330,34 +1410,64 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
/// another writer has advanced the store head behind our back.
|
||||
async fn ensure_session_head(&mut self) -> Result<(), PodError> {
|
||||
let w = self.worker.as_ref().unwrap();
|
||||
let state = SessionStartState {
|
||||
system_prompt: w.get_system_prompt(),
|
||||
config: w.request_config(),
|
||||
history: w.history(),
|
||||
let prev_session_id;
|
||||
let initial_state = {
|
||||
let head = self.session_head.lock().await;
|
||||
prev_session_id = head.session_id;
|
||||
head.head_hash.is_none()
|
||||
};
|
||||
let mut head = self.session_head.lock().await;
|
||||
if head.head_hash.is_none() {
|
||||
let hash =
|
||||
session_store::create_session_with_id(&self.store, head.session_id, state).await?;
|
||||
head.head_hash = Some(hash);
|
||||
drop(head);
|
||||
if initial_state {
|
||||
let initial = LogEntry::SessionStart {
|
||||
ts: session_log::now_millis(),
|
||||
system_prompt: w.get_system_prompt().map(String::from),
|
||||
config: w.request_config().clone(),
|
||||
history: to_logged(w.history()),
|
||||
forked_from: None,
|
||||
compacted_from: None,
|
||||
};
|
||||
self.commit_entry(initial).await?;
|
||||
self.persist_scope_snapshot().await?;
|
||||
return Ok(());
|
||||
}
|
||||
let prev_session_id = head.session_id;
|
||||
let mut session_id = head.session_id;
|
||||
let mut head_hash = head.head_hash.clone();
|
||||
session_store::ensure_head_or_fork(&self.store, &mut session_id, &mut head_hash, state)
|
||||
.await?;
|
||||
head.session_id = session_id;
|
||||
head.head_hash = head_hash;
|
||||
self.session_id = session_id;
|
||||
// ensure_head_or_fork mints a fresh session_id when it auto-
|
||||
// forks. Sync that to pods.json so a concurrent
|
||||
// restore_from_manifest can't see "no live writer" for the new
|
||||
// session and grab it.
|
||||
if session_id != prev_session_id && self.scope_allocation.is_some() {
|
||||
pod_registry::update_session(&self.manifest.pod.name, session_id)?;
|
||||
// Check store head + auto-fork if it drifted.
|
||||
let store_head = self
|
||||
.store
|
||||
.read_head_hash(prev_session_id)
|
||||
.await
|
||||
.map_err(PodError::from)?;
|
||||
let mut head = self.session_head.lock().await;
|
||||
if store_head == head.head_hash {
|
||||
return Ok(());
|
||||
}
|
||||
// Fork: mint a fresh session and switch to it. The new
|
||||
// SessionStart entry replaces the mirror and is broadcast
|
||||
// through the sink so existing subscribers reset their view.
|
||||
let fork_id = session_store::new_session_id();
|
||||
let entry = LogEntry::SessionStart {
|
||||
ts: session_log::now_millis(),
|
||||
system_prompt: w.get_system_prompt().map(String::from),
|
||||
config: w.request_config().clone(),
|
||||
history: to_logged(w.history()),
|
||||
forked_from: None,
|
||||
compacted_from: None,
|
||||
};
|
||||
let hash = session_log::compute_hash(None, &entry);
|
||||
let hashed = HashedEntry {
|
||||
hash: hash.clone(),
|
||||
prev_hash: None,
|
||||
entry: entry.clone(),
|
||||
};
|
||||
self.store
|
||||
.create_session(fork_id, &[hashed])
|
||||
.await
|
||||
.map_err(PodError::from)?;
|
||||
head.session_id = fork_id;
|
||||
head.head_hash = Some(hash);
|
||||
self.session_id = fork_id;
|
||||
self.sink.reset_with_initial(entry);
|
||||
drop(head);
|
||||
if self.scope_allocation.is_some() {
|
||||
pod_registry::update_session(&self.manifest.pod.name, fork_id)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
@@ -1493,28 +1603,84 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
history_before: usize,
|
||||
result: &Result<WorkerResult, WorkerError>,
|
||||
) -> Result<(), StoreError> {
|
||||
// Use direct field access for split borrows (worker immutable,
|
||||
// head_hash mutable).
|
||||
let w = self.worker.as_ref().unwrap();
|
||||
let new_items = &w.history()[history_before..];
|
||||
let mut head = self.session_head.lock().await;
|
||||
self.session_id = head.session_id;
|
||||
session_store::save_delta(&self.store, head.session_id, &mut head.head_hash, new_items)
|
||||
.await?;
|
||||
// Per-item commits for AssistantItems / ToolResults /
|
||||
// HookInjectedItems already landed mid-turn through the
|
||||
// controller-spawned drain task, fed by
|
||||
// `Worker::on_history_append`. Drain the queue here so every
|
||||
// in-flight item has actually been committed before the
|
||||
// trailing `TurnEnd` entry. When no drain is wired (low-level
|
||||
// tests / direct `Pod::new` usage) we fall back to a synchronous
|
||||
// pass that replicates the legacy `save_delta` classification —
|
||||
// those code paths don't fire `on_history_append`, so the items
|
||||
// would otherwise be lost.
|
||||
let _ = history_before; // referenced only by the fallback below.
|
||||
self.session_id = self.session_head.lock().await.session_id;
|
||||
if let Some(tx) = self.log_cmd_tx.as_ref() {
|
||||
let (ack_tx, ack_rx) = tokio::sync::oneshot::channel();
|
||||
if tx.send(LogCommand::Flush(ack_tx)).is_ok() {
|
||||
let _ = ack_rx.await;
|
||||
}
|
||||
} else {
|
||||
// Fallback path for tests / Pod::new: classify and commit
|
||||
// the post-`history_before` slice inline, matching the old
|
||||
// `save_delta` shape.
|
||||
let new_items: Vec<Item> = self.worker.as_ref().unwrap().history()[history_before..]
|
||||
.iter()
|
||||
.cloned()
|
||||
.collect();
|
||||
let ts = session_log::now_millis();
|
||||
let mut i = 0;
|
||||
while i < new_items.len() {
|
||||
let item = &new_items[i];
|
||||
if item.is_user_message() {
|
||||
i += 1;
|
||||
} else if item.is_tool_result() {
|
||||
let start = i;
|
||||
while i < new_items.len() && new_items[i].is_tool_result() {
|
||||
i += 1;
|
||||
}
|
||||
let items = new_items[start..i]
|
||||
.iter()
|
||||
.map(session_store::LoggedItem::from)
|
||||
.collect();
|
||||
self.commit_entry(LogEntry::ToolResults { ts, items }).await?;
|
||||
} else if item.is_assistant_message()
|
||||
|| item.is_tool_call()
|
||||
|| item.is_reasoning()
|
||||
{
|
||||
let start = i;
|
||||
while i < new_items.len()
|
||||
&& (new_items[i].is_assistant_message()
|
||||
|| new_items[i].is_tool_call()
|
||||
|| new_items[i].is_reasoning())
|
||||
{
|
||||
i += 1;
|
||||
}
|
||||
let items = new_items[start..i]
|
||||
.iter()
|
||||
.map(session_store::LoggedItem::from)
|
||||
.collect();
|
||||
self.commit_entry(LogEntry::AssistantItems { ts, items })
|
||||
.await?;
|
||||
} else {
|
||||
self.commit_entry(LogEntry::HookInjectedItems {
|
||||
ts,
|
||||
items: vec![session_store::LoggedItem::from(&new_items[i])],
|
||||
})
|
||||
.await?;
|
||||
i += 1;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
drop(head);
|
||||
self.flush_pending_scope_snapshot().await?;
|
||||
|
||||
let turn_count = self.worker.as_ref().unwrap().turn_count();
|
||||
let mut head = self.session_head.lock().await;
|
||||
session_store::save_turn_end(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
self.commit_entry(LogEntry::TurnEnd {
|
||||
ts: session_log::now_millis(),
|
||||
turn_count,
|
||||
)
|
||||
})
|
||||
.await?;
|
||||
drop(head);
|
||||
|
||||
// Flush any sync-buffered metrics from this run first
|
||||
// (currently `prune.fire` / `prune.skip` from the prune observer).
|
||||
@@ -1547,19 +1713,15 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
record,
|
||||
correlation_id,
|
||||
} = recorded;
|
||||
let mut head = self.session_head.lock().await;
|
||||
session_store::save_usage(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
record.history_len,
|
||||
record.input_total_tokens,
|
||||
record.cache_read_tokens,
|
||||
record.cache_write_tokens,
|
||||
record.output_tokens,
|
||||
)
|
||||
self.commit_entry(LogEntry::LlmUsage {
|
||||
ts: session_log::now_millis(),
|
||||
history_len: record.history_len,
|
||||
input_total_tokens: record.input_total_tokens,
|
||||
cache_read_tokens: record.cache_read_tokens,
|
||||
cache_write_tokens: record.cache_write_tokens,
|
||||
output_tokens: record.output_tokens,
|
||||
})
|
||||
.await?;
|
||||
drop(head);
|
||||
if let Some(id) = correlation_id {
|
||||
let metric = session_metrics::Metric::now("prune.post_request")
|
||||
.with_correlation_id(&id)
|
||||
@@ -1577,25 +1739,19 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
let interrupted = self.worker.as_ref().unwrap().last_run_interrupted();
|
||||
match result {
|
||||
Ok(r) => {
|
||||
let mut head = self.session_head.lock().await;
|
||||
session_store::save_run_completed(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
r.clone(),
|
||||
self.commit_entry(LogEntry::RunCompleted {
|
||||
ts: session_log::now_millis(),
|
||||
interrupted,
|
||||
)
|
||||
result: r.clone(),
|
||||
})
|
||||
.await?;
|
||||
}
|
||||
Err(e) => {
|
||||
let mut head = self.session_head.lock().await;
|
||||
session_store::save_run_errored(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
e.to_string(),
|
||||
self.commit_entry(LogEntry::RunErrored {
|
||||
ts: session_log::now_millis(),
|
||||
interrupted,
|
||||
)
|
||||
message: e.to_string(),
|
||||
})
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
@@ -1824,34 +1980,46 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
task_snapshot_text.clone(),
|
||||
));
|
||||
|
||||
// Persist as a new compacted session.
|
||||
let mut head = self.session_head.lock().await;
|
||||
let old_session_id = head.session_id;
|
||||
let old_head_hash = head
|
||||
.head_hash
|
||||
.clone()
|
||||
.expect("head_hash should be set after at least one entry");
|
||||
|
||||
let w = self.worker.as_ref().unwrap();
|
||||
let state = SessionStartState {
|
||||
system_prompt: w.get_system_prompt(),
|
||||
config: w.request_config(),
|
||||
history: &new_history,
|
||||
// Build the SessionStart entry for the new compacted session,
|
||||
// then atomically rotate to it: create on disk, swap head, reset
|
||||
// the broadcast sink so existing subscribers see the new
|
||||
// `SessionStart { compacted_from }` and reset their view.
|
||||
let new_session_id = session_store::new_session_id();
|
||||
let session_start = {
|
||||
let mut head = self.session_head.lock().await;
|
||||
let old_session_id = head.session_id;
|
||||
let old_head_hash = head
|
||||
.head_hash
|
||||
.clone()
|
||||
.expect("head_hash should be set after at least one entry");
|
||||
let w = self.worker.as_ref().unwrap();
|
||||
let entry = LogEntry::SessionStart {
|
||||
ts: session_log::now_millis(),
|
||||
system_prompt: w.get_system_prompt().map(String::from),
|
||||
config: w.request_config().clone(),
|
||||
history: to_logged(&new_history),
|
||||
forked_from: None,
|
||||
compacted_from: Some(session_store::SessionOrigin {
|
||||
session_id: old_session_id,
|
||||
at_hash: old_head_hash,
|
||||
}),
|
||||
};
|
||||
let hash = session_log::compute_hash(None, &entry);
|
||||
let hashed = HashedEntry {
|
||||
hash: hash.clone(),
|
||||
prev_hash: None,
|
||||
entry: entry.clone(),
|
||||
};
|
||||
self.store.create_session(new_session_id, &[hashed]).await?;
|
||||
head.session_id = new_session_id;
|
||||
head.head_hash = Some(hash);
|
||||
self.session_id = new_session_id;
|
||||
entry
|
||||
};
|
||||
let (new_session_id, new_head_hash) = session_store::create_compacted_session(
|
||||
&self.store,
|
||||
state,
|
||||
old_session_id,
|
||||
old_head_hash,
|
||||
)
|
||||
.await?;
|
||||
|
||||
// Swap in the new session state. usage_history belongs to the old
|
||||
// session — the new compacted session starts with no measurements
|
||||
// until its first LLM call.
|
||||
self.session_id = new_session_id;
|
||||
head.session_id = new_session_id;
|
||||
head.head_hash = Some(new_head_hash);
|
||||
// Broadcast the SessionStart through the sink. This atomically
|
||||
// resets the mirror to `[SessionStart]` so any subscriber
|
||||
// querying after this point sees the post-compaction prefix.
|
||||
self.sink.reset_with_initial(session_start);
|
||||
// Keep pods.json pointing at the live session_id. Without this
|
||||
// a concurrent `restore_from_manifest(new_session_id)` would
|
||||
// see no live writer and grab the session this Pod just moved
|
||||
@@ -1861,7 +2029,6 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
if self.scope_allocation.is_some() {
|
||||
pod_registry::update_session(&self.manifest.pod.name, new_session_id)?;
|
||||
}
|
||||
drop(head);
|
||||
// Align user_segments with the post-compaction history. Items
|
||||
// before `retain_from` (now folded into the summary) lose their
|
||||
// segments; only the user_messages surviving in retained_items
|
||||
@@ -1873,9 +2040,11 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
}
|
||||
|
||||
self.worker.as_mut().unwrap().set_history(new_history);
|
||||
for item in &compact_introduced_system_messages {
|
||||
self.broadcast_system_message_item(item);
|
||||
}
|
||||
// Compaction-introduced system messages are part of the new
|
||||
// SessionStart's history (broadcast above) — clients derive
|
||||
// their blocks from `SessionStart.history`. No per-item
|
||||
// broadcast is required.
|
||||
let _ = &compact_introduced_system_messages;
|
||||
let worker = self.worker.as_mut().unwrap();
|
||||
// Anchor the prompt cache at the summary item so that Anthropic
|
||||
// can place a durable `cache_control` breakpoint there — our
|
||||
@@ -2139,18 +2308,13 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
};
|
||||
let payload_value = serde_json::to_value(&pointer_payload)
|
||||
.expect("ExtractPointerPayload is always JSON-serializable");
|
||||
{
|
||||
let mut head = self.session_head.lock().await;
|
||||
session_store::save_extension(
|
||||
&self.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
extract::EXTRACT_DOMAIN,
|
||||
payload_value,
|
||||
)
|
||||
.await?;
|
||||
self.session_id = head.session_id;
|
||||
}
|
||||
self.commit_entry(LogEntry::Extension {
|
||||
ts: session_log::now_millis(),
|
||||
domain: extract::EXTRACT_DOMAIN.into(),
|
||||
payload: payload_value,
|
||||
})
|
||||
.await?;
|
||||
self.session_id = self.session_head.lock().await.session_id;
|
||||
|
||||
*self
|
||||
.extract_pointer
|
||||
@@ -2488,6 +2652,8 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
extract_pointer: Arc::new(Mutex::new(None)),
|
||||
memory_task: None,
|
||||
user_segments: Vec::new(),
|
||||
sink: SessionLogSink::new(),
|
||||
log_cmd_tx: None,
|
||||
};
|
||||
pod.apply_permissions_from_manifest();
|
||||
pod.apply_prune_from_manifest();
|
||||
@@ -2559,6 +2725,8 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
extract_pointer: Arc::new(Mutex::new(None)),
|
||||
memory_task: None,
|
||||
user_segments: Vec::new(),
|
||||
sink: SessionLogSink::new(),
|
||||
log_cmd_tx: None,
|
||||
};
|
||||
pod.apply_permissions_from_manifest();
|
||||
pod.apply_prune_from_manifest();
|
||||
@@ -2590,10 +2758,16 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
store: St,
|
||||
loader: PromptLoader,
|
||||
) -> Result<Self, PodError> {
|
||||
let state = session_store::restore(&store, session_id).await?;
|
||||
// Read raw entries once so we can both reconstruct state and
|
||||
// seed the broadcast sink's mirror with the same prefix that
|
||||
// sits on disk.
|
||||
let raw_entries = store.read_all(session_id).await?;
|
||||
let state = session_store::collect_state(&raw_entries);
|
||||
if state.head_hash.is_none() {
|
||||
return Err(PodError::SessionEmpty { session_id });
|
||||
}
|
||||
let mirror_entries: Vec<LogEntry> =
|
||||
raw_entries.iter().map(|e| e.entry.clone()).collect();
|
||||
let scope_snapshot = state
|
||||
.pod_scope
|
||||
.clone()
|
||||
@@ -2696,6 +2870,11 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
extract_pointer: Arc::new(Mutex::new(extract_pointer)),
|
||||
memory_task: None,
|
||||
user_segments: state.user_segments,
|
||||
// Seed the mirror with the entries we just replayed so a
|
||||
// late-attaching client sees the full prefix without an
|
||||
// extra round trip.
|
||||
sink: SessionLogSink::with_initial(mirror_entries),
|
||||
log_cmd_tx: None,
|
||||
};
|
||||
pod.apply_permissions_from_manifest();
|
||||
pod.apply_prune_from_manifest();
|
||||
|
||||
@@ -73,12 +73,6 @@ impl RuntimeDir {
|
||||
atomic_write(&self.path.join("manifest.toml"), toml.as_bytes()).await
|
||||
}
|
||||
|
||||
/// Write history.json atomically.
|
||||
pub async fn write_history(&self, state: &PodSharedState) -> Result<(), io::Error> {
|
||||
let content = state.history_json();
|
||||
atomic_write(&self.path.join("history.json"), content.as_bytes()).await
|
||||
}
|
||||
|
||||
/// Write `spawned_pods.json` atomically. The entries are the full
|
||||
/// set of spawned children known to this Pod — callers pass the
|
||||
/// replacement list, no incremental merge.
|
||||
@@ -223,18 +217,6 @@ mod tests {
|
||||
assert_eq!(parsed[0].pod_name, "child");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn write_history_creates_file() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let rt = RuntimeDir::create(tmp.path(), "my-pod").await.unwrap();
|
||||
let state = test_state();
|
||||
|
||||
rt.write_history(&state).await.unwrap();
|
||||
|
||||
let content = std::fs::read_to_string(rt.path().join("history.json")).unwrap();
|
||||
assert_eq!(content, "[]");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn socket_path() {
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
|
||||
@@ -0,0 +1,474 @@
|
||||
//! Pod-side session-log mirror + broadcast.
|
||||
//!
|
||||
//! Owns the in-memory `Vec<LogEntry>` mirror that backs `Event::Snapshot`
|
||||
//! delivery to newly connected clients and the
|
||||
//! `broadcast::Sender<LogEntry>` that fans out per-entry commits to
|
||||
//! existing subscribers. Disk writes remain the responsibility of the
|
||||
//! Pod (which still owns the `Store` handle); the sink stays focused on
|
||||
//! the wire-side fan-out.
|
||||
//!
|
||||
//! Atomicity contract (see ticket `tickets/pod-state-from-session-log.md`):
|
||||
//!
|
||||
//! 1. Pod writes the entry to disk via the `Store`.
|
||||
//! 2. Pod calls [`SessionLogSink::publish`] which acquires the mirror
|
||||
//! mutex, pushes the entry, and fires `broadcast::send` — all under
|
||||
//! the same critical section.
|
||||
//!
|
||||
//! [`SessionLogSink::subscribe_with_snapshot`] takes the same mutex,
|
||||
//! so the `(snapshot, receiver)` pair returned to a connecting client
|
||||
//! splits the entry sequence cleanly: every entry shows up in exactly
|
||||
//! one of `snapshot` or on `receiver`.
|
||||
//!
|
||||
//! Disk-write failures short-circuit before `publish`, so a failed
|
||||
//! entry never appears in the mirror or on the broadcast.
|
||||
|
||||
use std::sync::{Arc, Mutex as StdMutex};
|
||||
|
||||
use session_store::{
|
||||
EntryHash, HashedEntry, LogEntry, SessionId, SessionStartState, Store, StoreError, session_log,
|
||||
};
|
||||
use tokio::sync::{Mutex as AsyncMutex, MutexGuard, broadcast};
|
||||
|
||||
/// Broadcast capacity for the live receiver. Slow subscribers that
|
||||
/// fall behind will see `RecvError::Lagged` and are expected to drop
|
||||
/// the connection so that the next reconnect's `subscribe_with_snapshot`
|
||||
/// re-seeds the prefix.
|
||||
const BROADCAST_CAPACITY: usize = 256;
|
||||
|
||||
/// In-memory mirror + broadcast fan-out for the active session log.
|
||||
///
|
||||
/// Clone is cheap (`Arc` clone) — the Pod hands one to the IPC layer
|
||||
/// for read-only `subscribe_with_snapshot` access and keeps one for
|
||||
/// its own write path.
|
||||
#[derive(Clone)]
|
||||
pub struct SessionLogSink {
|
||||
inner: Arc<SinkInner>,
|
||||
}
|
||||
|
||||
struct SinkInner {
|
||||
/// Full session log mirror in commit order. Reset on session swap
|
||||
/// (compaction / fork) via [`SessionLogSink::reset_with_initial`].
|
||||
mirror: StdMutex<Vec<LogEntry>>,
|
||||
/// Broadcast channel for live entry updates. The same `Sender`
|
||||
/// survives session swaps so existing subscribers keep their
|
||||
/// receiver — they observe the swap as a freshly broadcast
|
||||
/// `LogEntry::SessionStart` and reset their view accordingly.
|
||||
broadcast_tx: broadcast::Sender<LogEntry>,
|
||||
}
|
||||
|
||||
impl SessionLogSink {
|
||||
/// Create a fresh sink with an empty mirror. Used before any entry
|
||||
/// has been written (deferred SessionStart) or as a placeholder in
|
||||
/// tests.
|
||||
pub fn new() -> Self {
|
||||
let (broadcast_tx, _) = broadcast::channel(BROADCAST_CAPACITY);
|
||||
Self {
|
||||
inner: Arc::new(SinkInner {
|
||||
mirror: StdMutex::new(Vec::new()),
|
||||
broadcast_tx,
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
/// Create a sink seeded with a prefix of entries already on disk.
|
||||
/// Used by restore / fork-at-restore code paths that materialise
|
||||
/// the existing log before the sink starts taking new commits.
|
||||
pub fn with_initial(entries: Vec<LogEntry>) -> Self {
|
||||
let (broadcast_tx, _) = broadcast::channel(BROADCAST_CAPACITY);
|
||||
Self {
|
||||
inner: Arc::new(SinkInner {
|
||||
mirror: StdMutex::new(entries),
|
||||
broadcast_tx,
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
/// Push `entry` to the mirror and broadcast it.
|
||||
///
|
||||
/// MUST be called only after the Pod has successfully persisted the
|
||||
/// entry to the underlying `Store` — disk write is the gate. Failed
|
||||
/// disk writes must not call `publish`.
|
||||
pub fn publish(&self, entry: LogEntry) {
|
||||
let mut mirror = self
|
||||
.inner
|
||||
.mirror
|
||||
.lock()
|
||||
.expect("session log mirror mutex poisoned");
|
||||
mirror.push(entry.clone());
|
||||
// SendError means there are zero subscribers; harmless. We hold
|
||||
// the mirror lock across `send` so that `subscribe_with_snapshot`
|
||||
// cannot observe an inconsistent (snapshot, receiver) pair.
|
||||
let _ = self.inner.broadcast_tx.send(entry);
|
||||
}
|
||||
|
||||
/// Atomically swap the mirror to `[initial]` and broadcast the new
|
||||
/// session-start entry. Used during compaction / fork: the new
|
||||
/// `LogEntry::SessionStart` is the first entry of the replacement
|
||||
/// session, and existing subscribers transition by replaying it
|
||||
/// like any other live entry.
|
||||
///
|
||||
/// Existing snapshot prefixes seen by old subscribers stay valid
|
||||
/// for the prior session; the new `SessionStart` on the broadcast
|
||||
/// is the signal to reset their derived view.
|
||||
pub fn reset_with_initial(&self, initial: LogEntry) {
|
||||
let mut mirror = self
|
||||
.inner
|
||||
.mirror
|
||||
.lock()
|
||||
.expect("session log mirror mutex poisoned");
|
||||
mirror.clear();
|
||||
mirror.push(initial.clone());
|
||||
let _ = self.inner.broadcast_tx.send(initial);
|
||||
}
|
||||
|
||||
/// Replace the mirror with the supplied prefix without broadcasting.
|
||||
///
|
||||
/// Used by restore paths that load a session's complete log into
|
||||
/// the mirror before any subscriber is connected. Callers that need
|
||||
/// to notify existing subscribers should use [`reset_with_initial`].
|
||||
pub fn replace_silent(&self, entries: Vec<LogEntry>) {
|
||||
let mut mirror = self
|
||||
.inner
|
||||
.mirror
|
||||
.lock()
|
||||
.expect("session log mirror mutex poisoned");
|
||||
*mirror = entries;
|
||||
}
|
||||
|
||||
/// Atomically read the current mirror and subscribe to subsequent
|
||||
/// commits. The returned snapshot and receiver split the entry
|
||||
/// timeline into a duplicate-free, gap-free prefix/suffix pair.
|
||||
pub fn subscribe_with_snapshot(&self) -> (Vec<LogEntry>, broadcast::Receiver<LogEntry>) {
|
||||
let mirror = self
|
||||
.inner
|
||||
.mirror
|
||||
.lock()
|
||||
.expect("session log mirror mutex poisoned");
|
||||
let snapshot = mirror.clone();
|
||||
let rx = self.inner.broadcast_tx.subscribe();
|
||||
(snapshot, rx)
|
||||
}
|
||||
|
||||
/// Current entry count. Useful for tests / diagnostics.
|
||||
pub fn len(&self) -> usize {
|
||||
self.inner
|
||||
.mirror
|
||||
.lock()
|
||||
.expect("session log mirror mutex poisoned")
|
||||
.len()
|
||||
}
|
||||
|
||||
/// Whether the mirror is empty.
|
||||
pub fn is_empty(&self) -> bool {
|
||||
self.len() == 0
|
||||
}
|
||||
}
|
||||
|
||||
impl Default for SessionLogSink {
|
||||
fn default() -> Self {
|
||||
Self::new()
|
||||
}
|
||||
}
|
||||
|
||||
/// Active session head for the Pod's persistent log: session id +
|
||||
/// last-committed entry hash. Replaces the previous `SessionHead`
|
||||
/// struct local to `Pod`; bundled here so the writer can hand a
|
||||
/// cloneable handle to background tasks (e.g. the per-item drain
|
||||
/// task spawned by the controller).
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SessionHeadState {
|
||||
pub session_id: SessionId,
|
||||
pub head_hash: Option<EntryHash>,
|
||||
}
|
||||
|
||||
/// Pod-side session-log writer.
|
||||
///
|
||||
/// Bundles the (1) persistent store, (2) the in-memory session-head
|
||||
/// state (id + hash), and (3) the broadcast sink. `append_entry`
|
||||
/// chains the hash on disk, advances the head, then publishes the
|
||||
/// entry through the sink — under a single async mutex so two writers
|
||||
/// cannot interleave the chain.
|
||||
///
|
||||
/// `Clone` is a cheap `Arc` clone. The Pod keeps one writer for its
|
||||
/// inline commits (UserInput, TurnEnd, Usage, RunCompleted/Errored,
|
||||
/// scope snapshots, metrics) and hands clones to background tasks
|
||||
/// (e.g. the controller's per-item history drain task).
|
||||
pub struct SessionLogWriter<St> {
|
||||
inner: Arc<WriterInner<St>>,
|
||||
}
|
||||
|
||||
struct WriterInner<St> {
|
||||
store: St,
|
||||
head: AsyncMutex<SessionHeadState>,
|
||||
sink: SessionLogSink,
|
||||
}
|
||||
|
||||
impl<St> Clone for SessionLogWriter<St> {
|
||||
fn clone(&self) -> Self {
|
||||
Self {
|
||||
inner: self.inner.clone(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<St> SessionLogWriter<St>
|
||||
where
|
||||
St: Store + Clone,
|
||||
{
|
||||
/// Create a writer for a fresh Pod with no entries on disk yet.
|
||||
/// `head_hash` is `None` until the first `append_entry` (typically
|
||||
/// the deferred `SessionStart` written by `ensure_session_head`).
|
||||
pub fn new(store: St, session_id: SessionId) -> Self {
|
||||
Self {
|
||||
inner: Arc::new(WriterInner {
|
||||
store,
|
||||
head: AsyncMutex::new(SessionHeadState {
|
||||
session_id,
|
||||
head_hash: None,
|
||||
}),
|
||||
sink: SessionLogSink::new(),
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
/// Create a writer seeded with a session already on disk. The
|
||||
/// mirror is populated with `mirror` (typically loaded via
|
||||
/// `Store::read_all`), and `head_hash` should be the hash of the
|
||||
/// last entry.
|
||||
pub fn restored(
|
||||
store: St,
|
||||
session_id: SessionId,
|
||||
head_hash: Option<EntryHash>,
|
||||
mirror: Vec<LogEntry>,
|
||||
) -> Self {
|
||||
Self {
|
||||
inner: Arc::new(WriterInner {
|
||||
store,
|
||||
head: AsyncMutex::new(SessionHeadState {
|
||||
session_id,
|
||||
head_hash,
|
||||
}),
|
||||
sink: SessionLogSink::with_initial(mirror),
|
||||
}),
|
||||
}
|
||||
}
|
||||
|
||||
/// Append `entry` to the log: disk write → in-memory mirror push →
|
||||
/// broadcast — atomic w.r.t. `subscribe_with_snapshot` callers.
|
||||
pub async fn append_entry(&self, entry: LogEntry) -> Result<EntryHash, StoreError> {
|
||||
let mut head = self.inner.head.lock().await;
|
||||
let hash = session_store::append_entry_with_hash(
|
||||
&self.inner.store,
|
||||
head.session_id,
|
||||
&mut head.head_hash,
|
||||
entry.clone(),
|
||||
)
|
||||
.await?;
|
||||
self.inner.sink.publish(entry);
|
||||
Ok(hash)
|
||||
}
|
||||
|
||||
/// Atomically swap to a new compacted session.
|
||||
///
|
||||
/// Creates the new session on disk with `initial` as its
|
||||
/// `SessionStart`, advances the head, and resets the sink mirror
|
||||
/// to `[initial]` while broadcasting the entry. Existing
|
||||
/// subscribers observe the swap as a freshly broadcast
|
||||
/// `SessionStart` (with `compacted_from` set), which is their
|
||||
/// signal to reset their derived view.
|
||||
pub async fn swap_session(
|
||||
&self,
|
||||
new_session_id: SessionId,
|
||||
initial: LogEntry,
|
||||
) -> Result<EntryHash, StoreError> {
|
||||
let hash = session_log::compute_hash(None, &initial);
|
||||
let hashed = HashedEntry {
|
||||
hash: hash.clone(),
|
||||
prev_hash: None,
|
||||
entry: initial.clone(),
|
||||
};
|
||||
self.inner
|
||||
.store
|
||||
.create_session(new_session_id, &[hashed])
|
||||
.await?;
|
||||
let mut head = self.inner.head.lock().await;
|
||||
head.session_id = new_session_id;
|
||||
head.head_hash = Some(hash.clone());
|
||||
self.inner.sink.reset_with_initial(initial);
|
||||
Ok(hash)
|
||||
}
|
||||
|
||||
/// If the store's head no longer matches our cached head, mint a
|
||||
/// fresh session that forks from the current state and switch to
|
||||
/// it. Returns `true` when a fork happened.
|
||||
pub async fn ensure_head_or_fork(
|
||||
&self,
|
||||
state: SessionStartState<'_>,
|
||||
) -> Result<bool, StoreError> {
|
||||
let mut head = self.inner.head.lock().await;
|
||||
let store_head = self.inner.store.read_head_hash(head.session_id).await?;
|
||||
if store_head == head.head_hash {
|
||||
return Ok(false);
|
||||
}
|
||||
let fork_id = session_store::new_session_id();
|
||||
let entry = LogEntry::SessionStart {
|
||||
ts: session_log::now_millis(),
|
||||
system_prompt: state.system_prompt.map(String::from),
|
||||
config: state.config.clone(),
|
||||
history: session_store::to_logged(state.history),
|
||||
forked_from: None,
|
||||
compacted_from: None,
|
||||
};
|
||||
let hash = session_log::compute_hash(None, &entry);
|
||||
let hashed = HashedEntry {
|
||||
hash: hash.clone(),
|
||||
prev_hash: None,
|
||||
entry: entry.clone(),
|
||||
};
|
||||
self.inner.store.create_session(fork_id, &[hashed]).await?;
|
||||
head.session_id = fork_id;
|
||||
head.head_hash = Some(hash);
|
||||
self.inner.sink.reset_with_initial(entry);
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
/// Cloneable handle to the broadcast sink. Used by the IPC layer
|
||||
/// for `subscribe_with_snapshot` and by tests that just want the
|
||||
/// non-write side.
|
||||
pub fn sink(&self) -> SessionLogSink {
|
||||
self.inner.sink.clone()
|
||||
}
|
||||
|
||||
/// Underlying store handle. Direct access is preserved for callers
|
||||
/// that read state (`read_all`, `read_head_hash`) without going
|
||||
/// through the writer's hash chain.
|
||||
pub fn store(&self) -> &St {
|
||||
&self.inner.store
|
||||
}
|
||||
|
||||
/// Cheap snapshot of the current session id.
|
||||
pub async fn current_session_id(&self) -> SessionId {
|
||||
self.inner.head.lock().await.session_id
|
||||
}
|
||||
|
||||
/// Cheap snapshot of the current head hash.
|
||||
pub async fn current_head_hash(&self) -> Option<EntryHash> {
|
||||
self.inner.head.lock().await.head_hash.clone()
|
||||
}
|
||||
|
||||
/// Direct lock on the head. Used by paths that need to coordinate
|
||||
/// custom writes with the hash chain (currently
|
||||
/// `session_metrics::record_metric`).
|
||||
pub async fn lock_head(&self) -> MutexGuard<'_, SessionHeadState> {
|
||||
self.inner.head.lock().await
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use llm_worker::llm_client::RequestConfig;
|
||||
use session_store::session_log::now_millis;
|
||||
|
||||
fn session_start() -> LogEntry {
|
||||
LogEntry::SessionStart {
|
||||
ts: now_millis(),
|
||||
system_prompt: None,
|
||||
config: RequestConfig::default(),
|
||||
history: vec![],
|
||||
forked_from: None,
|
||||
compacted_from: None,
|
||||
}
|
||||
}
|
||||
|
||||
fn turn_end(n: usize) -> LogEntry {
|
||||
LogEntry::TurnEnd {
|
||||
ts: now_millis(),
|
||||
turn_count: n,
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn publish_then_subscribe_returns_history_in_snapshot() {
|
||||
let sink = SessionLogSink::new();
|
||||
sink.publish(session_start());
|
||||
sink.publish(turn_end(1));
|
||||
|
||||
let (snapshot, mut rx) = sink.subscribe_with_snapshot();
|
||||
assert_eq!(snapshot.len(), 2);
|
||||
assert!(matches!(snapshot[0], LogEntry::SessionStart { .. }));
|
||||
assert!(matches!(snapshot[1], LogEntry::TurnEnd { turn_count: 1, .. }));
|
||||
assert!(rx.try_recv().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn subscribe_then_publish_delivers_live_entries() {
|
||||
let sink = SessionLogSink::new();
|
||||
sink.publish(session_start());
|
||||
|
||||
let (snapshot, mut rx) = sink.subscribe_with_snapshot();
|
||||
assert_eq!(snapshot.len(), 1);
|
||||
|
||||
sink.publish(turn_end(1));
|
||||
match rx.try_recv() {
|
||||
Ok(LogEntry::TurnEnd { turn_count: 1, .. }) => {}
|
||||
other => panic!("unexpected: {other:?}"),
|
||||
}
|
||||
assert!(rx.try_recv().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn snapshot_and_live_never_overlap() {
|
||||
let sink = SessionLogSink::new();
|
||||
sink.publish(session_start());
|
||||
let (snapshot, mut rx) = sink.subscribe_with_snapshot();
|
||||
sink.publish(turn_end(1));
|
||||
|
||||
assert_eq!(snapshot.len(), 1);
|
||||
match rx.try_recv() {
|
||||
Ok(LogEntry::TurnEnd { turn_count: 1, .. }) => {}
|
||||
other => panic!("unexpected: {other:?}"),
|
||||
}
|
||||
assert!(rx.try_recv().is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn reset_with_initial_clears_and_broadcasts() {
|
||||
let sink = SessionLogSink::new();
|
||||
sink.publish(session_start());
|
||||
sink.publish(turn_end(1));
|
||||
|
||||
let (_pre_snapshot, mut rx) = sink.subscribe_with_snapshot();
|
||||
sink.reset_with_initial(session_start());
|
||||
|
||||
match rx.try_recv() {
|
||||
Ok(LogEntry::SessionStart { .. }) => {}
|
||||
other => panic!("expected SessionStart broadcast, got {other:?}"),
|
||||
}
|
||||
|
||||
let (post_snapshot, _) = sink.subscribe_with_snapshot();
|
||||
assert_eq!(post_snapshot.len(), 1);
|
||||
assert!(matches!(post_snapshot[0], LogEntry::SessionStart { .. }));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn replace_silent_does_not_broadcast() {
|
||||
let sink = SessionLogSink::new();
|
||||
sink.publish(session_start());
|
||||
let (_pre_snapshot, mut rx) = sink.subscribe_with_snapshot();
|
||||
|
||||
sink.replace_silent(vec![session_start(), turn_end(1)]);
|
||||
|
||||
// No broadcast fired.
|
||||
assert!(rx.try_recv().is_err());
|
||||
let (post_snapshot, _) = sink.subscribe_with_snapshot();
|
||||
assert_eq!(post_snapshot.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn with_initial_seeds_the_mirror() {
|
||||
let sink = SessionLogSink::with_initial(vec![session_start(), turn_end(1)]);
|
||||
let (snapshot, _) = sink.subscribe_with_snapshot();
|
||||
assert_eq!(snapshot.len(), 2);
|
||||
}
|
||||
}
|
||||
@@ -1,7 +1,6 @@
|
||||
use std::sync::{OnceLock, RwLock};
|
||||
|
||||
use llm_worker::llm_client::types::Item;
|
||||
use protocol::{PodStatus, Segment};
|
||||
use protocol::PodStatus;
|
||||
use serde_json::json;
|
||||
use session_store::SessionId;
|
||||
|
||||
@@ -19,21 +18,20 @@ pub struct KnowledgeCandidate {
|
||||
|
||||
/// Shared state between PodController and runtime directory.
|
||||
///
|
||||
/// Controller updates this in-memory; RuntimeDir writes it to disk.
|
||||
/// Wrapped in `Arc` for sharing.
|
||||
/// Controller updates this in-memory; RuntimeDir writes the status
|
||||
/// snapshot to disk. Wrapped in `Arc` for sharing.
|
||||
///
|
||||
/// History and typed user-segment mirrors used to live here so the
|
||||
/// IPC layer could answer `Method::GetHistory`. Those reads now go
|
||||
/// directly through the session-log sink (`Event::Snapshot` +
|
||||
/// `Event::Entry`), so this struct holds only status, identity,
|
||||
/// greeting, and completion lookup hubs.
|
||||
pub struct PodSharedState {
|
||||
pub pod_name: String,
|
||||
pub session_id: SessionId,
|
||||
pub manifest_toml: String,
|
||||
pub greeting: protocol::Greeting,
|
||||
pub status: RwLock<PodStatus>,
|
||||
pub history: RwLock<Vec<Item>>,
|
||||
/// Typed user submissions in submit order. The K-th entry corresponds
|
||||
/// to the K-th `Item::user_message` in `history` (modulo seed history
|
||||
/// loaded from a pre-compaction `SessionStart.history`, whose original
|
||||
/// segments are not preserved). Surfaced via `Event::History` so
|
||||
/// clients can re-render typed atoms on session restore.
|
||||
pub user_segments: RwLock<Vec<Vec<Segment>>>,
|
||||
/// Pod-from-the-inside view of the filesystem. Set once in
|
||||
/// `PodController::start` after the `ScopedFs` is materialised, and
|
||||
/// read from the IPC server layer to answer `ListCompletions`
|
||||
@@ -58,8 +56,6 @@ impl PodSharedState {
|
||||
manifest_toml,
|
||||
greeting,
|
||||
status: RwLock::new(PodStatus::Idle),
|
||||
history: RwLock::new(Vec::new()),
|
||||
user_segments: RwLock::new(Vec::new()),
|
||||
fs_view: OnceLock::new(),
|
||||
workflows: OnceLock::new(),
|
||||
knowledge: OnceLock::new(),
|
||||
@@ -112,25 +108,6 @@ impl PodSharedState {
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
pub fn user_segments(&self) -> Vec<Vec<Segment>> {
|
||||
self.user_segments
|
||||
.read()
|
||||
.map(|s| s.clone())
|
||||
.unwrap_or_default()
|
||||
}
|
||||
|
||||
pub fn set_user_segments(&self, segments: Vec<Vec<Segment>>) {
|
||||
if let Ok(mut s) = self.user_segments.write() {
|
||||
*s = segments;
|
||||
}
|
||||
}
|
||||
|
||||
pub fn push_user_segments(&self, segments: Vec<Segment>) {
|
||||
if let Ok(mut s) = self.user_segments.write() {
|
||||
s.push(segments);
|
||||
}
|
||||
}
|
||||
|
||||
pub fn set_status(&self, status: PodStatus) {
|
||||
if let Ok(mut s) = self.status.write() {
|
||||
*s = status;
|
||||
@@ -141,16 +118,6 @@ impl PodSharedState {
|
||||
self.status.read().map(|s| *s).unwrap_or(PodStatus::Idle)
|
||||
}
|
||||
|
||||
pub fn history(&self) -> Vec<Item> {
|
||||
self.history.read().map(|h| h.clone()).unwrap_or_default()
|
||||
}
|
||||
|
||||
pub fn update_history(&self, items: Vec<Item>) {
|
||||
if let Ok(mut h) = self.history.write() {
|
||||
*h = items;
|
||||
}
|
||||
}
|
||||
|
||||
/// Serialize status as JSON.
|
||||
pub fn status_json(&self) -> String {
|
||||
let status = self.get_status();
|
||||
@@ -161,21 +128,11 @@ impl PodSharedState {
|
||||
})
|
||||
.to_string()
|
||||
}
|
||||
|
||||
/// Serialize history as JSON.
|
||||
pub fn history_json(&self) -> String {
|
||||
if let Ok(h) = self.history.read() {
|
||||
serde_json::to_string(&*h).unwrap_or_else(|_| "[]".into())
|
||||
} else {
|
||||
"[]".into()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use llm_worker::llm_client::types::{ContentPart, Item, Role};
|
||||
|
||||
fn test_state() -> PodSharedState {
|
||||
PodSharedState::new(
|
||||
@@ -231,29 +188,6 @@ mod tests {
|
||||
assert_eq!(parsed["state"], "running");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn history_json_empty_initially() {
|
||||
let state = test_state();
|
||||
assert_eq!(state.history_json(), "[]");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn history_json_after_update() {
|
||||
let state = test_state();
|
||||
let items = vec![Item::Message {
|
||||
id: None,
|
||||
role: Role::Assistant,
|
||||
content: vec![ContentPart::Text {
|
||||
text: "Hello".into(),
|
||||
}],
|
||||
status: None,
|
||||
}];
|
||||
state.update_history(items);
|
||||
let json = state.history_json();
|
||||
let parsed: serde_json::Value = serde_json::from_str(&json).unwrap();
|
||||
assert!(parsed.is_array());
|
||||
assert_eq!(parsed[0]["role"], "assistant");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn knowledge_completions_empty_when_unset() {
|
||||
|
||||
@@ -19,6 +19,7 @@ use llm_worker::tool::{Tool, ToolDefinition, ToolError, ToolMeta, ToolOutput};
|
||||
use protocol::stream::{JsonLineReader, JsonLineWriter};
|
||||
use protocol::{ErrorCode, Event, Method};
|
||||
use serde::Deserialize;
|
||||
use session_store::LogEntry;
|
||||
use tokio::net::UnixStream;
|
||||
|
||||
use crate::runtime::dir::SpawnedPodRecord;
|
||||
@@ -385,33 +386,31 @@ async fn send_run_and_confirm(socket: &Path, input: String) -> Result<(), SendRu
|
||||
}
|
||||
}
|
||||
|
||||
/// Connect and ask the Pod for its conversation history. Skips
|
||||
/// pre-History events (such as buffered alerts replayed to new
|
||||
/// clients). Returns the raw JSON items as `serde_json::Value` since
|
||||
/// the pod crate already round-trips via `Value` on the wire.
|
||||
/// Connect to a Pod's socket and read the connect-time `Event::Snapshot`.
|
||||
///
|
||||
/// Pods deliver the session-log mirror as the first non-Alert event on
|
||||
/// every new connection, so consuming it is sufficient — no explicit
|
||||
/// `GetHistory` method round trip. Returns the entries as raw JSON
|
||||
/// values; callers deserialize as `session_store::LogEntry` if they
|
||||
/// need typed access.
|
||||
async fn fetch_history(socket: &Path) -> std::io::Result<Vec<serde_json::Value>> {
|
||||
let stream = tokio::time::timeout(SOCKET_OP_TIMEOUT, UnixStream::connect(socket))
|
||||
.await
|
||||
.map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "connect timed out"))??;
|
||||
let (r, w) = stream.into_split();
|
||||
let mut writer = JsonLineWriter::new(w);
|
||||
let (r, _w) = stream.into_split();
|
||||
let mut reader = JsonLineReader::new(r);
|
||||
|
||||
tokio::time::timeout(SOCKET_OP_TIMEOUT, writer.write(&Method::GetHistory))
|
||||
.await
|
||||
.map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "write timed out"))??;
|
||||
|
||||
loop {
|
||||
let event = tokio::time::timeout(SOCKET_OP_TIMEOUT, reader.next::<Event>())
|
||||
.await
|
||||
.map_err(|_| std::io::Error::new(std::io::ErrorKind::TimedOut, "read timed out"))??;
|
||||
match event {
|
||||
Some(Event::History { items, .. }) => return Ok(items),
|
||||
Some(Event::Snapshot { entries, .. }) => return Ok(entries),
|
||||
Some(_) => continue,
|
||||
None => {
|
||||
return Err(std::io::Error::new(
|
||||
std::io::ErrorKind::UnexpectedEof,
|
||||
"pod closed connection before History event",
|
||||
"pod closed connection before Snapshot event",
|
||||
));
|
||||
}
|
||||
}
|
||||
@@ -426,24 +425,36 @@ async fn is_reachable(socket: &Path) -> bool {
|
||||
.unwrap_or(false)
|
||||
}
|
||||
|
||||
fn extract_assistant_text(items: &[serde_json::Value]) -> String {
|
||||
fn extract_assistant_text(entries: &[serde_json::Value]) -> String {
|
||||
let mut out = String::new();
|
||||
for value in items {
|
||||
let Ok(item) = serde_json::from_value::<Item>(value.clone()) else {
|
||||
for value in entries {
|
||||
// The wire payload is the JSON form of `session_store::LogEntry`.
|
||||
// Walk Assistant items inside each entry that can carry them:
|
||||
// post-compaction `SessionStart.history` (seed) and per-LLM-call
|
||||
// `AssistantItems` deltas.
|
||||
let Ok(entry) = serde_json::from_value::<LogEntry>(value.clone()) else {
|
||||
continue;
|
||||
};
|
||||
if let Item::Message {
|
||||
role: Role::Assistant,
|
||||
content,
|
||||
..
|
||||
} = item
|
||||
{
|
||||
for part in content {
|
||||
if let ContentPart::Text { text } = part {
|
||||
if !out.is_empty() {
|
||||
out.push_str("\n\n");
|
||||
let logged_items = match entry {
|
||||
LogEntry::SessionStart { history, .. } => history,
|
||||
LogEntry::AssistantItems { items, .. } => items,
|
||||
_ => continue,
|
||||
};
|
||||
for logged in logged_items {
|
||||
let item: Item = logged.into();
|
||||
if let Item::Message {
|
||||
role: Role::Assistant,
|
||||
content,
|
||||
..
|
||||
} = item
|
||||
{
|
||||
for part in content {
|
||||
if let ContentPart::Text { text } = part {
|
||||
if !out.is_empty() {
|
||||
out.push_str("\n\n");
|
||||
}
|
||||
out.push_str(&text);
|
||||
}
|
||||
out.push_str(&text);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user