|
|
|
@@ -51,7 +51,7 @@ use tokio::task::JoinHandle;
|
|
|
|
|
/// without taking a mutex on the append hot path. `entries_written` is
|
|
|
|
|
/// an `AtomicUsize` bumped on every successful append; the writer's
|
|
|
|
|
/// tally is compared against the store's on-disk count to detect
|
|
|
|
|
/// concurrent writers in `ensure_session_head`.
|
|
|
|
|
/// concurrent writers in `ensure_segment_head`.
|
|
|
|
|
pub struct SegmentState {
|
|
|
|
|
segment_id: ArcSwap<SegmentId>,
|
|
|
|
|
entries_written: AtomicUsize,
|
|
|
|
@@ -69,7 +69,7 @@ impl SegmentState {
|
|
|
|
|
**self.segment_id.load()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn set_session_id(&self, id: SegmentId) {
|
|
|
|
|
pub fn set_segment_id(&self, id: SegmentId) {
|
|
|
|
|
self.segment_id.store(Arc::new(id));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -163,8 +163,8 @@ pub struct Pod<C: LlmClient, St: Store> {
|
|
|
|
|
store: St,
|
|
|
|
|
/// Shared session pointer. Source of truth for the Pod's current
|
|
|
|
|
/// `segment_id` and append tally. `self.segment_id()` is a thin
|
|
|
|
|
/// wrapper over `session_state.segment_id()`.
|
|
|
|
|
session_state: Arc<SegmentState>,
|
|
|
|
|
/// wrapper over `segment_state.segment_id()`.
|
|
|
|
|
segment_state: Arc<SegmentState>,
|
|
|
|
|
/// Absolute working directory of the Pod.
|
|
|
|
|
pwd: PathBuf,
|
|
|
|
|
/// Shared, atomically-swappable view of the Pod's resolved scope.
|
|
|
|
@@ -338,7 +338,7 @@ impl<C: LlmClient + Clone + 'static, St: Store + Clone + 'static> Pod<C, St> {
|
|
|
|
|
manifest: self.manifest.clone(),
|
|
|
|
|
worker: Some(worker),
|
|
|
|
|
store: self.store.clone(),
|
|
|
|
|
session_state: self.session_state.clone(),
|
|
|
|
|
segment_state: self.segment_state.clone(),
|
|
|
|
|
pwd: self.pwd.clone(),
|
|
|
|
|
scope: self.scope.clone(),
|
|
|
|
|
hook_builder: HookRegistryBuilder::new(),
|
|
|
|
@@ -382,7 +382,7 @@ impl<C: LlmClient + Clone + 'static, St: Store + Clone + 'static> Pod<C, St> {
|
|
|
|
|
pub fn log_writer_handle(&self) -> LogWriterHandle<St> {
|
|
|
|
|
LogWriterHandle {
|
|
|
|
|
store: self.store.clone(),
|
|
|
|
|
state: self.session_state.clone(),
|
|
|
|
|
state: self.segment_state.clone(),
|
|
|
|
|
sink: self.sink.clone(),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@@ -469,7 +469,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
pwd: PathBuf,
|
|
|
|
|
scope: Scope,
|
|
|
|
|
) -> Result<Self, PodError> {
|
|
|
|
|
// Segment creation is deferred to `ensure_session_head` at first
|
|
|
|
|
// Segment creation is deferred to `ensure_segment_head` at first
|
|
|
|
|
// run so a later-installed system-prompt template (see
|
|
|
|
|
// `set_system_prompt_template`) can be captured by `SegmentStart`.
|
|
|
|
|
let segment_id = session_store::new_segment_id();
|
|
|
|
@@ -478,7 +478,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
manifest,
|
|
|
|
|
worker: Some(worker),
|
|
|
|
|
store,
|
|
|
|
|
session_state: SegmentState::new(segment_id, 0),
|
|
|
|
|
segment_state: SegmentState::new(segment_id, 0),
|
|
|
|
|
pwd,
|
|
|
|
|
scope: SharedScope::new(scope),
|
|
|
|
|
hook_builder: HookRegistryBuilder::new(),
|
|
|
|
@@ -545,7 +545,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
/// The session ID used for persistence. Read lock-free from the
|
|
|
|
|
/// shared session pointer so fork-time swaps are observed immediately.
|
|
|
|
|
pub fn segment_id(&self) -> SegmentId {
|
|
|
|
|
self.session_state.segment_id()
|
|
|
|
|
self.segment_state.segment_id()
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// The Pod's manifest.
|
|
|
|
@@ -604,7 +604,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
/// can restore the narrowed scope instead of reclaiming delegated
|
|
|
|
|
/// writes.
|
|
|
|
|
pub fn persist_scope_snapshot(&mut self) -> Result<(), StoreError> {
|
|
|
|
|
if self.session_state.entries_written() == 0 {
|
|
|
|
|
if self.segment_state.entries_written() == 0 {
|
|
|
|
|
return Ok(());
|
|
|
|
|
}
|
|
|
|
|
let snapshot = {
|
|
|
|
@@ -627,9 +627,9 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
/// concurrent appenders — the kernel orders `O_APPEND` writes for
|
|
|
|
|
/// lines smaller than `PIPE_BUF`.
|
|
|
|
|
pub(crate) fn commit_entry(&self, entry: LogEntry) -> Result<(), StoreError> {
|
|
|
|
|
let segment_id = self.session_state.segment_id();
|
|
|
|
|
let segment_id = self.segment_state.segment_id();
|
|
|
|
|
self.store.append(segment_id, &entry)?;
|
|
|
|
|
self.session_state.increment_entries();
|
|
|
|
|
self.segment_state.increment_entries();
|
|
|
|
|
self.sink.publish(entry);
|
|
|
|
|
Ok(())
|
|
|
|
|
}
|
|
|
|
@@ -1143,7 +1143,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
self.ensure_interceptor_installed();
|
|
|
|
|
self.ensure_system_prompt_materialized()?;
|
|
|
|
|
self.cleanup_finished_memory_task();
|
|
|
|
|
self.ensure_session_head()?;
|
|
|
|
|
self.ensure_segment_head()?;
|
|
|
|
|
if self.should_pre_run_compact() {
|
|
|
|
|
self.join_memory_task().await;
|
|
|
|
|
}
|
|
|
|
@@ -1616,10 +1616,10 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
/// `ensure_system_prompt_materialized` has just rendered. Subsequent
|
|
|
|
|
/// calls fall through to entry-count comparison, which auto-forks
|
|
|
|
|
/// when another writer has appended behind our back.
|
|
|
|
|
fn ensure_session_head(&mut self) -> Result<(), PodError> {
|
|
|
|
|
fn ensure_segment_head(&mut self) -> Result<(), PodError> {
|
|
|
|
|
let w = self.worker.as_ref().unwrap();
|
|
|
|
|
let prev_session_id = self.session_state.segment_id();
|
|
|
|
|
let entries_written = self.session_state.entries_written();
|
|
|
|
|
let prev_segment_id = self.segment_state.segment_id();
|
|
|
|
|
let entries_written = self.segment_state.entries_written();
|
|
|
|
|
if entries_written == 0 {
|
|
|
|
|
let initial = LogEntry::SegmentStart {
|
|
|
|
|
ts: segment_log::now_millis(),
|
|
|
|
@@ -1636,7 +1636,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
// Check store count + auto-fork if it drifted.
|
|
|
|
|
let store_count = self
|
|
|
|
|
.store
|
|
|
|
|
.read_entry_count(prev_session_id)
|
|
|
|
|
.read_entry_count(prev_segment_id)
|
|
|
|
|
.map_err(PodError::from)?;
|
|
|
|
|
if store_count == entries_written {
|
|
|
|
|
return Ok(());
|
|
|
|
@@ -1656,8 +1656,8 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
self.store
|
|
|
|
|
.create_segment(fork_id, &[entry.clone()])
|
|
|
|
|
.map_err(PodError::from)?;
|
|
|
|
|
self.session_state.set_session_id(fork_id);
|
|
|
|
|
self.session_state.set_entries_written(1);
|
|
|
|
|
self.segment_state.set_segment_id(fork_id);
|
|
|
|
|
self.segment_state.set_entries_written(1);
|
|
|
|
|
self.sink.reset_with_initial(entry);
|
|
|
|
|
if self.scope_allocation.is_some() {
|
|
|
|
|
pod_registry::update_segment(&self.manifest.pod.name, fork_id)?;
|
|
|
|
@@ -2145,7 +2145,7 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
// the broadcast sink so existing subscribers see the new
|
|
|
|
|
// `SegmentStart { compacted_from }` and reset their view.
|
|
|
|
|
let new_segment_id = session_store::new_segment_id();
|
|
|
|
|
let old_session_id = self.session_state.segment_id();
|
|
|
|
|
let old_session_id = self.segment_state.segment_id();
|
|
|
|
|
let source_turn_count = self.worker.as_ref().unwrap().turn_count();
|
|
|
|
|
let w = self.worker.as_ref().unwrap();
|
|
|
|
|
let entry = LogEntry::SegmentStart {
|
|
|
|
@@ -2160,8 +2160,8 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
}),
|
|
|
|
|
};
|
|
|
|
|
self.store.create_segment(new_segment_id, &[entry.clone()])?;
|
|
|
|
|
self.session_state.set_session_id(new_segment_id);
|
|
|
|
|
self.session_state.set_entries_written(1);
|
|
|
|
|
self.segment_state.set_segment_id(new_segment_id);
|
|
|
|
|
self.segment_state.set_entries_written(1);
|
|
|
|
|
let session_start = entry;
|
|
|
|
|
// Broadcast the SegmentStart through the sink. This atomically
|
|
|
|
|
// resets the mirror to `[SegmentStart]` so any subscriber
|
|
|
|
@@ -2435,12 +2435,12 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
|
|
|
|
extract::ExtractedPayload::default()
|
|
|
|
|
});
|
|
|
|
|
|
|
|
|
|
let source_session_id = self.session_state.segment_id();
|
|
|
|
|
let source_segment_id = self.segment_state.segment_id();
|
|
|
|
|
let staging_id = if payload.is_empty() {
|
|
|
|
|
String::new()
|
|
|
|
|
} else {
|
|
|
|
|
let source = memory::schema::SourceRef {
|
|
|
|
|
segment_id: source_session_id.to_string(),
|
|
|
|
|
segment_id: source_segment_id.to_string(),
|
|
|
|
|
range: [start_entry as u64, end_entry as u64],
|
|
|
|
|
};
|
|
|
|
|
let (id, _) = extract::write_staging(&layout, source, payload)
|
|
|
|
@@ -2736,7 +2736,7 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
|
|
|
|
let skill_shadows = std::mem::take(&mut common.skill_shadows);
|
|
|
|
|
|
|
|
|
|
// Segment creation is deferred to the first run (see
|
|
|
|
|
// `ensure_session_head`) so the SegmentStart entry can capture
|
|
|
|
|
// `ensure_segment_head`) so the SegmentStart entry can capture
|
|
|
|
|
// the rendered system prompt, not the raw template source. The
|
|
|
|
|
// segment_id is allocated here so the pod-registry registration
|
|
|
|
|
// can record it from the start.
|
|
|
|
@@ -2765,7 +2765,7 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
|
|
|
|
manifest,
|
|
|
|
|
worker: Some(worker),
|
|
|
|
|
store,
|
|
|
|
|
session_state: SegmentState::new(segment_id, 0),
|
|
|
|
|
segment_state: SegmentState::new(segment_id, 0),
|
|
|
|
|
pwd: common.pwd,
|
|
|
|
|
scope: SharedScope::new(common.scope),
|
|
|
|
|
hook_builder: HookRegistryBuilder::new(),
|
|
|
|
@@ -2835,7 +2835,7 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
|
|
|
|
manifest,
|
|
|
|
|
worker: Some(worker),
|
|
|
|
|
store,
|
|
|
|
|
session_state: SegmentState::new(segment_id, 0),
|
|
|
|
|
segment_state: SegmentState::new(segment_id, 0),
|
|
|
|
|
pwd: common.pwd,
|
|
|
|
|
scope: SharedScope::new(common.scope),
|
|
|
|
|
hook_builder: HookRegistryBuilder::new(),
|
|
|
|
@@ -2903,13 +2903,13 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
|
|
|
|
let raw_entries = store.read_all(segment_id)?;
|
|
|
|
|
let state = session_store::collect_state(&raw_entries);
|
|
|
|
|
if state.entries_count == 0 {
|
|
|
|
|
return Err(PodError::SessionEmpty { segment_id });
|
|
|
|
|
return Err(PodError::SegmentEmpty { segment_id });
|
|
|
|
|
}
|
|
|
|
|
let mirror_entries: Vec<LogEntry> = raw_entries.clone();
|
|
|
|
|
let scope_snapshot = state
|
|
|
|
|
.pod_scope
|
|
|
|
|
.clone()
|
|
|
|
|
.ok_or(PodError::SessionScopeMissing { segment_id })?;
|
|
|
|
|
.ok_or(PodError::SegmentScopeMissing { segment_id })?;
|
|
|
|
|
|
|
|
|
|
let mut common = prepare_pod_common_with_scope(
|
|
|
|
|
&manifest,
|
|
|
|
@@ -2974,7 +2974,7 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
|
|
|
|
manifest,
|
|
|
|
|
worker: Some(worker),
|
|
|
|
|
store,
|
|
|
|
|
session_state: SegmentState::new(segment_id, state.entries_count),
|
|
|
|
|
segment_state: SegmentState::new(segment_id, state.entries_count),
|
|
|
|
|
pwd: common.pwd,
|
|
|
|
|
scope: SharedScope::new(common.scope),
|
|
|
|
|
hook_builder: HookRegistryBuilder::new(),
|
|
|
|
@@ -3235,12 +3235,12 @@ pub enum PodError {
|
|
|
|
|
WorkflowResolve(#[from] WorkflowResolveError),
|
|
|
|
|
|
|
|
|
|
#[error("session {segment_id} has no entries to restore")]
|
|
|
|
|
SessionEmpty { segment_id: SegmentId },
|
|
|
|
|
SegmentEmpty { segment_id: SegmentId },
|
|
|
|
|
|
|
|
|
|
#[error(
|
|
|
|
|
"session {segment_id} has no persisted scope snapshot; refusing resume without explicit scope"
|
|
|
|
|
)]
|
|
|
|
|
SessionScopeMissing { segment_id: SegmentId },
|
|
|
|
|
SegmentScopeMissing { segment_id: SegmentId },
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Bundle of resources that every high-level Pod constructor needs:
|
|
|
|
|