resumeの実装
This commit is contained in:
+17
-1
@@ -4,7 +4,7 @@ use std::process::ExitCode;
|
||||
use clap::Parser;
|
||||
use manifest::paths;
|
||||
use pod::{Pod, PodController, PodFactory};
|
||||
use session_store::FsStore;
|
||||
use session_store::{FsStore, SessionId};
|
||||
|
||||
#[derive(Parser)]
|
||||
#[command(
|
||||
@@ -43,6 +43,14 @@ struct Cli {
|
||||
/// callbacks upward. Required alongside `--adopt`.
|
||||
#[arg(long, value_name = "PATH", requires = "adopt")]
|
||||
callback: Option<PathBuf>,
|
||||
|
||||
/// Restore a Pod from an existing session. The source session log
|
||||
/// is forked at its head into a new session id, so the original
|
||||
/// jsonl is left untouched and double-write races are impossible.
|
||||
/// Mutually exclusive with `--adopt` (spawned children always start
|
||||
/// fresh).
|
||||
#[arg(long, value_name = "UUID", conflicts_with = "adopt")]
|
||||
session: Option<SessionId>,
|
||||
}
|
||||
|
||||
async fn build_factory(cli: &Cli) -> Result<PodFactory, String> {
|
||||
@@ -136,6 +144,14 @@ async fn main() -> ExitCode {
|
||||
return ExitCode::FAILURE;
|
||||
}
|
||||
}
|
||||
} else if let Some(source_session_id) = cli.session {
|
||||
match Pod::restore_from_manifest(source_session_id, manifest, store, loader).await {
|
||||
Ok(p) => p,
|
||||
Err(e) => {
|
||||
eprintln!("error: failed to restore pod: {e}");
|
||||
return ExitCode::FAILURE;
|
||||
}
|
||||
}
|
||||
} else {
|
||||
match Pod::from_manifest(manifest, store, loader).await {
|
||||
Ok(p) => p,
|
||||
|
||||
+225
-119
@@ -210,75 +210,11 @@ impl<C: LlmClient, St: Store> Pod<C, St> {
|
||||
self.inject_resident_knowledge = enabled;
|
||||
}
|
||||
|
||||
/// Restore a Pod from a persisted session.
|
||||
/// Shared handle to the prompt catalog. Cheap to clone (`Arc`).
|
||||
pub fn prompts(&self) -> &Arc<PromptCatalog> {
|
||||
&self.prompts
|
||||
}
|
||||
|
||||
pub async fn restore(
|
||||
session_id: SessionId,
|
||||
manifest: PodManifest,
|
||||
client: C,
|
||||
store: St,
|
||||
pwd: PathBuf,
|
||||
scope: Scope,
|
||||
) -> Result<Self, PodError> {
|
||||
let state = session_store::restore(&store, session_id).await?;
|
||||
let mut worker = Worker::new(client);
|
||||
if let Some(ref prompt) = state.system_prompt {
|
||||
worker.set_system_prompt(prompt);
|
||||
}
|
||||
// A leading `Role::System` item can only come from `compact`
|
||||
// (the Pod's one and only write path that prepends a summary at
|
||||
// history[0]). Restoring the anchor lets Anthropic re-use a
|
||||
// stable cache prefix for long-lived restored sessions.
|
||||
let anchored_on_summary = matches!(
|
||||
state.history.first(),
|
||||
Some(Item::Message {
|
||||
role: llm_worker::Role::System,
|
||||
..
|
||||
})
|
||||
);
|
||||
worker.set_history(state.history);
|
||||
worker.set_request_config(state.config);
|
||||
worker.set_turn_count(state.turn_count);
|
||||
worker.set_last_run_interrupted(state.last_run_interrupted);
|
||||
if anchored_on_summary {
|
||||
worker.set_cache_anchor(Some(0));
|
||||
}
|
||||
|
||||
let prompts = PromptCatalog::builtins_only()?;
|
||||
let extract_pointer = memory::extract::fold_pointer(&state.extensions);
|
||||
let mut pod = Self {
|
||||
manifest,
|
||||
worker: Some(worker),
|
||||
store,
|
||||
session_id,
|
||||
head_hash: state.head_hash,
|
||||
pwd,
|
||||
scope,
|
||||
hook_builder: HookRegistryBuilder::new(),
|
||||
interceptor_installed: false,
|
||||
compact_state: None,
|
||||
usage_tracker: Arc::new(UsageTracker::new()),
|
||||
usage_history: Arc::new(Mutex::new(state.usage_history)),
|
||||
tracker: None,
|
||||
system_prompt_template: None,
|
||||
alerter: None,
|
||||
event_tx: None,
|
||||
pending_notifies: NotifyBuffer::new(),
|
||||
scope_allocation: None,
|
||||
callback_socket: None,
|
||||
prompts,
|
||||
inject_resident_knowledge: true,
|
||||
extract_in_flight: Arc::new(AtomicBool::new(false)),
|
||||
extract_pointer: Mutex::new(extract_pointer),
|
||||
};
|
||||
pod.apply_prune_from_manifest();
|
||||
Ok(pod)
|
||||
}
|
||||
|
||||
/// The session ID used for persistence.
|
||||
pub fn session_id(&self) -> SessionId {
|
||||
self.session_id
|
||||
@@ -1534,15 +1470,18 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
store: St,
|
||||
loader: PromptLoader,
|
||||
) -> Result<Self, PodError> {
|
||||
let pwd = current_pwd()?;
|
||||
let scope = build_scope_with_memory(&manifest, &pwd)?;
|
||||
if !scope.is_readable(&pwd) {
|
||||
return Err(PodError::PwdOutsideScope { pwd });
|
||||
}
|
||||
let common = prepare_pod_common(&manifest, &loader, /* parse_template */ true)?;
|
||||
|
||||
// Session creation is deferred to the first run (see
|
||||
// `ensure_session_head`) so the SessionStart entry can capture
|
||||
// the rendered system prompt, not the raw template source. The
|
||||
// session_id is allocated here so the scope-lock registration
|
||||
// can record it from the start.
|
||||
let session_id = session_store::new_session_id();
|
||||
|
||||
// Register this Pod in the machine-wide scope-lock registry
|
||||
// before building anything else, so a spawn that conflicts on
|
||||
// scope fails fast (and without having paid for client setup).
|
||||
// scope fails fast.
|
||||
let socket_path = dir::default_base()
|
||||
.map_err(ScopeLockError::from)?
|
||||
.join(&manifest.pod.name)
|
||||
@@ -1551,50 +1490,34 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
manifest.pod.name.clone(),
|
||||
std::process::id(),
|
||||
socket_path,
|
||||
scope.allow_rules(),
|
||||
common.scope.allow_rules(),
|
||||
session_id,
|
||||
)?;
|
||||
|
||||
let client = provider::build_client(&manifest.model)?;
|
||||
let mut worker = Worker::new(client);
|
||||
let mut worker = Worker::new(common.client);
|
||||
apply_worker_manifest(&mut worker, &manifest.worker);
|
||||
|
||||
// Resolve the instruction reference and parse the resulting
|
||||
// template eagerly (syntax check only). Rendering is deferred
|
||||
// to `ensure_system_prompt_materialized` at first turn so
|
||||
// runtime values (date, tools, scope summary, ...) can be
|
||||
// injected.
|
||||
let system_prompt_template = Some(
|
||||
SystemPromptTemplate::parse(&manifest.worker.instruction, loader.clone())
|
||||
.map_err(|source| PodError::InvalidSystemPromptTemplate { source })?,
|
||||
);
|
||||
|
||||
let prompts = PromptCatalog::load(&loader, manifest.pod.prompt_pack.as_deref())?;
|
||||
|
||||
// Session creation is deferred to the first run (see
|
||||
// `ensure_session_head`) so the SessionStart entry can capture
|
||||
// the rendered system prompt, not the raw template source.
|
||||
let session_id = session_store::new_session_id();
|
||||
let mut pod = Self {
|
||||
manifest,
|
||||
worker: Some(worker),
|
||||
store,
|
||||
session_id,
|
||||
head_hash: None,
|
||||
pwd,
|
||||
scope,
|
||||
pwd: common.pwd,
|
||||
scope: common.scope,
|
||||
hook_builder: HookRegistryBuilder::new(),
|
||||
interceptor_installed: false,
|
||||
compact_state: None,
|
||||
usage_tracker: Arc::new(UsageTracker::new()),
|
||||
usage_history: Arc::new(Mutex::new(Vec::new())),
|
||||
tracker: None,
|
||||
system_prompt_template,
|
||||
system_prompt_template: common.system_prompt_template,
|
||||
alerter: None,
|
||||
event_tx: None,
|
||||
pending_notifies: NotifyBuffer::new(),
|
||||
scope_allocation: Some(scope_allocation),
|
||||
callback_socket: None,
|
||||
prompts,
|
||||
prompts: common.prompts,
|
||||
inject_resident_knowledge: true,
|
||||
extract_in_flight: Arc::new(AtomicBool::new(false)),
|
||||
extract_pointer: Mutex::new(None),
|
||||
@@ -1610,57 +1533,43 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
/// [`scope_lock::delegate_scope`], rather than installing a new
|
||||
/// top-level entry. `callback_socket` carries the spawner's
|
||||
/// Unix-socket path so the spawned Pod can send `Method::Notify`
|
||||
/// back to the spawner; it is stored but unused in the
|
||||
/// `spawn-pod-tool` ticket — the receiving side lands in the
|
||||
/// follow-up `pod-callback` ticket.
|
||||
/// back to the spawner.
|
||||
pub async fn from_manifest_spawned(
|
||||
manifest: PodManifest,
|
||||
store: St,
|
||||
loader: PromptLoader,
|
||||
callback_socket: PathBuf,
|
||||
) -> Result<Self, PodError> {
|
||||
let pwd = current_pwd()?;
|
||||
let scope = build_scope_with_memory(&manifest, &pwd)?;
|
||||
if !scope.is_readable(&pwd) {
|
||||
return Err(PodError::PwdOutsideScope { pwd });
|
||||
}
|
||||
|
||||
let scope_allocation =
|
||||
scope_lock::adopt_allocation(manifest.pod.name.clone(), std::process::id())?;
|
||||
|
||||
let client = provider::build_client(&manifest.model)?;
|
||||
let mut worker = Worker::new(client);
|
||||
apply_worker_manifest(&mut worker, &manifest.worker);
|
||||
|
||||
let system_prompt_template = Some(
|
||||
SystemPromptTemplate::parse(&manifest.worker.instruction, loader.clone())
|
||||
.map_err(|source| PodError::InvalidSystemPromptTemplate { source })?,
|
||||
);
|
||||
|
||||
let prompts = PromptCatalog::load(&loader, manifest.pod.prompt_pack.as_deref())?;
|
||||
let common = prepare_pod_common(&manifest, &loader, /* parse_template */ true)?;
|
||||
|
||||
let session_id = session_store::new_session_id();
|
||||
let scope_allocation =
|
||||
scope_lock::adopt_allocation(manifest.pod.name.clone(), std::process::id(), session_id)?;
|
||||
|
||||
let mut worker = Worker::new(common.client);
|
||||
apply_worker_manifest(&mut worker, &manifest.worker);
|
||||
|
||||
let mut pod = Self {
|
||||
manifest,
|
||||
worker: Some(worker),
|
||||
store,
|
||||
session_id,
|
||||
head_hash: None,
|
||||
pwd,
|
||||
scope,
|
||||
pwd: common.pwd,
|
||||
scope: common.scope,
|
||||
hook_builder: HookRegistryBuilder::new(),
|
||||
interceptor_installed: false,
|
||||
compact_state: None,
|
||||
usage_tracker: Arc::new(UsageTracker::new()),
|
||||
usage_history: Arc::new(Mutex::new(Vec::new())),
|
||||
tracker: None,
|
||||
system_prompt_template,
|
||||
system_prompt_template: common.system_prompt_template,
|
||||
alerter: None,
|
||||
event_tx: None,
|
||||
pending_notifies: NotifyBuffer::new(),
|
||||
scope_allocation: Some(scope_allocation),
|
||||
callback_socket: Some(callback_socket),
|
||||
prompts,
|
||||
prompts: common.prompts,
|
||||
inject_resident_knowledge: true,
|
||||
extract_in_flight: Arc::new(AtomicBool::new(false)),
|
||||
extract_pointer: Mutex::new(None),
|
||||
@@ -1669,6 +1578,136 @@ impl<St: Store> Pod<Box<dyn LlmClient>, St> {
|
||||
Ok(pod)
|
||||
}
|
||||
|
||||
/// Restore a Pod from an existing session log.
|
||||
///
|
||||
/// Resolves the manifest cascade exactly like [`Self::from_manifest`]
|
||||
/// (pwd / scope / scope-lock / client / prompt catalog), then forks
|
||||
/// the source session at its current head and seeds a fresh Worker
|
||||
/// from the resulting `RestoredState`. The Pod writes to the new
|
||||
/// fork session's jsonl; the source session's log is left intact.
|
||||
///
|
||||
/// Refuses to resume if another live Pod is currently writing to
|
||||
/// `source_session_id` (detected via `scope.lock`).
|
||||
///
|
||||
/// `system_prompt` is replayed verbatim from the session log —
|
||||
/// templates are not re-rendered on restore so a long-running
|
||||
/// session keeps a stable cache prefix even when the manifest's
|
||||
/// instruction template would render differently today.
|
||||
pub async fn restore_from_manifest(
|
||||
source_session_id: SessionId,
|
||||
manifest: PodManifest,
|
||||
store: St,
|
||||
loader: PromptLoader,
|
||||
) -> Result<Self, PodError> {
|
||||
// Refuse to resume into a session that's already being written.
|
||||
if let Some(info) = scope_lock::lookup_session(source_session_id)? {
|
||||
return Err(PodError::SessionInUse {
|
||||
session_id: source_session_id,
|
||||
pod_name: info.pod_name,
|
||||
socket: info.socket,
|
||||
});
|
||||
}
|
||||
|
||||
// Read the source state, then fork it into a fresh session id.
|
||||
// The fork's SessionStart captures the full history with
|
||||
// `forked_from` provenance pointing back to the source, so the
|
||||
// source jsonl stays untouched and double-write races are
|
||||
// impossible by construction.
|
||||
let state = session_store::restore(&store, source_session_id).await?;
|
||||
let Some(source_head) = state.head_hash.clone() else {
|
||||
return Err(PodError::SessionEmpty {
|
||||
session_id: source_session_id,
|
||||
});
|
||||
};
|
||||
let session_id = session_store::fork_at(&store, source_session_id, &source_head).await?;
|
||||
|
||||
let common = prepare_pod_common(&manifest, &loader, /* parse_template */ false)?;
|
||||
|
||||
let socket_path = dir::default_base()
|
||||
.map_err(ScopeLockError::from)?
|
||||
.join(&manifest.pod.name)
|
||||
.join("sock");
|
||||
let scope_allocation = scope_lock::install_top_level(
|
||||
manifest.pod.name.clone(),
|
||||
std::process::id(),
|
||||
socket_path,
|
||||
common.scope.allow_rules(),
|
||||
session_id,
|
||||
)?;
|
||||
|
||||
// Build the worker and apply the manifest defaults first, then
|
||||
// overwrite the pieces the session log is authoritative for.
|
||||
let mut worker = Worker::new(common.client);
|
||||
apply_worker_manifest(&mut worker, &manifest.worker);
|
||||
if let Some(ref prompt) = state.system_prompt {
|
||||
worker.set_system_prompt(prompt);
|
||||
}
|
||||
// A leading `Role::System` item can only come from `compact`
|
||||
// (the Pod's one and only write path that prepends a summary at
|
||||
// history[0]). Restoring the anchor lets Anthropic re-use a
|
||||
// stable cache prefix for long-lived restored sessions.
|
||||
let anchored_on_summary = matches!(
|
||||
state.history.first(),
|
||||
Some(Item::Message {
|
||||
role: llm_worker::Role::System,
|
||||
..
|
||||
})
|
||||
);
|
||||
worker.set_history(state.history.clone());
|
||||
worker.set_request_config(state.config.clone());
|
||||
worker.set_turn_count(state.turn_count);
|
||||
worker.set_last_run_interrupted(state.last_run_interrupted);
|
||||
if anchored_on_summary {
|
||||
worker.set_cache_anchor(Some(0));
|
||||
}
|
||||
|
||||
let extract_pointer = memory::extract::fold_pointer(&state.extensions);
|
||||
|
||||
// The fork's SessionStart hash is the new head. We could
|
||||
// recompute it by reading the new session log, but
|
||||
// `session_store::fork_at` already returns the new session_id
|
||||
// and we know the chain starts fresh. The next `save_delta`
|
||||
// call will read head from store before appending, so leaving
|
||||
// `head_hash = None` here is safe but less efficient — we
|
||||
// refresh from the store to avoid a chain refresh on first
|
||||
// append.
|
||||
let head_hash = store
|
||||
.read_head_hash(session_id)
|
||||
.await
|
||||
.ok()
|
||||
.flatten();
|
||||
|
||||
let mut pod = Self {
|
||||
manifest,
|
||||
worker: Some(worker),
|
||||
store,
|
||||
session_id,
|
||||
head_hash,
|
||||
pwd: common.pwd,
|
||||
scope: common.scope,
|
||||
hook_builder: HookRegistryBuilder::new(),
|
||||
interceptor_installed: false,
|
||||
compact_state: None,
|
||||
usage_tracker: Arc::new(UsageTracker::new()),
|
||||
usage_history: Arc::new(Mutex::new(state.usage_history)),
|
||||
tracker: None,
|
||||
// Restore replays the saved system_prompt verbatim — no
|
||||
// template re-render on resume.
|
||||
system_prompt_template: None,
|
||||
alerter: None,
|
||||
event_tx: None,
|
||||
pending_notifies: NotifyBuffer::new(),
|
||||
scope_allocation: Some(scope_allocation),
|
||||
callback_socket: None,
|
||||
prompts: common.prompts,
|
||||
inject_resident_knowledge: true,
|
||||
extract_in_flight: Arc::new(AtomicBool::new(false)),
|
||||
extract_pointer: Mutex::new(extract_pointer),
|
||||
};
|
||||
pod.apply_prune_from_manifest();
|
||||
Ok(pod)
|
||||
}
|
||||
|
||||
/// Convenience: build a Pod from a single-layer TOML manifest string.
|
||||
///
|
||||
/// Parses the TOML into a [`PodManifestConfig`], converts to a
|
||||
@@ -1862,6 +1901,73 @@ pub enum PodError {
|
||||
|
||||
#[error("memory Phase 1 staging write failed: {0}")]
|
||||
ExtractStaging(#[source] memory::extract::StagingError),
|
||||
|
||||
#[error(
|
||||
"session {session_id} is currently in use by pod `{pod_name}` at {}",
|
||||
.socket.display()
|
||||
)]
|
||||
SessionInUse {
|
||||
session_id: SessionId,
|
||||
pod_name: String,
|
||||
socket: PathBuf,
|
||||
},
|
||||
|
||||
#[error("session {session_id} has no entries to restore")]
|
||||
SessionEmpty { session_id: SessionId },
|
||||
}
|
||||
|
||||
/// Bundle of resources that every high-level Pod constructor needs:
|
||||
/// pwd, scope, an LLM client, the prompt catalog, and (optionally) a
|
||||
/// parsed system-prompt template. Built once by [`prepare_pod_common`]
|
||||
/// from the manifest cascade and then split into Pod fields.
|
||||
struct PodCommon {
|
||||
pwd: PathBuf,
|
||||
scope: Scope,
|
||||
client: Box<dyn LlmClient>,
|
||||
prompts: Arc<PromptCatalog>,
|
||||
system_prompt_template: Option<SystemPromptTemplate>,
|
||||
}
|
||||
|
||||
/// Resolve pwd / scope / LLM client / prompt catalog from a validated
|
||||
/// manifest cascade. Used by `from_manifest`, `from_manifest_spawned`,
|
||||
/// and `restore_from_manifest` so they share one definition of "what
|
||||
/// pieces fall out of a manifest".
|
||||
///
|
||||
/// `parse_template` controls whether the manifest's instruction is
|
||||
/// parsed as a system-prompt template. New Pods always parse so the
|
||||
/// template is rendered at first turn; restored Pods skip parsing
|
||||
/// because the saved session log replays a previously-rendered
|
||||
/// `system_prompt` verbatim.
|
||||
fn prepare_pod_common(
|
||||
manifest: &PodManifest,
|
||||
loader: &PromptLoader,
|
||||
parse_template: bool,
|
||||
) -> Result<PodCommon, PodError> {
|
||||
let pwd = current_pwd()?;
|
||||
let scope = build_scope_with_memory(manifest, &pwd)?;
|
||||
if !scope.is_readable(&pwd) {
|
||||
return Err(PodError::PwdOutsideScope { pwd });
|
||||
}
|
||||
|
||||
let client = provider::build_client(&manifest.model)?;
|
||||
let prompts = PromptCatalog::load(loader, manifest.pod.prompt_pack.as_deref())?;
|
||||
|
||||
let system_prompt_template = if parse_template {
|
||||
Some(
|
||||
SystemPromptTemplate::parse(&manifest.worker.instruction, loader.clone())
|
||||
.map_err(|source| PodError::InvalidSystemPromptTemplate { source })?,
|
||||
)
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
Ok(PodCommon {
|
||||
pwd,
|
||||
scope,
|
||||
client,
|
||||
prompts,
|
||||
system_prompt_template,
|
||||
})
|
||||
}
|
||||
|
||||
/// Build the Pod's runtime [`Scope`] from the manifest, layering the
|
||||
|
||||
@@ -21,6 +21,7 @@ use std::path::{Path, PathBuf};
|
||||
use fs4::fs_std::FileExt;
|
||||
use manifest::{Permission, ScopeRule, paths};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use session_store::SessionId;
|
||||
|
||||
/// On-disk representation of the allocation table.
|
||||
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
|
||||
@@ -50,6 +51,12 @@ pub struct Allocation {
|
||||
/// Name of the Pod that delegated scope to this one, or `None` for
|
||||
/// a top-level Pod started directly by a human.
|
||||
pub delegated_from: Option<String>,
|
||||
/// Session ID this Pod is currently writing to. `None` means this
|
||||
/// is a pre-reservation made by a spawner via [`delegate_scope`]
|
||||
/// before the child has come up; the child fills it in at
|
||||
/// [`adopt_allocation`] time.
|
||||
#[serde(default)]
|
||||
pub session_id: Option<SessionId>,
|
||||
}
|
||||
|
||||
impl LockFile {
|
||||
@@ -60,6 +67,14 @@ impl LockFile {
|
||||
pub fn find_mut(&mut self, pod_name: &str) -> Option<&mut Allocation> {
|
||||
self.allocations.iter_mut().find(|a| a.pod_name == pod_name)
|
||||
}
|
||||
|
||||
/// Find the allocation currently writing to `session_id`. Skips
|
||||
/// pre-reservations whose `session_id` is still `None`.
|
||||
pub fn find_by_session(&self, session_id: SessionId) -> Option<&Allocation> {
|
||||
self.allocations
|
||||
.iter()
|
||||
.find(|a| a.session_id == Some(session_id))
|
||||
}
|
||||
}
|
||||
|
||||
/// Default on-disk path: `<runtime_dir>/scope.lock` resolved via
|
||||
@@ -294,6 +309,7 @@ pub fn register_pod(
|
||||
pid: u32,
|
||||
socket: PathBuf,
|
||||
scope_allow: Vec<ScopeRule>,
|
||||
session_id: SessionId,
|
||||
) -> Result<(), ScopeLockError> {
|
||||
reclaim_stale(guard);
|
||||
if guard.data().find(&pod_name).is_some() {
|
||||
@@ -316,6 +332,7 @@ pub fn register_pod(
|
||||
socket,
|
||||
scope_allow,
|
||||
delegated_from: None,
|
||||
session_id: Some(session_id),
|
||||
});
|
||||
guard.save()?;
|
||||
Ok(())
|
||||
@@ -361,6 +378,9 @@ pub fn delegate_scope(
|
||||
socket,
|
||||
scope_allow,
|
||||
delegated_from: Some(spawner.into()),
|
||||
// Pre-reservation. The child fills in its own session_id when
|
||||
// it calls `adopt_allocation` after the worker is built.
|
||||
session_id: None,
|
||||
});
|
||||
guard.save()?;
|
||||
Ok(())
|
||||
@@ -483,10 +503,18 @@ pub fn install_top_level(
|
||||
pid: u32,
|
||||
socket: PathBuf,
|
||||
scope_allow: Vec<ScopeRule>,
|
||||
session_id: SessionId,
|
||||
) -> Result<ScopeAllocationGuard, ScopeLockError> {
|
||||
let lock_path = default_lock_path()?;
|
||||
let mut guard = LockFileGuard::open(&lock_path)?;
|
||||
register_pod(&mut guard, pod_name.clone(), pid, socket, scope_allow)?;
|
||||
register_pod(
|
||||
&mut guard,
|
||||
pod_name.clone(),
|
||||
pid,
|
||||
socket,
|
||||
scope_allow,
|
||||
session_id,
|
||||
)?;
|
||||
Ok(ScopeAllocationGuard {
|
||||
pod_name,
|
||||
lock_path,
|
||||
@@ -497,13 +525,15 @@ pub fn install_top_level(
|
||||
/// a spawning Pod.
|
||||
///
|
||||
/// The spawning flow is two-stage: the spawner calls [`delegate_scope`]
|
||||
/// (with its own pid as a live placeholder), then exec's the child; the
|
||||
/// child, once running, calls this function to rewrite the allocation's
|
||||
/// pid to its own and claim the `ScopeAllocationGuard` so the entry is
|
||||
/// released when the child exits.
|
||||
/// (with its own pid as a live placeholder, `session_id = None`), then
|
||||
/// exec's the child; the child, once running, calls this function to
|
||||
/// rewrite the allocation's pid + session_id to its own and claim the
|
||||
/// `ScopeAllocationGuard` so the entry is released when the child
|
||||
/// exits.
|
||||
pub fn adopt_allocation(
|
||||
pod_name: String,
|
||||
new_pid: u32,
|
||||
session_id: SessionId,
|
||||
) -> Result<ScopeAllocationGuard, ScopeLockError> {
|
||||
let lock_path = default_lock_path()?;
|
||||
let mut guard = LockFileGuard::open(&lock_path)?;
|
||||
@@ -512,6 +542,7 @@ pub fn adopt_allocation(
|
||||
.find_mut(&pod_name)
|
||||
.ok_or_else(|| ScopeLockError::UnknownPod(pod_name.clone()))?;
|
||||
alloc.pid = new_pid;
|
||||
alloc.session_id = Some(session_id);
|
||||
guard.save()?;
|
||||
Ok(ScopeAllocationGuard {
|
||||
pod_name,
|
||||
@@ -519,6 +550,33 @@ pub fn adopt_allocation(
|
||||
})
|
||||
}
|
||||
|
||||
/// Information about a Pod that currently holds an allocation for a
|
||||
/// given session.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct SessionLockInfo {
|
||||
pub pod_name: String,
|
||||
pub socket: PathBuf,
|
||||
pub pid: u32,
|
||||
}
|
||||
|
||||
/// Open the default lock file, reclaim stale entries, and return the
|
||||
/// allocation currently writing to `session_id`, if any.
|
||||
///
|
||||
/// Used by `Pod::restore_from_manifest` to refuse a resume that would
|
||||
/// race a live writer on the same source session.
|
||||
pub fn lookup_session(session_id: SessionId) -> Result<Option<SessionLockInfo>, ScopeLockError> {
|
||||
let lock_path = default_lock_path()?;
|
||||
let mut guard = LockFileGuard::open(&lock_path)?;
|
||||
reclaim_stale(&mut guard);
|
||||
Ok(guard.data().find_by_session(session_id).map(|a| {
|
||||
SessionLockInfo {
|
||||
pod_name: a.pod_name.clone(),
|
||||
socket: a.socket.clone(),
|
||||
pid: a.pid,
|
||||
}
|
||||
}))
|
||||
}
|
||||
|
||||
/// Errors raised by the mutating scope-lock operations.
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum ScopeLockError {
|
||||
@@ -548,6 +606,10 @@ mod tests {
|
||||
/// harness runs tests on multiple threads inside a single process,
|
||||
/// so env-var writes from one test would otherwise leak into a
|
||||
/// parallel test's `default_lock_path()` lookup.
|
||||
fn sid() -> SessionId {
|
||||
session_store::new_session_id()
|
||||
}
|
||||
|
||||
static ENV_LOCK: LazyLock<Mutex<()>> = LazyLock::new(|| Mutex::new(()));
|
||||
|
||||
/// Sandbox `INSOMNIA_RUNTIME_DIR` to a tempdir for the duration of
|
||||
@@ -652,6 +714,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
@@ -699,6 +762,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
let err = register_pod(
|
||||
@@ -707,6 +771,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("b"),
|
||||
vec![write_rule("/src/core", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap_err();
|
||||
match err {
|
||||
@@ -726,6 +791,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
let err = register_pod(
|
||||
@@ -734,6 +800,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a2"),
|
||||
vec![write_rule("/docs", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap_err();
|
||||
assert!(matches!(err, ScopeLockError::DuplicatePodName(ref n) if n == "a"));
|
||||
@@ -750,6 +817,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
let err = delegate_scope(
|
||||
@@ -775,6 +843,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
delegate_scope(
|
||||
@@ -812,6 +881,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
delegate_scope(
|
||||
@@ -848,6 +918,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
delegate_scope(
|
||||
@@ -886,6 +957,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
delegate_scope(
|
||||
@@ -939,6 +1011,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
// B only reads under the same tree — allowed.
|
||||
@@ -948,6 +1021,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("b"),
|
||||
vec![read_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(g.data().allocations.len(), 2);
|
||||
@@ -964,6 +1038,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
release_pod(&mut g, "a").unwrap();
|
||||
@@ -973,6 +1048,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("b"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
@@ -988,6 +1064,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
delegate_scope(
|
||||
@@ -1023,6 +1100,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
{
|
||||
@@ -1047,7 +1125,7 @@ mod tests {
|
||||
delegate_placeholder(&mut g, "child", std::process::id());
|
||||
}
|
||||
let child_pid = std::process::id().wrapping_add(1);
|
||||
let guard = adopt_allocation("child".into(), child_pid).unwrap();
|
||||
let guard = adopt_allocation("child".into(), child_pid, sid()).unwrap();
|
||||
{
|
||||
let g = LockFileGuard::open(&lock_path).unwrap();
|
||||
let alloc = g.data().find("child").unwrap();
|
||||
@@ -1064,7 +1142,7 @@ mod tests {
|
||||
fn adopt_allocation_errors_on_unknown_pod() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let _sandbox = RuntimeDirSandbox::new(dir.path());
|
||||
let err = adopt_allocation("ghost".into(), 42).unwrap_err();
|
||||
let err = adopt_allocation("ghost".into(), 42, sid()).unwrap_err();
|
||||
assert!(matches!(err, ScopeLockError::UnknownPod(ref n) if n == "ghost"));
|
||||
}
|
||||
|
||||
@@ -1078,6 +1156,7 @@ mod tests {
|
||||
socket: sock(pod_name),
|
||||
scope_allow: vec![write_rule("/tmp/child", true)],
|
||||
delegated_from: None,
|
||||
session_id: None,
|
||||
});
|
||||
g.save().unwrap();
|
||||
}
|
||||
@@ -1093,6 +1172,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("a"),
|
||||
vec![write_rule("/src", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap();
|
||||
delegate_scope(
|
||||
@@ -1112,6 +1192,7 @@ mod tests {
|
||||
std::process::id(),
|
||||
sock("x"),
|
||||
vec![write_rule("/src/core/x", true)],
|
||||
sid(),
|
||||
)
|
||||
.unwrap_err();
|
||||
match err {
|
||||
|
||||
@@ -351,6 +351,7 @@ async fn stop_pod_sends_shutdown_and_releases_scope() {
|
||||
permission: Permission::Write,
|
||||
recursive: true,
|
||||
}],
|
||||
session_store::new_session_id(),
|
||||
)
|
||||
.unwrap();
|
||||
scope_lock::delegate_scope(
|
||||
|
||||
@@ -358,6 +358,7 @@ async fn shutdown_releases_scope_allocation_when_present() {
|
||||
std::process::id(),
|
||||
"/tmp/kid.sock".into(),
|
||||
vec![],
|
||||
session_store::new_session_id(),
|
||||
)
|
||||
.unwrap();
|
||||
std::mem::forget(guard);
|
||||
|
||||
@@ -73,6 +73,7 @@ async fn setup_spawner(
|
||||
permission: Permission::Write,
|
||||
recursive: true,
|
||||
}],
|
||||
session_store::new_session_id(),
|
||||
)
|
||||
.unwrap();
|
||||
// Leak the guard — the spawner allocation needs to outlive the
|
||||
|
||||
Reference in New Issue
Block a user