diff --git a/crates/tracedecay-cli/src/cli.rs b/crates/tracedecay-cli/src/cli.rs index 12b5da83b6..e04de1b441 100644 --- a/crates/tracedecay-cli/src/cli.rs +++ b/crates/tracedecay-cli/src/cli.rs @@ -658,8 +658,11 @@ pub enum Commands { #[command(long_about = WIPE_LONG_ABOUT, after_help = WIPE_AFTER_HELP)] Wipe { /// Wipe ALL profile database state while preserving identity and config - #[arg(short, long)] + #[arg(short, long, conflicts_with = "stale")] all: bool, + /// Delete only the stores the running daemon reports as requiring reset + #[arg(long)] + stale: bool, }, /// List tracedecay projects (current folder, parents, and children) #[command(long_about = LIST_LONG_ABOUT, after_help = LIST_AFTER_HELP)] diff --git a/crates/tracedecay-cli/src/cli/help.rs b/crates/tracedecay-cli/src/cli/help.rs index 0f3cc11949..a276104813 100644 --- a/crates/tracedecay-cli/src/cli/help.rs +++ b/crates/tracedecay-cli/src/cli/help.rs @@ -362,7 +362,8 @@ Exit status (a completed binary upgrade stays installed in every case): Code's `/plugins install`); act on the printed step, then rerun. Also 75 when the restored daemon serves a store whose persisted shape this binary does not open: it names the store and the exact reset command - (`tracedecay wipe --all --yes`), which nothing runs on your behalf + (`tracedecay wipe --stale --yes` for session stores), which nothing + runs on your behalf Related: tracedecay upgrade (refresh only after a real install), tracedecay update-plugin (plugins only), tracedecay channel."; @@ -653,8 +654,11 @@ pub(crate) const WIPE_LONG_ABOUT: &str = "\ Deletes .tracedecay stores (code graph, memory, sessions) for the current \ folder, its parents, and its children. With --all, deletes the complete \ profile-scoped database state, including global, user memory/session, project, \ -legacy, remote, Grafeo WAL, and host-admission stores. Profile identity, \ -configuration, and agent integration config remain untouched. \ +legacy, remote, Grafeo WAL, and host-admission stores. With --stale, deletes \ +exactly the stores the running daemon reports as requiring reset (a persisted \ +shape this binary does not open) and nothing else; the daemon recreates each \ +one empty. Profile identity, configuration, and agent integration config \ +remain untouched. \ Destructive and unrecoverable; re-create indexes with `tracedecay init`. \ Prompts for a `go!` confirmation unless `--yes` is passed. When the managed \ daemon holds the profile (even wedged or hung), wipe stops the installed \ @@ -666,6 +670,7 @@ Examples: tracedecay wipe Wipe stores around the cwd tracedecay wipe --all Wipe all profile database state tracedecay wipe --all --yes Confirm without the prompt + tracedecay wipe --stale --yes Reset only the stores the daemon refuses Related: tracedecay list (inspect nearby project stores), tracedecay init (re-index afterwards), tracedecay uninstall (remove agent config instead)."; diff --git a/crates/tracedecay-cli/src/commands.rs b/crates/tracedecay-cli/src/commands.rs index f21ce249b0..7a30438b58 100644 --- a/crates/tracedecay-cli/src/commands.rs +++ b/crates/tracedecay-cli/src/commands.rs @@ -6,6 +6,7 @@ mod index; mod profile_storage; mod scope; mod settings; +mod stale_store_reset; mod storage; pub(crate) use admin_project::{ @@ -26,6 +27,7 @@ pub(crate) use settings::{ handle_gitignore, handle_upload_counter, mutate_project_configuration, project_configuration_set, report_configuration_receipt, }; +pub(crate) use stale_store_reset::handle_wipe_stale; pub(crate) use storage::{ ProfileOfflineAuthority, annotate_reset_required, handle_list, handle_wipe, join_outcome_and_restore, process_error_text, take_profile_offline, try_admit_profile_registry, diff --git a/crates/tracedecay-cli/src/commands/stale_store_reset.rs b/crates/tracedecay-cli/src/commands/stale_store_reset.rs new file mode 100644 index 0000000000..0f1cc67657 --- /dev/null +++ b/crates/tracedecay-cli/src/commands/stale_store_reset.rs @@ -0,0 +1,252 @@ +use std::path::{Path, PathBuf}; +use std::time::Duration; + +use tracedecay_domain::errors::{ResettableStoreV1, Result, StoreResetRequiredV1, TraceDecayError}; +use tracedecay_global_db::profile_registry_maintenance::verify_store_path_absent; +use tracedecay_runtime_core::config::ProfileRoot; +use tracedecay_runtime_core::storage::{ + SESSIONS_DB_FILENAME, profile_sharded_data_root, validate_project_id, +}; +use tracedecay_sessions::runtime::USER_SESSIONS_DB_FILENAME; + +use super::storage::{ + ProfileOfflineAuthority, join_outcome_and_restore, remove_fixed_profile_path, + take_profile_offline, +}; + +const OPERATION: &str = "wipe --stale"; + +/// How long the reset waits for the operator of an unmanaged daemon to stop +/// it. ponytail: a fixed interactive bound; a daemon shutdown request would +/// let the reset stop an unmanaged daemon itself. +const UNMANAGED_DAEMON_STOP_TIMEOUT: Duration = Duration::from_secs(300); + +/// Deletes exactly the stores the running daemon holds in their typed +/// reset-required state, nothing else. The daemon is the authority that +/// refused them, so the list is read from it before the profile is taken +/// offline; each deleted store is recreated empty on its next open. +#[hotpath::measure(label = "cli.wipe.stale", future = true)] +pub(crate) async fn handle_wipe_stale(profile: &ProfileRoot, assume_yes: bool) -> Result<()> { + if !assume_yes { + return Err(TraceDecayError::Config { + message: "wipe --stale deletes every store the daemon reports as requiring reset; \ + pass --yes to confirm. Nothing was wiped." + .to_owned(), + }); + } + if !tracedecay_daemon_control::daemon_reachable(profile) { + return Err(TraceDecayError::Config { + message: "wipe --stale reads the stores that require reset from the running daemon; \ + start it (`tracedecay daemon start`) and re-run. Nothing was wiped." + .to_owned(), + }); + } + let stores = tracedecay_daemon_control::daemon_reset_required_stores( + profile, + crate::product_runtime::PRODUCT_BUILD_VERSION, + )?; + if stores.is_empty() { + eprintln!("No store requires reset. Nothing was wiped."); + return Ok(()); + } + let targets = resettable_targets(&stores)?; + let profile_root = profile.data_dir().to_path_buf(); + let profile_offline = take_stale_reset_offline(profile, &profile_root)?; + let outcome = reset_stores(&profile_root, &targets); + let restore = profile_offline.finish(); + join_outcome_and_restore(OPERATION, outcome, restore) +} + +/// The managed daemon is quiesced and restored afterwards. An unmanaged daemon +/// (`tracedecay daemon run`) has no service to stop and its census ends with +/// it, so the census is read first and the reset waits for its operator to +/// stop it. +fn take_stale_reset_offline( + profile: &ProfileRoot, + profile_root: &Path, +) -> Result { + if tracedecay_daemon_control::installed_service_state(profile)? + != tracedecay_daemon_control::DaemonServiceState::Missing + { + return take_profile_offline(profile, profile_root, OPERATION); + } + eprintln!( + "An unmanaged TraceDecay daemon holds the profile; stop it and wipe --stale continues \ + (waiting up to {}s).", + UNMANAGED_DAEMON_STOP_TIMEOUT.as_secs() + ); + let lease = tracedecay_runtime_core::lifecycle_lease::acquire_exclusive_with_timeout( + profile_root, + OPERATION, + UNMANAGED_DAEMON_STOP_TIMEOUT, + )?; + Ok(ProfileOfflineAuthority::Lease(lease)) +} + +/// Every reported store as one this command resets on its own, or the +/// refusal naming the reset a store needs instead. Nothing is deleted unless +/// all of them are resettable. +fn resettable_targets(stores: &[StoreResetRequiredV1]) -> Result> { + stores + .iter() + .map(|store| { + ResettableStoreV1::from_label(&store.store).ok_or_else(|| TraceDecayError::Config { + message: format!( + "{} requires reset ({}) and is not reset on its own; run `{}`. Nothing was \ + wiped.", + store.store, store.reason, store.remedy + ), + }) + }) + .collect() +} + +fn reset_stores(profile_root: &Path, targets: &[ResettableStoreV1]) -> Result<()> { + for target in targets { + let (directory, database) = store_location(profile_root, target)?; + let removed = remove_store_family(&directory, database)?; + println!( + "reset {} ({removed} entries removed from {}); the daemon recreates it empty", + target.label(), + directory.display() + ); + } + Ok(()) +} + +fn store_location( + profile_root: &Path, + target: &ResettableStoreV1, +) -> Result<(PathBuf, &'static str)> { + match target { + ResettableStoreV1::ProfileSessions => { + Ok((profile_root.to_path_buf(), USER_SESSIONS_DB_FILENAME)) + } + ResettableStoreV1::ProjectSessions { project_id } => { + validate_project_id(project_id).map_err(|message| TraceDecayError::Config { + message: format!("daemon reported an invalid project id `{project_id}`: {message}"), + })?; + Ok(( + profile_sharded_data_root(profile_root, project_id), + SESSIONS_DB_FILENAME, + )) + } + } +} + +/// Removes `database` and every entry named after it in `directory`: its +/// SQLite sidecars and spools (`.db…`), its host-admission directory +/// (`..db…`), and its session relation graph (`.grafeo…`). +fn remove_store_family(directory: &Path, database: &str) -> Result { + let stem = database.strip_suffix(".db").unwrap_or(database); + let owned_prefix = format!("{stem}."); + let hidden_prefix = format!(".{database}."); + let entries = std::fs::read_dir(directory).map_err(|error| TraceDecayError::Config { + message: format!("failed to list '{}': {error}", directory.display()), + })?; + let mut members = Vec::new(); + for entry in entries { + let entry = entry.map_err(|error| TraceDecayError::Config { + message: format!("failed to list '{}': {error}", directory.display()), + })?; + let name = entry.file_name().to_string_lossy().into_owned(); + if name.starts_with(&owned_prefix) || name.starts_with(&hidden_prefix) { + members.push(name); + } + } + let mut removed = 0; + for name in &members { + removed += usize::from(remove_fixed_profile_path(directory, name)?); + verify_store_path_absent(&directory.join(name))?; + } + Ok(removed) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn store_family_removal_keeps_every_other_store_in_the_directory() { + let directory = tempfile::TempDir::new().expect("store directory"); + let root = directory.path(); + for file in [ + "sessions.db", + "sessions.db-wal", + "sessions.db-shm", + "sessions.grafeo", + "tracedecay.db", + "tracedecay.db-wal", + "tracedecay.grafeo", + "store_manifest.json", + "user-sessions.db", + ] { + std::fs::write(root.join(file), file).expect("write store file"); + } + for dir in [ + ".sessions.db.host-admission", + "sessions.db.delivery-settlement-spool-v1", + "code-index-v1", + ] { + std::fs::create_dir(root.join(dir)).expect("create store directory"); + std::fs::write(root.join(dir).join("entry"), dir).expect("write store entry"); + } + + assert_eq!(remove_store_family(root, "sessions.db").unwrap(), 6); + + let mut kept: Vec = std::fs::read_dir(root) + .unwrap() + .map(|entry| entry.unwrap().file_name().to_string_lossy().into_owned()) + .collect(); + kept.sort(); + assert_eq!( + kept, + [ + "code-index-v1", + "store_manifest.json", + "tracedecay.db", + "tracedecay.db-wal", + "tracedecay.grafeo", + "user-sessions.db", + ] + ); + } + + #[test] + fn a_store_that_is_not_resettable_on_its_own_refuses_the_whole_reset() { + let store = |label: &str, remedy: &str| StoreResetRequiredV1 { + store: label.to_owned(), + authority: "observations".to_owned(), + found_version: None, + required_version: None, + reason: "rows predate the unified identity".to_owned(), + remedy: remedy.to_owned(), + }; + assert_eq!( + resettable_targets(&[ + store("profile sessions", "tracedecay wipe --stale --yes"), + store( + "project sessions proj_a5b3d7e3ebe14ca7", + "tracedecay wipe --stale --yes" + ), + ]) + .unwrap(), + [ + ResettableStoreV1::ProfileSessions, + ResettableStoreV1::ProjectSessions { + project_id: "proj_a5b3d7e3ebe14ca7".to_owned() + }, + ] + ); + assert_eq!( + resettable_targets(&[ + store("profile sessions", "tracedecay wipe --stale --yes"), + store("profile authority", "tracedecay wipe --all --yes"), + ]) + .unwrap_err() + .to_string(), + "config error: profile authority requires reset (rows predate the unified identity) \ + and is not reset on its own; run `tracedecay wipe --all --yes`. Nothing was wiped." + ); + } +} diff --git a/crates/tracedecay-cli/src/commands/storage.rs b/crates/tracedecay-cli/src/commands/storage.rs index 26967acd40..58b896c311 100644 --- a/crates/tracedecay-cli/src/commands/storage.rs +++ b/crates/tracedecay-cli/src/commands/storage.rs @@ -104,7 +104,7 @@ fn validate_complete_wipe_profile_root( Ok(()) } -fn remove_fixed_profile_path( +pub(super) fn remove_fixed_profile_path( profile_root: &Path, name: &str, ) -> tracedecay_domain::errors::Result { @@ -330,7 +330,7 @@ mod wipe_safety_tests { assert!(profile.contains("session temporal persisted shape requires reset")); assert!(profile.contains("refused authority: session temporal")); assert!( - profile.contains("\n tracedecay wipe --all --yes"), + profile.contains("\n tracedecay wipe --stale --yes"), "{profile}" ); diff --git a/crates/tracedecay-cli/src/main.rs b/crates/tracedecay-cli/src/main.rs index d39d191bb9..78bbe60f5d 100644 --- a/crates/tracedecay-cli/src/main.rs +++ b/crates/tracedecay-cli/src/main.rs @@ -1309,7 +1309,10 @@ async fn dispatch_project_command( Commands::Storage { action } => { commands::handle_profile_storage_action(profile, action, assume_yes).await?; } - Commands::Wipe { all } => { + Commands::Wipe { stale: true, .. } => { + commands::handle_wipe_stale(profile, assume_yes).await?; + } + Commands::Wipe { all, stale: false } => { commands::handle_wipe(profile, all, assume_yes).await?; } Commands::List { all } => { diff --git a/crates/tracedecay-cli/src/update_cmd.rs b/crates/tracedecay-cli/src/update_cmd.rs index e825dac802..86026fb203 100644 --- a/crates/tracedecay-cli/src/update_cmd.rs +++ b/crates/tracedecay-cli/src/update_cmd.rs @@ -840,7 +840,10 @@ mod tests { found_version: Some(5), required_version: 6, } - .store_reset_required("profile sessions") + .store_reset_required( + "profile sessions", + tracedecay_domain::errors::STALE_STORE_RESET_COMMAND, + ) .expect("a versioned profile refusal is a reset-required store") } @@ -866,7 +869,7 @@ mod tests { pending_reset_line(&reset), " \x1b[33mpending operator action:\x1b[0m profile sessions requires reset \ (git correlation profile schema 5 is incompatible with required schema 6; reset \ - the profile); run `tracedecay wipe --all --yes`" + the profile); run `tracedecay wipe --stale --yes`" ); let failed = update_completion(Some(PluginRefreshOutcome::Failed), &[reset]) .expect_err("a failed refresh still fails `update`") diff --git a/crates/tracedecay-daemon-control/src/service/update_restore_tests.rs b/crates/tracedecay-daemon-control/src/service/update_restore_tests.rs index a8a1112f4e..df09f3fbec 100644 --- a/crates/tracedecay-daemon-control/src/service/update_restore_tests.rs +++ b/crates/tracedecay-daemon-control/src/service/update_restore_tests.rs @@ -282,7 +282,7 @@ fn restore_check_accepts_a_reset_required_daemon_and_returns_its_pending_reset() "found_version": 5, "required_version": 6, "reason": "git correlation profile schema 5 is incompatible with required schema 6; reset the profile", - "remedy": "tracedecay wipe --all --yes", + "remedy": "tracedecay wipe --stale --yes", }], }, }), @@ -317,7 +317,7 @@ fn restore_check_accepts_a_reset_required_daemon_and_returns_its_pending_reset() reason: "git correlation profile schema 5 is incompatible with required schema 6; \ reset the profile" .to_owned(), - remedy: "tracedecay wipe --all --yes".to_owned(), + remedy: "tracedecay wipe --stale --yes".to_owned(), }], } ); diff --git a/crates/tracedecay-daemon-service/src/invocation/retained.rs b/crates/tracedecay-daemon-service/src/invocation/retained.rs index fe0416c253..8ec255e34c 100644 --- a/crates/tracedecay-daemon-service/src/invocation/retained.rs +++ b/crates/tracedecay-daemon-service/src/invocation/retained.rs @@ -17,7 +17,11 @@ pub(super) fn missing_retained_runtime_problem( ) -> DaemonInvocationResponse { if matches!( publication, - Some(ProjectRuntimePublicationStateV1::Warming | ProjectRuntimePublicationStateV1::Failed) + Some( + ProjectRuntimePublicationStateV1::Warming + | ProjectRuntimePublicationStateV1::Failed + | ProjectRuntimePublicationStateV1::ResetRequired(_) + ) ) { return missing_registered_owner_problem(publication, request_id); } diff --git a/crates/tracedecay-daemon-service/src/invocation/work.rs b/crates/tracedecay-daemon-service/src/invocation/work.rs index 0b8f2639c0..98a501dba1 100644 --- a/crates/tracedecay-daemon-service/src/invocation/work.rs +++ b/crates/tracedecay-daemon-service/src/invocation/work.rs @@ -71,10 +71,18 @@ pub(super) fn missing_registered_owner_problem( publication: Option, request_id: String, ) -> DaemonInvocationResponse { - if publication == Some(crate::project_runtime::ProjectRuntimePublicationStateV1::Failed) { - return runtime_publication_failed_problem(request_id); + match publication { + Some(crate::project_runtime::ProjectRuntimePublicationStateV1::Failed) => { + runtime_publication_failed_problem(request_id) + } + Some(crate::project_runtime::ProjectRuntimePublicationStateV1::ResetRequired(refusal)) => { + application_problem( + request_id, + ApplicationProblem::from_detail(refusal.as_ref().clone()), + ) + } + _ => runtime_mounting_problem(request_id), } - runtime_mounting_problem(request_id) } /// Dispatches one Work invocation through the product authority and publishes diff --git a/crates/tracedecay-daemon-service/src/project_runtime.rs b/crates/tracedecay-daemon-service/src/project_runtime.rs index 751afd4caa..2519ec16fa 100644 --- a/crates/tracedecay-daemon-service/src/project_runtime.rs +++ b/crates/tracedecay-daemon-service/src/project_runtime.rs @@ -143,7 +143,7 @@ impl RecoveryCancelProbe { /// stage distinguishes a still-mounting owner from a finished publication /// that will never grow the missing slot, so a permanent composition error /// is not reported as endless retryable pre-admission. -#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +#[derive(Clone, Debug, Default, PartialEq, Eq)] pub enum ProjectRuntimePublicationStateV1 { /// Owner registration is still in progress, or no explicit terminal has /// been recorded yet. @@ -153,6 +153,10 @@ pub enum ProjectRuntimePublicationStateV1 { Ready, /// Project-open publication failed. Missing owners stay missing. Failed, + /// Project-open publication stopped at a store held in its typed + /// reset-required state. Missing owners stay missing and answer this + /// refusal until the store is reset. + ResetRequired(Arc), } /// Exact project-open publication attempt. @@ -272,7 +276,11 @@ impl ProjectRuntime { /// like a missing runtime again instead of a finished failure. fn retain_after_reservation_release(&self) -> bool { self.has_components() - || matches!(self.publication, ProjectRuntimePublicationStateV1::Failed) + || matches!( + self.publication, + ProjectRuntimePublicationStateV1::Failed + | ProjectRuntimePublicationStateV1::ResetRequired(_) + ) } } @@ -1322,7 +1330,7 @@ impl ProjectRuntimeRegistryV1 { let canonical = canonical_existing_identity(project_root).ok(); let runtimes = self.lock_runtimes(); runtime_for_lookup(&runtimes, project_root, canonical.as_deref()) - .map(|runtime| runtime.publication) + .map(|runtime| runtime.publication.clone()) } /// Begin mandatory owner publication for an exact registered root. @@ -1385,7 +1393,7 @@ impl ProjectRuntimeRegistryV1 { let runtime = runtimes.get(project_root); ( runtime.and_then(|runtime| runtime.advisory_cycle.clone()), - runtime.map(|runtime| runtime.publication), + runtime.map(|runtime| runtime.publication.clone()), changed, ) } @@ -1395,6 +1403,19 @@ impl ProjectRuntimeRegistryV1 { self.finish_publication(attempt, ProjectRuntimePublicationStateV1::Failed) } + /// Record that the current attempt stopped at a store held reset-required; + /// every owner it left missing answers `refusal`. + pub fn mark_publication_reset_required( + &self, + attempt: &ProjectRuntimePublicationAttemptV1, + refusal: tracedecay_contracts::ApplicationProblemDetailV1, + ) -> bool { + self.finish_publication( + attempt, + ProjectRuntimePublicationStateV1::ResetRequired(Arc::new(refusal)), + ) + } + /// Record successful mandatory owner publication for the current attempt. pub fn mark_publication_ready(&self, attempt: &ProjectRuntimePublicationAttemptV1) -> bool { self.finish_publication(attempt, ProjectRuntimePublicationStateV1::Ready) diff --git a/crates/tracedecay-daemon-service/src/project_runtime/request_snapshot.rs b/crates/tracedecay-daemon-service/src/project_runtime/request_snapshot.rs index 9166ff0e70..159e9cbbb2 100644 --- a/crates/tracedecay-daemon-service/src/project_runtime/request_snapshot.rs +++ b/crates/tracedecay-daemon-service/src/project_runtime/request_snapshot.rs @@ -58,7 +58,7 @@ impl AdmittedProjectRuntimeV1 { fn capture(runtime: &ProjectRuntime) -> Self { let feedback = runtime.feedback.as_ref(); Self { - publication: runtime.publication, + publication: runtime.publication.clone(), feedback: feedback.map(RegisteredFeedbackRuntime::runtime), feedback_owner: feedback.map(RegisteredFeedbackRuntime::invocation_owner), advisory_cycle: runtime.advisory_cycle.clone(), @@ -219,7 +219,7 @@ impl ProjectRequestRuntimesV1 { Self { admitted: true, resolved_root: Some(request_lease.inner.registered_root.clone()), - publication: Some(admitted.publication), + publication: Some(admitted.publication.clone()), feedback: admitted.feedback.clone(), feedback_owner: admitted.feedback_owner.clone(), advisory_cycle: admitted.advisory_cycle.clone(), diff --git a/crates/tracedecay-domain/src/errors.rs b/crates/tracedecay-domain/src/errors.rs index 331478db3d..e1bea1a2d5 100644 --- a/crates/tracedecay-domain/src/errors.rs +++ b/crates/tracedecay-domain/src/errors.rs @@ -121,6 +121,49 @@ pub type Result = std::result::Result; /// the next open creates the shape the running binary writes. pub const PROFILE_RESET_COMMAND: &str = "tracedecay wipe --all --yes"; +/// Deletes exactly the stores the daemon reports in their typed +/// reset-required state and nothing else; the next open recreates each one +/// empty. Stores it cannot reset on their own name [`PROFILE_RESET_COMMAND`]. +pub const STALE_STORE_RESET_COMMAND: &str = "tracedecay wipe --stale --yes"; + +/// A registered store [`STALE_STORE_RESET_COMMAND`] deletes on its own, named +/// by the `store` label [`StoreResetRequiredV1`] carries. +#[derive(Clone, Debug, PartialEq, Eq)] +pub enum ResettableStoreV1 { + ProfileSessions, + ProjectSessions { project_id: String }, +} + +impl ResettableStoreV1 { + const PROFILE_SESSIONS_LABEL: &'static str = "profile sessions"; + const PROJECT_SESSIONS_PREFIX: &'static str = "project sessions "; + + #[must_use] + pub fn label(&self) -> String { + match self { + Self::ProfileSessions => Self::PROFILE_SESSIONS_LABEL.to_owned(), + Self::ProjectSessions { project_id } => { + format!("{}{project_id}", Self::PROJECT_SESSIONS_PREFIX) + } + } + } + + /// The store a [`Self::label`] names, `None` for a store that is not + /// resettable on its own. + #[must_use] + pub fn from_label(label: &str) -> Option { + if label == Self::PROFILE_SESSIONS_LABEL { + return Some(Self::ProfileSessions); + } + label + .strip_prefix(Self::PROJECT_SESSIONS_PREFIX) + .filter(|project_id| !project_id.is_empty()) + .map(|project_id| Self::ProjectSessions { + project_id: project_id.to_owned(), + }) + } +} + /// A persisted store the daemon keeps mounted in a typed reset-required state /// instead of refusing to serve: every read against it returns the typed /// refusal, stores that admit keep serving, and `remedy` is the exact command @@ -201,9 +244,24 @@ impl TraceDecayError { Some((authority, reason)) } + /// Whether this is a persisted-shape refusal a store is served in until + /// its reset. + #[must_use] + pub fn is_store_reset_required(&self) -> bool { + matches!( + self, + Self::ResetRequired { .. } | Self::ProfileResetRequired { .. } + ) + } + /// The typed reset-required state `store` is mounted in when this error is /// a profile-scoped persisted-shape refusal, `None` for any other failure. - pub fn store_reset_required(&self, store: impl Into) -> Option { + /// `remedy` is the command that resets `store`. + pub fn store_reset_required( + &self, + store: impl Into, + remedy: &str, + ) -> Option { let (authority, found_version, required_version) = match self { Self::ResetRequired { authority, .. } => (authority.clone(), None, None), Self::ProfileResetRequired { @@ -223,7 +281,7 @@ impl TraceDecayError { found_version, required_version, reason: self.to_string(), - remedy: PROFILE_RESET_COMMAND.to_owned(), + remedy: remedy.to_owned(), }) } diff --git a/crates/tracedecay-global-db/src/observation/schema.rs b/crates/tracedecay-global-db/src/observation/schema.rs index 3dfc462409..7323d641e0 100644 --- a/crates/tracedecay-global-db/src/observation/schema.rs +++ b/crates/tracedecay-global-db/src/observation/schema.rs @@ -10,10 +10,18 @@ const OBSERVATION_SCHEMA_MIGRATION: &str = "observations-v2-canonical-autoincrem /// Marker proving `observations` was created by a binary that derives every /// provider's id on `tracedecay.observation.v1`. Recorded at creation. An /// existing store that already holds rows and lacks the marker keeps those -/// rows unread: admission returns a typed reset instead of decoding them. +/// rows unread: admission reports [`OBSERVATIONS_PREDATE_UNIFIED_IDENTITY`] +/// instead of decoding them. const OBSERVATION_UNIFIED_IDENTITY_MIGRATION: &str = "observations-unified-identity-v1"; -const OBSERVATION_UNIFIED_IDENTITY_RESET_REASON: &str = "observation rows predate the unified observation identity and cannot be read; reset the profile so ingestion can rebuild them from host transcripts"; +/// The observation authority of a store written before the unified +/// observation identity. The store's other authorities stay admissible; its +/// session features are refused until the store is reset. +pub(crate) const OBSERVATIONS_PREDATE_UNIFIED_IDENTITY: crate::registered::RefusedAuthorityV1 = + crate::registered::RefusedAuthorityV1 { + authority: "observations", + reason: "observation rows predate the unified observation identity and cannot be read; reset the profile so ingestion can rebuild them from host transcripts", + }; /// Marker proving every `observation_repository_provenance` row references /// its repository capture through `observation_repository_captures` instead of @@ -214,9 +222,12 @@ const OBSERVATION_AUTHORITY_SCHEMA_SQL: &str = FOREIGN KEY(observation_id) REFERENCES observations(observation_id) );"; +/// Installs the observation authority. Returns the refused authority, with +/// the unified-identity marker and every observation-row rewrite skipped, when +/// the store holds rows written before the unified identity. pub async fn ensure_observation_schema( conn: &(impl Executor + Sync), -) -> tracedecay_domain::errors::Result<()> { +) -> tracedecay_domain::errors::Result> { let table_preexisted = observation_table_exists(conn).await?; tracedecay_runtime_core::db::retrieval_anchor_schema::install_retrieval_anchor_schema( conn, @@ -240,10 +251,7 @@ pub async fn ensure_observation_schema( } } else if !migration_recorded(conn, OBSERVATION_UNIFIED_IDENTITY_MIGRATION).await? { if observation_rows_exist(conn).await? { - return Err(tracedecay_domain::errors::TraceDecayError::reset_required( - "observations", - OBSERVATION_UNIFIED_IDENTITY_RESET_REASON, - )); + return Ok(Some(OBSERVATIONS_PREDATE_UNIFIED_IDENTITY)); } conn.execute( "INSERT OR IGNORE INTO global_schema_migrations(migration) VALUES (?1)", @@ -263,7 +271,7 @@ pub async fn ensure_observation_schema( .await .map_err(|error| global_db_operation_error(OBSERVATION_SCHEMA_OPERATION, error))?; } - Ok(()) + Ok(None) } #[cfg(test)] @@ -385,17 +393,35 @@ mod tests { .unwrap(); drop(conn); - let Err(error) = reopen(&path).await else { - panic!("a store holding pre-unified Claude observations must be refused, not decoded"); - }; - let TraceDecayError::ResetRequired { authority, reason } = error else { - panic!("old observation rows must be a typed reset, got {error}"); - }; - assert_eq!(authority, "observations"); + let (lease, owner) = reopen(&path) + .await + .expect("the store's other authorities stay admissible"); + for refusal in [lease.reset_required(), owner.reset_required()] { + let Some(TraceDecayError::ResetRequired { authority, reason }) = refusal else { + panic!("old observation rows must be a typed reset, got {refusal:?}"); + }; + assert_eq!(authority, "observations"); + assert_eq!( + reason, + "observation rows predate the unified observation identity and cannot be read; reset the profile so ingestion can rebuild them from host transcripts" + ); + } + drop((lease, owner)); + let conn = TestConnection::open(&path); + let mut rows = conn + .query( + "SELECT COUNT(*) FROM global_schema_migrations WHERE migration = ?1", + params![OBSERVATION_UNIFIED_IDENTITY_MIGRATION], + ) + .await + .unwrap(); + let recorded: i64 = rows.next().await.unwrap().unwrap().get(0).unwrap(); assert_eq!( - reason, - "observation rows predate the unified observation identity and cannot be read; reset the profile so ingestion can rebuild them from host transcripts" + recorded, 0, + "a refused store never adopts the unified identity" ); + drop(rows); + drop(conn); let conn = TestConnection::open(&path); let mut rows = conn diff --git a/crates/tracedecay-global-db/src/registered.rs b/crates/tracedecay-global-db/src/registered.rs index 62a0ab9802..bd1d423823 100644 --- a/crates/tracedecay-global-db/src/registered.rs +++ b/crates/tracedecay-global-db/src/registered.rs @@ -32,6 +32,21 @@ type SessionRelationGraphStateV1 = RwLock< )>, >; +/// An authority inside an admitted store whose persisted rows this binary +/// refuses to read. The store serves its other authorities; every feature +/// that reads the refused one gets [`Self::error`] until the store is reset. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub(crate) struct RefusedAuthorityV1 { + pub(crate) authority: &'static str, + pub(crate) reason: &'static str, +} + +impl RefusedAuthorityV1 { + pub(crate) fn error(self) -> TraceDecayError { + TraceDecayError::reset_required(self.authority, self.reason) + } +} + /// The sole map owner for one registered global-database publication. /// /// It can issue independently counted client leases, but cannot be cloned or @@ -41,6 +56,7 @@ pub struct RegisteredGlobalDbOwnerV1 { database: DatabaseOwnerV1, project_graph: Arc>, session_relation_graph: Arc, + refused_authority: Option, } /// Cloneable, weak issuance route for one registered global-database owner. @@ -53,6 +69,7 @@ pub struct RegisteredGlobalDbWeakLeaseIssuerV1 { database: DatabaseOwnerWeakLeaseIssuerV1, project_graph: Arc>, session_relation_graph: Weak, + refused_authority: Option, } impl RegisteredGlobalDbOwnerV1 { @@ -81,13 +98,18 @@ impl RegisteredGlobalDbOwnerV1 { ) -> tracedecay_domain::errors::Result { let temporary = database.issue_lease().map_err(registered_owner_error)?; let registered = RegisteredGlobalDb::from_owned_database(temporary); - super::schema_stages::ensure_attached_registered_schema(®istered.database).await?; - super::schema_stages::converge_attached_registered_schema(®istered.database).await?; + let (_, refused_authority) = + super::schema_stages::ensure_attached_registered_schema(®istered.database).await?; + // A store refused for reset is never converged: its reset deletes it. + if refused_authority.is_none() { + super::schema_stages::converge_attached_registered_schema(®istered.database).await?; + } drop(registered); Ok(Self { database, project_graph: Arc::new(OnceLock::new()), session_relation_graph: Arc::new(RwLock::new(None)), + refused_authority, }) } @@ -100,7 +122,7 @@ impl RegisteredGlobalDbOwnerV1 { { let temporary = database.issue_lease().map_err(registered_owner_error)?; let registered = RegisteredGlobalDb::from_owned_database(temporary); - let convergence = + let (convergence, refused_authority) = super::schema_stages::ensure_attached_registered_schema(®istered.database).await?; drop(registered); Ok(( @@ -108,6 +130,7 @@ impl RegisteredGlobalDbOwnerV1 { database, project_graph: Arc::new(OnceLock::new()), session_relation_graph: Arc::new(RwLock::new(None)), + refused_authority, }, convergence, )) @@ -122,10 +145,18 @@ impl RegisteredGlobalDbOwnerV1 { self.database.issue_lease()?, Arc::clone(&self.project_graph), Arc::clone(&self.session_relation_graph), + self.refused_authority, ), )) } + /// The typed reset refusal of the authority this admitted store refuses + /// to read, `None` when every authority it holds is admissible. + #[must_use] + pub fn reset_required(&self) -> Option { + self.refused_authority.map(RefusedAuthorityV1::error) + } + /// Issues a mode-reduced client that can never regain write authority. pub fn issue_read_only_lease(&self) -> Result { Ok(RegisteredGlobalDbLeaseV1::from_database( @@ -133,6 +164,7 @@ impl RegisteredGlobalDbOwnerV1 { self.database.issue_read_only_lease()?, Arc::clone(&self.project_graph), Arc::clone(&self.session_relation_graph), + self.refused_authority, ), )) } @@ -145,6 +177,7 @@ impl RegisteredGlobalDbOwnerV1 { database: self.database.weak_lease_issuer(), project_graph: Arc::clone(&self.project_graph), session_relation_graph: Arc::downgrade(&self.session_relation_graph), + refused_authority: self.refused_authority, } } @@ -193,6 +226,7 @@ impl RegisteredGlobalDbWeakLeaseIssuerV1 { self.database.issue_lease()?, Arc::clone(&self.project_graph), session_relation_graph, + self.refused_authority, ), )) } @@ -260,6 +294,7 @@ pub struct RegisteredGlobalDb { database: Database, project_graph: Arc>, session_relation_graph: Arc, + refused_authority: Option, } impl RegisteredGlobalDb { @@ -306,6 +341,7 @@ impl RegisteredGlobalDb { database, Arc::new(OnceLock::new()), Arc::new(RwLock::new(None)), + None, ) } @@ -313,14 +349,23 @@ impl RegisteredGlobalDb { database: Database, project_graph: Arc>, session_relation_graph: Arc, + refused_authority: Option, ) -> Self { Self { database, project_graph, session_relation_graph, + refused_authority, } } + /// The typed reset refusal of the authority the store behind this client + /// refuses to read, `None` when every authority it holds is admissible. + #[must_use] + pub fn reset_required(&self) -> Option { + self.refused_authority.map(RefusedAuthorityV1::error) + } + /// Wraps an already-published guarded database for WAL maintenance tests. /// /// The exclusive-maintenance truncation lane cannot be exercised through diff --git a/crates/tracedecay-global-db/src/schema_stages.rs b/crates/tracedecay-global-db/src/schema_stages.rs index 9ea53962b3..d78ae4ad53 100644 --- a/crates/tracedecay-global-db/src/schema_stages.rs +++ b/crates/tracedecay-global-db/src/schema_stages.rs @@ -11,6 +11,7 @@ use super::{ global_db_operation_message, managed_test_runs, observability_rollup, observation, observation_projection, project_registry, session_temporal_schema, stack_delivery, }; +use crate::registered::RefusedAuthorityV1; use tracedecay_runtime_core::{ db::{ Database, DatabaseWriteTransaction, @@ -536,7 +537,7 @@ pub async fn ensure_registered_schema_for_admission( .await .map_err(|error| global_db_operation_error(OPERATION, error))?; - install_and_commit_registered_schema( + if let Some(refused) = install_and_commit_registered_schema( transaction, configuration_fresh.as_ref(), temporal_admission, @@ -545,7 +546,10 @@ pub async fn ensure_registered_schema_for_admission( "commit registered global schema", "roll back registered global schema", ) - .await?; + .await? + { + return Err(refused.error()); + } observation_projection::ensure_observation_projection_performance_indexes(installation) .await @@ -584,7 +588,7 @@ async fn install_registered_schema_stages( temporal_admission: session_temporal_schema::SessionTemporalSchemaAdmission, workflow_admission: WorkflowSchemaAdmission, force_exhaustive: bool, -) -> tracedecay_domain::errors::Result<()> { +) -> tracedecay_domain::errors::Result> { Box::pin(install_registered_schema_stage_sequence( transaction, configuration_fresh, @@ -603,7 +607,7 @@ async fn install_and_commit_registered_schema( force_exhaustive: bool, commit_operation: &'static str, rollback_operation: &'static str, -) -> tracedecay_domain::errors::Result<()> +) -> tracedecay_domain::errors::Result> where T: Executor + Sync + SchemaInstallTransaction, { @@ -616,9 +620,12 @@ where ) .await; match admission { - Ok(()) => transaction.commit().await.map_err(|error| { - global_db_operation_error(commit_operation, std::io::Error::other(error)) - }), + Ok(refused_authority) => { + transaction.commit().await.map_err(|error| { + global_db_operation_error(commit_operation, std::io::Error::other(error)) + })?; + Ok(refused_authority) + } Err(error) => match transaction.rollback().await { Ok(()) => Err(error), Err(rollback_error) => Err(global_db_operation_error( @@ -668,7 +675,7 @@ async fn install_registered_schema_stage_sequence( temporal_admission: session_temporal_schema::SessionTemporalSchemaAdmission, workflow_admission: WorkflowSchemaAdmission, force_exhaustive: bool, -) -> tracedecay_domain::errors::Result<()> { +) -> tracedecay_domain::errors::Result> { crate::hotpath_observe::record_transaction_rows(1); let is_fresh = configuration_fresh.is_some(); configuration::ensure_configuration_schema(transaction, configuration_fresh) @@ -803,7 +810,7 @@ async fn install_registered_schema_stage_sequence( } session_temporal_schema::SessionTemporalSchemaAdmission::Current => {} } - observation::ensure_observation_schema(transaction).await?; + let refused_authority = observation::ensure_observation_schema(transaction).await?; observation_projection::ensure_observation_projection_schema(transaction) .await .map_err(|error| global_db_operation_error("initialize observation projection", error))?; @@ -857,7 +864,7 @@ async fn install_registered_schema_stage_sequence( tracedecay_sessions::runtime::workflow_index::ensure_workflow_index_schema(transaction) .await .map_err(|error| global_db_operation_error("initialize workflow index schema", error))?; - Ok(()) + Ok(refused_authority) } /// Completes resumable authority convergence after the registered runtime is @@ -976,10 +983,12 @@ pub async fn converge_attached_registered_schema( /// initialization. The returned convergence plan carries the LCM status-index /// work for lifecycle-owned daemon maintenance; short-lived callers run that /// same work synchronously through [`converge_attached_registered_schema`]. +/// A store whose observation rows this binary refuses is still admitted for +/// its other authorities and returns that refused authority beside the plan. #[hotpath::measure(future = true, label = "global_db.schema.persist.attach")] -pub async fn ensure_attached_registered_schema( +pub(crate) async fn ensure_attached_registered_schema( database: &Database, -) -> tracedecay_domain::errors::Result { +) -> tracedecay_domain::errors::Result<(RegisteredSchemaConvergence, Option)> { let read_connection = database.read_connection(); let RegisteredSchemaAdmissionClassification { configuration_fresh, @@ -990,7 +999,7 @@ pub async fn ensure_attached_registered_schema( let transaction = database .begin_bulk_write_transaction("install attached registered global database schema") .await?; - install_and_commit_registered_schema( + let refused_authority = install_and_commit_registered_schema( transaction, configuration_fresh.as_ref(), temporal_admission, @@ -1013,11 +1022,14 @@ pub async fn ensure_attached_registered_schema( transaction.commit().await?; } validate_admitted_authority_schema(&read_connection, configuration_fresh.is_some()).await?; - Ok(RegisteredSchemaConvergence { - force_exhaustive, - is_fresh: configuration_fresh.is_some(), - lcm_status_performance_indexes: true, - }) + Ok(( + RegisteredSchemaConvergence { + force_exhaustive, + is_fresh: configuration_fresh.is_some(), + lcm_status_performance_indexes: true, + }, + refused_authority, + )) } #[derive(Clone, Copy, Debug, PartialEq, Eq)] diff --git a/crates/tracedecay-mcp/src/tool_errors.rs b/crates/tracedecay-mcp/src/tool_errors.rs index b460d8e133..12de89f17f 100644 --- a/crates/tracedecay-mcp/src/tool_errors.rs +++ b/crates/tracedecay-mcp/src/tool_errors.rs @@ -2,7 +2,9 @@ use serde_json::{Value, json}; use tracedecay_contracts::ApplicationProblem; -use tracedecay_domain::errors::{PROFILE_RESET_COMMAND, TraceDecayError}; +use tracedecay_domain::errors::{ + PROFILE_RESET_COMMAND, STALE_STORE_RESET_COMMAND, TraceDecayError, +}; use tracedecay_domain::{ CURSOR_INVALID_CODE, CURSOR_PARAMETER_CHANGED_CODE, CursorBindingMismatchV1, }; @@ -387,6 +389,22 @@ fn is_project_store_authority(authority: &str) -> bool { PROJECT_STORE_AUTHORITIES.contains(&authority) } +/// Authorities refused while admitting a registered store. The daemon holds +/// that store in its typed reset-required state, so the scoped reset deletes +/// exactly the refused stores. +const REGISTERED_STORE_AUTHORITIES: [&str; 6] = [ + "observations", + "session temporal", + "workflow", + "authority schema", + "LCM profile schema", + "git correlation profile schema", +]; + +fn is_registered_store_authority(authority: &str) -> bool { + REGISTERED_STORE_AUTHORITIES.contains(&authority) +} + /// Authorities whose refused shape was written into an agent host's files. /// No profile reset reaches it; its reset deletes exactly the block or /// package the refusal reason names. @@ -418,6 +436,8 @@ pub fn reset_required_command(authority: &str, project_root: Option<&std::path:: ) } else if is_host_artifact_authority(authority) { "delete the block or package directory named in the refusal".to_string() + } else if is_registered_store_authority(authority) { + STALE_STORE_RESET_COMMAND.to_string() } else { PROFILE_RESET_COMMAND.to_string() } @@ -447,6 +467,15 @@ pub fn reset_required_remedy(authority: &str, project_root: Option<&std::path::P the next managed-skill export writes the current shape" ); } + if is_registered_store_authority(authority) { + return format!( + "refused authority: {authority}\n\ + this binary does not open or migrate that shape; reset it (its old data is \ + deleted, nothing is backed up):\n \ + {command} deletes only the stores the daemon reports as requiring reset\n\ + the daemon recreates each one empty" + ); + } format!( "refused authority: {authority}\n\ this binary does not open or migrate that shape; reset it (its old data is deleted, \ @@ -531,11 +560,11 @@ mod tests { let remedy = wire["error"]["data"]["remedy"] .as_str() .expect("the refusal names its reset command"); - assert!(remedy.contains("tracedecay wipe --all --yes"), "{remedy}"); + assert!(remedy.contains("tracedecay wipe --stale --yes"), "{remedy}"); assert!( wire["error"]["message"] .as_str() - .is_some_and(|message| message.contains("tracedecay wipe --all --yes")), + .is_some_and(|message| message.contains("tracedecay wipe --stale --yes")), "the human-readable message must carry the reset command too: {wire}" ); assert!( @@ -567,9 +596,17 @@ mod tests { ); assert!(!project.contains("wipe --all"), "{project}"); - let profile = super::reset_required_remedy("session temporal", None); - assert!(profile.contains("refused authority: session temporal")); - assert!(!profile.contains("tracedecay update"), "{profile}"); + let stale = super::reset_required_remedy("session temporal", None); + assert!(stale.contains("refused authority: session temporal")); + assert!(!stale.contains("tracedecay update"), "{stale}"); + assert!( + stale.contains("\n tracedecay wipe --stale --yes"), + "{stale}" + ); + assert!(!stale.contains("wipe --all"), "{stale}"); + + let profile = super::reset_required_remedy("project registry", None); + assert!(profile.contains("refused authority: project registry")); assert!( profile.contains("\n tracedecay wipe --all --yes"), "{profile}" diff --git a/crates/tracedecay-project/src/project/lifecycle/branches.rs b/crates/tracedecay-project/src/project/lifecycle/branches.rs index f8ad8e4a7e..7cc951eca6 100644 --- a/crates/tracedecay-project/src/project/lifecycle/branches.rs +++ b/crates/tracedecay-project/src/project/lifecycle/branches.rs @@ -128,7 +128,7 @@ impl TraceDecay { ) .await?; let configuration_database = runtime_registry - .project_sessions(project_id, enrollment_roots) + .project_session_store(project_id, enrollment_roots) .await?; Self::open_branch_with_registered_configuration( project_root, diff --git a/crates/tracedecay-project/src/project/lifecycle/mod.rs b/crates/tracedecay-project/src/project/lifecycle/mod.rs index cbfd7e54fc..a3c6111e5c 100644 --- a/crates/tracedecay-project/src/project/lifecycle/mod.rs +++ b/crates/tracedecay-project/src/project/lifecycle/mod.rs @@ -308,7 +308,7 @@ impl TraceDecay { project_id.as_str(), )?; let configuration_database = runtime_registry - .project_sessions(project_id, vec![project_root.to_path_buf()]) + .project_session_store(project_id, vec![project_root.to_path_buf()]) .await?; Self::init_with_registered_configuration( project_root, @@ -569,7 +569,7 @@ impl TraceDecay { ) .await?; let configuration_database = runtime_registry - .project_sessions(project_id, enrollment_roots) + .project_session_store(project_id, enrollment_roots) .await?; Self::open_with_registered_configuration( project_root, @@ -764,7 +764,7 @@ impl TraceDecay { ) .await?; let configuration_database = runtime_registry - .project_sessions(project_id, enrollment_roots) + .project_session_store(project_id, enrollment_roots) .await?; Self::open_read_only_with_registered_configuration( project_root, diff --git a/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs b/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs index c6603d753c..2340b8285b 100644 --- a/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs +++ b/crates/tracedecay-store-runtime/src/session_registry/maintenance.rs @@ -11,7 +11,10 @@ use tokio::sync::Semaphore; use tracedecay_contracts::storage::{ SchemaConvergenceFindingV1, SchemaConvergenceStageV1, SchemaConvergenceStateV1, }; -use tracedecay_domain::errors::{StoreResetRequiredV1, TraceDecayError}; +use tracedecay_domain::errors::{ + PROFILE_RESET_COMMAND, ResettableStoreV1, STALE_STORE_RESET_COMMAND, StoreResetRequiredV1, + TraceDecayError, +}; use tracedecay_global_db::schema_stages::RegisteredSchemaConvergence; use tracedecay_store::{StoreRuntimeBindingV1, StoreShardIdV1, StoreShardScopeV1}; @@ -547,7 +550,14 @@ impl DaemonSessionRuntimeRegistryV1 { ) -> Result { let shard_id = runtime.binding().shard_id.clone(); let attached = self.attach_registered_inner(runtime).await; - self.record_registered_admission(shard_id, attached.as_ref().err()); + let refused_authority = attached + .as_ref() + .ok() + .and_then(RegisteredGlobalDbOwnerV1::reset_required); + self.record_registered_admission( + shard_id, + attached.as_ref().err().or(refused_authority.as_ref()), + ); attached } @@ -557,15 +567,25 @@ impl DaemonSessionRuntimeRegistryV1 { refusal: Option<&TraceDecayError>, ) { let reset_required = refusal.and_then(|error| { - let store = match &shard_id.scope { - StoreShardScopeV1::Profile => "profile authority".to_owned(), - StoreShardScopeV1::ProfileSessions => "profile sessions".to_owned(), + let resettable = match &shard_id.scope { + StoreShardScopeV1::ProfileSessions => Some(ResettableStoreV1::ProfileSessions), StoreShardScopeV1::ProjectSessions { project_id } => { - format!("project sessions {project_id}") + Some(ResettableStoreV1::ProjectSessions { + project_id: project_id.to_string(), + }) } - scope => format!("{scope:?}"), + _ => None, }; - error.store_reset_required(store) + match resettable { + Some(store) => error.store_reset_required(store.label(), STALE_STORE_RESET_COMMAND), + None => { + let store = match &shard_id.scope { + StoreShardScopeV1::Profile => "profile authority".to_owned(), + scope => format!("{scope:?}"), + }; + error.store_reset_required(store, PROFILE_RESET_COMMAND) + } + } }); let mut stores = self .reset_required_stores @@ -614,9 +634,10 @@ impl DaemonSessionRuntimeRegistryV1 { Box::pin(async move { let database = Database::publish_runtime(runtime, DatabaseAccessMode::ReadWrite).await?; - let long_lived = self.long_lived_session_maintenance; // Only long-lived daemons defer schema convergence to resumable - // maintenance. + // maintenance. A store refused for reset is never converged: its + // reset deletes it. + let long_lived = self.long_lived_session_maintenance; let (database, convergence) = if long_lived { let (database, convergence) = RegisteredGlobalDbOwnerV1::admit_and_attach_for_daemon(database).await?; @@ -627,7 +648,7 @@ impl DaemonSessionRuntimeRegistryV1 { None, ) }; - if long_lived { + if long_lived && database.reset_required().is_none() { let lease = database.issue_lease().map_err(|error| { session_registry_error( "issue registered schema convergence client", diff --git a/crates/tracedecay-store-runtime/src/session_registry/mounts.rs b/crates/tracedecay-store-runtime/src/session_registry/mounts.rs index c2fe5e2922..ee45f078f3 100644 --- a/crates/tracedecay-store-runtime/src/session_registry/mounts.rs +++ b/crates/tracedecay-store-runtime/src/session_registry/mounts.rs @@ -483,8 +483,15 @@ impl DaemonSessionRuntimeRegistryV1 { Ok(lease) } + /// The profile session store. A store held in its typed reset-required + /// state answers that refusal instead of a client. #[hotpath::skip] pub async fn profile_sessions(&self) -> Result { + refuse_reset_required(Box::pin(self.mount_profile_session_store()).await?) + } + + #[hotpath::skip] + async fn mount_profile_session_store(&self) -> Result { let existing = { let mounted = self .profile_sessions @@ -1172,10 +1179,24 @@ impl DaemonSessionRuntimeRegistryV1 { Ok(identities) } + /// The mounted project session store for session features, `None` while + /// it is held in its typed reset-required state. #[hotpath::skip] pub async fn mounted_project_sessions( &self, project_id: &ProjectId, + ) -> Option { + self.mounted_project_session_store(project_id) + .await + .filter(|lease| lease.reset_required().is_none()) + } + + /// The mounted project session store for the authorities it holds + /// besides sessions; see [`Self::project_session_store`]. + #[hotpath::skip] + pub async fn mounted_project_session_store( + &self, + project_id: &ProjectId, ) -> Option { let mounted = self .project_owners @@ -1198,6 +1219,30 @@ impl DaemonSessionRuntimeRegistryV1 { project_id: ProjectId, enrollment_roots: impl IntoIterator, ) -> Result { + self.register_project_session_authority(&project_id, enrollment_roots)?; + self.mount_registered_project_sessions(project_id).await + } + + /// The project session store for the authorities it holds besides + /// sessions: the project's configuration and query cursor keys. Unlike + /// [`Self::project_sessions`], a store held in its typed reset-required + /// state still serves them, so the project opens and code intelligence + /// serves while every session feature answers the refusal. + #[hotpath::skip] + pub async fn project_session_store( + &self, + project_id: ProjectId, + enrollment_roots: impl IntoIterator, + ) -> Result { + self.register_project_session_authority(&project_id, enrollment_roots)?; + self.mount_project_session_store(project_id).await + } + + fn register_project_session_authority( + &self, + project_id: &ProjectId, + enrollment_roots: impl IntoIterator, + ) -> Result<()> { self.resolver .register_project_authority(LocalProjectEnrollmentAuthorityV1::new( project_id.clone(), @@ -1205,14 +1250,23 @@ impl DaemonSessionRuntimeRegistryV1 { )) .map_err(|error| { session_registry_error("register project session authority", format!("{error:?}")) - })?; - self.mount_registered_project_sessions(project_id).await + }) } + /// The project session store. A store held in its typed reset-required + /// state answers that refusal instead of a client. #[hotpath::skip] pub async fn mount_registered_project_sessions( &self, project_id: ProjectId, + ) -> Result { + refuse_reset_required(Box::pin(self.mount_project_session_store(project_id)).await?) + } + + #[hotpath::skip] + async fn mount_project_session_store( + &self, + project_id: ProjectId, ) -> Result { let has_entry = { let mounted = self @@ -1746,6 +1800,14 @@ impl DaemonSessionRuntimeRegistryV1 { } } +/// A session client, or the typed reset refusal of the store behind it. +fn refuse_reset_required(lease: RegisteredGlobalDbLeaseV1) -> Result { + match lease.reset_required() { + Some(refusal) => Err(refusal), + None => Ok(lease), + } +} + fn issue_mounted_memory_lease( owner: &MemoryStoreOwnerV1, access: DatabaseAccessMode, diff --git a/crates/tracedecay/src/daemon/bootstrap.rs b/crates/tracedecay/src/daemon/bootstrap.rs index 43429be26c..2e2423c176 100644 --- a/crates/tracedecay/src/daemon/bootstrap.rs +++ b/crates/tracedecay/src/daemon/bootstrap.rs @@ -839,7 +839,7 @@ async fn install_profile_worker_plan( // dead daemon. Its persisted worker selection is unreadable, so the // plan runs on the selection a reset profile initializes to until the // operator resets it. - Err(error) if error.store_reset_required("profile sessions").is_some() => { + Err(error) if error.is_store_reset_required() => { log_daemon_event( "profile_worker_plan_reset_required", &[ diff --git a/crates/tracedecay/src/daemon/branch_admin.rs b/crates/tracedecay/src/daemon/branch_admin.rs index 40104d364d..d44bcfb94d 100644 --- a/crates/tracedecay/src/daemon/branch_admin.rs +++ b/crates/tracedecay/src/daemon/branch_admin.rs @@ -104,6 +104,14 @@ pub(super) fn graph_writer_scope( store_writer_scope(&cg.store_layout().data_root, class) } +/// What a caller reads from a project session store: session features are +/// refused while the store is held reset-required, its configuration is not. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum ProjectSessionShardUse { + Sessions, + Configuration, +} + #[cfg(unix)] #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] pub(super) enum MaintenanceReaperKind { @@ -1048,11 +1056,48 @@ impl StoreAdministration { graphs } + /// The project session store for session features. A store held in its + /// typed reset-required state answers that refusal. #[hotpath::measure(label = "daemon.branch_admin.project_session_database", future = true)] pub(super) async fn registered_project_session_database( &self, project_root: &Path, store_layout: &tracedecay_runtime_core::storage::StoreLayout, + ) -> Result { + Box::pin(self.registered_project_session_shard( + project_root, + store_layout, + ProjectSessionShardUse::Sessions, + )) + .await + } + + /// The project session store as the project's configuration authority, + /// served even while the store's session features answer a reset + /// refusal. + #[hotpath::measure( + label = "daemon.branch_admin.project_configuration_database", + future = true + )] + pub(super) async fn registered_project_configuration_database( + &self, + project_root: &Path, + store_layout: &tracedecay_runtime_core::storage::StoreLayout, + ) -> Result { + Box::pin(self.registered_project_session_shard( + project_root, + store_layout, + ProjectSessionShardUse::Configuration, + )) + .await + } + + #[hotpath::skip] + async fn registered_project_session_shard( + &self, + project_root: &Path, + store_layout: &tracedecay_runtime_core::storage::StoreLayout, + shard_use: ProjectSessionShardUse, ) -> Result { let project_id = store_layout .identity @@ -1098,7 +1143,14 @@ impl StoreAdministration { &project_id, )) .await?; - Box::pin(registry.project_sessions(project_id, enrollment_roots)).await + match shard_use { + ProjectSessionShardUse::Sessions => { + Box::pin(registry.project_sessions(project_id, enrollment_roots)).await + } + ProjectSessionShardUse::Configuration => { + Box::pin(registry.project_session_store(project_id, enrollment_roots)).await + } + } } #[cfg(test)] @@ -1804,7 +1856,7 @@ impl StoreAdministration { }); }; let configuration_database = self - .registered_project_session_database(project_root, &layout) + .registered_project_configuration_database(project_root, &layout) .await?; // Branch administration runs inside the daemon, which owns the durable // configuration store. Resolve the pinned snapshot on demand when this diff --git a/crates/tracedecay/src/daemon/http_application_router.rs b/crates/tracedecay/src/daemon/http_application_router.rs index 1f665eb4ef..da957fdcac 100644 --- a/crates/tracedecay/src/daemon/http_application_router.rs +++ b/crates/tracedecay/src/daemon/http_application_router.rs @@ -119,7 +119,7 @@ pub(super) async fn install_remote_http_application_router( // store is reset-required no remote request can be admitted, so the // remote protocol stays unmounted until the operator's reset restarts // the daemon on a fresh profile. - Err(error) if error.store_reset_required("profile authority").is_some() => { + Err(error) if error.is_store_reset_required() => { log_daemon_event( "remote_protocol_router_unmounted", &[ diff --git a/crates/tracedecay/src/daemon/project_composition.rs b/crates/tracedecay/src/daemon/project_composition.rs index dbcb1dd8d5..444df435ae 100644 --- a/crates/tracedecay/src/daemon/project_composition.rs +++ b/crates/tracedecay/src/daemon/project_composition.rs @@ -1699,6 +1699,49 @@ impl ProjectOpenInputs<'_> { } } + /// Code reads (primitives and callable code) on a route whose full + /// upgrade a reset-required session store refused. A refusal to mount + /// them leaves those reads answering the store's reset refusal. + async fn register_reset_required_code_reads( + &self, + core: &ComposedCoreServer, + resolved: &Arc, + ) { + let registered = match tracedecay_domain::ProjectId::new(core.project_id.clone()) { + Ok(project_id) => match core + .graph_runtime + .mounted_project_session_store(&project_id) + .await + { + Some(session_store) => { + Box::pin( + project_open_owners::register_reset_required_route_code_read_owners( + self.invocation, + self.canonical_project_path, + &core.project_id, + resolved, + session_store, + ), + ) + .await + } + None => Err(TraceDecayError::Config { + message: "the project session store is not mounted".to_owned(), + }), + }, + Err(error) => Err(TraceDecayError::Config { + message: format!("invalid project identity: {error}"), + }), + }; + if let Err(error) = registered { + self.log_phase( + "reset_required_code_reads_unavailable", + Some(("error", error.to_string())), + self.started, + ); + } + } + /// A failed upgrade either degrades the route to the still-published core /// (logging the failure and retiring the full server if it had gone live) /// or, when the core cannot be reclaimed, retires every server this @@ -1734,11 +1777,21 @@ impl ProjectOpenInputs<'_> { mutation.mark_failed(); } if core_retained { + let reset_refusal = + tracedecay_mcp::reset_required_context(&error).and_then(|(authority, _)| { + tracedecay_contracts::ApplicationProblemDetailV1::from_reset_required( + &error, + tracedecay_mcp::reset_required_command(&authority, None), + ) + }); if let Some(attempt) = &activation.publication_attempt { - self.invocation - .service - .project_runtimes - .mark_publication_failed(attempt); + let project_runtimes = &self.invocation.service.project_runtimes; + match &reset_refusal { + Some(refusal) => { + project_runtimes.mark_publication_reset_required(attempt, refusal.clone()) + } + None => project_runtimes.mark_publication_failed(attempt), + }; } if let Some(failed_full_server) = failed_full_server { failed_full_server.revoke_project_server_responses(); @@ -1753,8 +1806,9 @@ impl ProjectOpenInputs<'_> { // A session store in its typed reset-required state keeps the // full upgrade refused until the operator resets it, so the // retained core serves the code index meanwhile. - if tracedecay_mcp::reset_required_context(&error).is_some() { + if reset_refusal.is_some() { let code_index_status = self.activate_code_index(core); + Box::pin(self.register_reset_required_code_reads(core, resolved)).await; self.log_phase( "full_upgrade_reset_required", Some(("code_index", code_index_status.to_owned())), diff --git a/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs b/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs index d521c861e4..aaffa8ca57 100644 --- a/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs +++ b/crates/tracedecay/src/daemon/project_composition/code_index_activation.rs @@ -6,7 +6,11 @@ use super::*; use tracedecay_code_index_runtime::code_index_scheduler::{ - CodeIndexDemandAdmissionV1, CodeIndexDemandV1, query_runtime::QueryRuntimeMountErrorV1, + CodeIndexDemandAdmissionV1, CodeIndexDemandV1, + query_runtime::{ + DeferredMountAttemptV1, QueryRuntimeMountErrorV1, + retry_deferred_query_authority_until_serving, + }, }; use tracedecay_runtime_core::logging::log_daemon_event; use tracedecay_session_temporal_store::SessionTemporalAccess; @@ -64,11 +68,6 @@ pub(super) fn code_index_activation_mount( } let query_project_id = project_id.clone(); let query_graph_runtime = Arc::clone(&graph_runtime); - // Order-sensitive: subscribing before the mount is what keeps the - // first generation publication observable by the waiter below. - let publications = invocation - .code_index_schedulers - .subscribe_generation_publications(); let mount = invocation.mount_code_index( project_id, &project_root, @@ -92,7 +91,6 @@ pub(super) fn code_index_activation_mount( // mounted. Keep that wait in its own route-fenced task. spawn_query_authority_when_generation_ready(QueryAuthorityWaitInputs { invocation: invocation.clone(), - publications, project_root: project_root.clone(), project_id: query_project_id, graph_runtime: query_graph_runtime, @@ -111,8 +109,6 @@ pub(super) fn code_index_activation_mount( /// Route-fenced inputs for the post-mount query-authority wait. struct QueryAuthorityWaitInputs { invocation: DaemonInvocationState, - publications: - tokio::sync::broadcast::Receiver, project_root: PathBuf, project_id: tracedecay_domain::ProjectId, graph_runtime: Arc, @@ -121,99 +117,72 @@ struct QueryAuthorityWaitInputs { cancellation: CancellationToken, } -/// Wait for this project's first sealed generation, then mount the checked-in -/// core query policy from the project's durable cursor-key authority. Route revocation (which cancels this route's own token) and a -/// closed publication channel each end the wait without mounting. +/// Wait for this project's first retained generation, then mount the +/// checked-in core query policy from the project's durable cursor-key +/// authority. Route revocation (which cancels this route's own token) ends +/// the wait without mounting. /// -/// A publication is announced before its generation is seated, and a retained -/// `Noop` restore seats without announcing at all, so the subscription alone -/// can miss the first serving generation. The serving-seat signal covers both: -/// every slot write records a seat, and the loop re-reads the exact state -/// (`latest_generation_id`) on each wake. Subscribing before the first read is -/// what keeps a seat that lands during it observable. +/// The wait is the same one the full route's deferred mount uses: a restart +/// restores its retained generation without announcing it and may leave the +/// decoded serving slot empty, so only the retained text owner proves the +/// generation a mount needs. fn spawn_query_authority_when_generation_ready(inputs: QueryAuthorityWaitInputs) { let QueryAuthorityWaitInputs { - invocation: authority_invocation, - mut publications, - project_root: authority_project, - project_id: authority_project_id, - graph_runtime: authority_graph_runtime, - scope: authority_scope, - route_registered: authority_route_registered, - cancellation: authority_cancellation, + invocation, + project_root, + project_id, + graph_runtime, + scope, + route_registered, + cancellation, } = inputs; tokio::spawn(hotpath::future!( async move { - // Order-sensitive: subscribe before the first exact-state read. - let mut seats = authority_invocation - .code_index_schedulers - .subscribe_serving_seats(); - let generation_ready = loop { - if authority_invocation - .code_index_schedulers - .latest_generation_id(&authority_project) - .await - .is_some() - { - break true; - } - tokio::select! { - () = authority_cancellation.cancelled() => break false, - Ok(()) = seats.changed() => {} - publication = publications.recv() => match publication { - Ok(publication) if publication.project_root == authority_project => { - break true; - } - Ok(_) => {} - Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {} - Err(tokio::sync::broadcast::error::RecvError::Closed) => break false, - } - } - }; - if !generation_ready - || authority_cancellation.is_cancelled() - || !authority_route_registered.load(Ordering::Acquire) - { - return; - } + let schedulers = invocation.code_index_schedulers.clone(); let mut awaiting_generation_logged = false; - loop { - let outcome = tokio::select! { - biased; - () = authority_cancellation.cancelled() => return, - outcome = mount_core_query_authority_from_project_sessions( - &authority_invocation, - &authority_graph_runtime, - &authority_project_id, - &authority_project, - &authority_scope, - ) => outcome, - }; - if authority_cancellation.is_cancelled() - || !authority_route_registered.load(Ordering::Acquire) - { - return; - } - match outcome { - Err(error @ QueryRuntimeMountErrorV1::GenerationUnavailable) => { - if !awaiting_generation_logged { - log_query_authority_activation_outcome(&authority_project, Err(error)); - awaiting_generation_logged = true; + let retry = retry_deferred_query_authority_until_serving( + &schedulers, + project_root.clone(), + || { + let first_await = !awaiting_generation_logged; + awaiting_generation_logged = true; + let invocation = invocation.clone(); + let graph_runtime = Arc::clone(&graph_runtime); + let project_id = project_id.clone(); + let project_root = project_root.clone(); + let scope = scope.clone(); + let route_registered = Arc::clone(&route_registered); + async move { + if !route_registered.load(Ordering::Acquire) { + return DeferredMountAttemptV1::Terminal; } - tokio::select! { - () = authority_cancellation.cancelled() => return, - Ok(()) = seats.changed() => {} - publication = publications.recv() => match publication { - Ok(_) | Err(tokio::sync::broadcast::error::RecvError::Lagged(_)) => {} - Err(tokio::sync::broadcast::error::RecvError::Closed) => return, - } + let outcome = mount_core_query_authority_from_project_sessions( + &invocation, + &graph_runtime, + &project_id, + &project_root, + &scope, + ) + .await; + let awaiting = matches!( + outcome, + Err(QueryRuntimeMountErrorV1::GenerationUnavailable) + ); + if !awaiting || first_await { + log_query_authority_activation_outcome(&project_root, outcome); + } + if awaiting { + DeferredMountAttemptV1::AwaitNextPublication + } else { + DeferredMountAttemptV1::Terminal } } - outcome => { - log_query_authority_activation_outcome(&authority_project, outcome); - return; - } - } + }, + ); + tokio::select! { + biased; + () = cancellation.cancelled() => {} + () = retry => {} } }, label = "daemon.project.activate.query_authority" @@ -229,7 +198,10 @@ async fn mount_core_query_authority_from_project_sessions( project_root: &Path, scope: &tracedecay_contracts::ResolvedScope, ) -> std::result::Result<(), QueryRuntimeMountErrorV1> { - let Some(session_db) = graph_runtime.mounted_project_sessions(project_id).await else { + let Some(session_db) = graph_runtime + .mounted_project_session_store(project_id) + .await + else { return Err(QueryRuntimeMountErrorV1::Mount( "project session database is not mounted".to_owned(), )); diff --git a/crates/tracedecay/src/daemon/project_open_handshake.rs b/crates/tracedecay/src/daemon/project_open_handshake.rs index 62bb3e30b6..5651865f41 100644 --- a/crates/tracedecay/src/daemon/project_open_handshake.rs +++ b/crates/tracedecay/src/daemon/project_open_handshake.rs @@ -90,7 +90,7 @@ pub(super) async fn open_project_for_handshake( )?; } let configuration_database = Box::pin( - store_administration.registered_project_session_database(project_path, &store_layout), + store_administration.registered_project_configuration_database(project_path, &store_layout), ) .await?; #[cfg(any(test, feature = "test-helpers"))] diff --git a/crates/tracedecay/src/daemon/project_open_owners.rs b/crates/tracedecay/src/daemon/project_open_owners.rs index 142a6f55e5..e07e9ab26b 100644 --- a/crates/tracedecay/src/daemon/project_open_owners.rs +++ b/crates/tracedecay/src/daemon/project_open_owners.rs @@ -12,7 +12,9 @@ use std::time::{Duration, Instant}; use tracedecay_application::advisory::github_runtime::github_repository_from_remote_v1; use tracedecay_application::project_open_authorization::project_open_work_grant; -use tracedecay_contracts::{ApplicationContractError, ResolvedScope, now_micros}; +use tracedecay_contracts::{ + ApplicationContractError, ResolvedScope, TemporalRetrievalPort, now_micros, +}; use tracedecay_domain::{ProjectId, UtcMicros, canonical_sha256}; use super::DaemonInvocationState; @@ -22,7 +24,6 @@ use tracedecay_agent_hosts::native_integration::{ NativeIntegrationTargetV1, }; use tracedecay_application::lsp_runtime::DaemonLspSessionFactory; -use tracedecay_application::primitives::admitted_root_uri_for_project; use tracedecay_application::source_authorization::ProjectSourceAccessSnapshot; use tracedecay_code_index_runtime::git_transactions::DaemonGitIndexTransactionServiceRegistry; use tracedecay_daemon_service::{ @@ -37,6 +38,7 @@ use tracedecay_daemon_service::{ use tracedecay_domain::errors::{Result, TraceDecayError}; use tracedecay_lsp::analyzer::broker::AdmittedLspProvider; use tracedecay_lsp::analyzer::client::LspRefreshTimeouts; +use tracedecay_session_runtime::session_retrieval::DaemonSessionLookupPrimitiveV1; use tracedecay_session_temporal_store::SessionTemporalAccess; mod advisory_runtime; @@ -163,6 +165,116 @@ pub(crate) async fn install_project_open_source_edit_owners_for_test( Ok(true) } +/// Registers the owners that serve code reads on an admitted route: the +/// callable-code authorization and the primitive runtime. Returns the +/// admitted root URI they were bound to. +#[hotpath::measure(label = "daemon.project.owners.code_reads", future = true)] +async fn register_project_code_read_owners( + invocation: &DaemonInvocationState, + project_root: &Path, + server: &McpServer, + graph: &tracedecay_project::project::TraceDecay, + access: &ProjectSourceAccessSnapshot, + session_db: tracedecay_global_db::RegisteredGlobalDbLeaseV1, + temporal: Arc, +) -> Result { + // Primitive reads are part of the admitted core route. Publish their + // runtime and callable-code authorization before the slower mutation, + // delivery, native-integration, and Work owners finish mounting. + match hotpath::future!( + invocation.feedback_runtime_registrar().open_and_register( + graph.db().clone(), + project_root.to_path_buf(), + graph.store_layout().response_handle_root.clone(), + access.scope.clone(), + access.clone(), + Arc::new(DaemonCallableCodeAuthorizationSource::production( + project_root.to_path_buf(), + access.scope.clone(), + Arc::clone(graph.configuration_runtime()), + )), + ), + label = "daemon.project.open.owners.feedback" + ) + .await + { + Ok(_) | Err(DaemonFeedbackRuntimeRegistrationError::AlreadyRegistered) => {} + Err(error) => { + return Err(TraceDecayError::Config { + message: format!("project-open feedback runtime registration failed: {error:?}"), + }); + } + } + let source = graph + .source_read_context() + .ok_or_else(|| TraceDecayError::Config { + message: "project-open primitive runtime requires an exact registered source identity" + .to_owned(), + })?; + open_and_register_project_primitive_runtime( + invocation, + project_root, + source, + server, + session_db, + temporal, + access.clone(), + ) + .await +} + +/// Registers the code-read owners on a route whose full upgrade is refused +/// because a session store is held in its typed reset-required state. Code +/// reads serve from the retained core; session lookups answer the refusal. +/// `session_store` is the project session store for its non-session +/// authorities (the primitive cursor keys). +#[hotpath::measure( + label = "daemon.project.owners.reset_required_code_reads", + future = true +)] +pub(super) async fn register_reset_required_route_code_read_owners( + invocation: &DaemonInvocationState, + project_root: &Path, + project_id: &str, + server: &McpServer, + session_store: tracedecay_global_db::RegisteredGlobalDbLeaseV1, +) -> Result<()> { + let project_id = + ProjectId::new(project_id.to_owned()).map_err(|_| TraceDecayError::Config { + message: "project-open owners require an authoritative project identity".to_owned(), + })?; + let graph = server.cg().await; + let scope = + tracedecay_code_index_runtime::resolved_scope_for_project(project_root, &project_id) + .map_err(|error| TraceDecayError::Config { + message: format!("project-open resolved scope denied: {error}"), + })?; + let configuration = graph + .configuration_runtime() + .client() + .current() + .await + .map_err(|error| TraceDecayError::Config { + message: format!("project-open configuration currentness failed: {error}"), + })?; + let access = + daemon_owned_project_source_access_at(&scope, project_root, &configuration, now_micros()) + .map_err(|error| TraceDecayError::Config { + message: format!("project-open source access denied: {error}"), + })?; + Box::pin(register_project_code_read_owners( + invocation, + project_root, + server, + &graph, + &access, + session_store, + Arc::new(primitive_runtime::ResetRequiredSessionLookupV1), + )) + .await + .map(drop) +} + /// Registers code-index-independent owners for one newly inserted project. #[hotpath::measure(label = "daemon.project.owners.register", future = true)] #[cfg_attr( @@ -292,66 +404,23 @@ pub(super) async fn register_project_open_production_owners( })?; let grant_expires_at = access.grant_expires_at; let requester = access.requester.clone(); - // Primitive reads are part of the admitted core route. Publish their - // runtime and callable-code authorization before the slower mutation, - // delivery, native-integration, and Work owners finish mounting. - match hotpath::future!( - invocation.feedback_runtime_registrar().open_and_register( - database.clone(), - project_root.to_path_buf(), - graph.store_layout().response_handle_root.clone(), - scope.clone(), - access.clone(), - Arc::new(DaemonCallableCodeAuthorizationSource::production( - project_root.to_path_buf(), - scope.clone(), - Arc::clone(graph.configuration_runtime()), - )), - ), - label = "daemon.project.open.owners.feedback" - ) - .await - { - Ok(_) | Err(DaemonFeedbackRuntimeRegistrationError::AlreadyRegistered) => {} - Err(error) => { - return Err(TraceDecayError::Config { - message: format!("project-open feedback runtime registration failed: {error:?}"), - }); - } - } - tracing::info!( - event = "project_open_owner_phase", - project = %project_root.display(), - phase = "feedback_runtime_registered", - step_elapsed_ms = owner_phase_started.elapsed().as_millis(), - elapsed_ms = owner_registration_started.elapsed().as_millis(), - ); - owner_phase_started = Instant::now(); - - let admitted_root_uri = - admitted_root_uri_for_project(project_root).map_err(|error| TraceDecayError::Config { - message: format!("project-open admitted root URI denied: {error}"), - })?; - let source = graph - .source_read_context() - .ok_or_else(|| TraceDecayError::Config { - message: "project-open primitive runtime requires an exact registered source identity" - .to_owned(), - })?; - open_and_register_project_primitive_runtime( + let temporal = Arc::new(DaemonSessionLookupPrimitiveV1::new( + server.project_session_application_retrieval_service(&access.scope)?, + )); + let admitted_root_uri = Box::pin(register_project_code_read_owners( invocation, project_root, - source, server, + &graph, + &access, session_db.clone(), - access.clone(), - &admitted_root_uri, - ) + temporal, + )) .await?; tracing::info!( event = "project_open_owner_phase", project = %project_root.display(), - phase = "primitive_runtime_registered", + phase = "code_read_owners_registered", step_elapsed_ms = owner_phase_started.elapsed().as_millis(), elapsed_ms = owner_registration_started.elapsed().as_millis(), ); diff --git a/crates/tracedecay/src/daemon/project_open_owners/primitive_runtime.rs b/crates/tracedecay/src/daemon/project_open_owners/primitive_runtime.rs index fa57b7c91d..7ba1cb2001 100644 --- a/crates/tracedecay/src/daemon/project_open_owners/primitive_runtime.rs +++ b/crates/tracedecay/src/daemon/project_open_owners/primitive_runtime.rs @@ -5,15 +5,33 @@ use std::sync::Arc; use tracedecay_application::primitives::{ ProductionPrimitiveCodeAuthoritiesV1, ProductionPrimitiveOpenRequestV1, + admitted_root_uri_for_project, }; use tracedecay_application::source_authorization::ProjectSourceAccessSnapshot; use crate::daemon::DaemonInvocationState; use crate::mcp::McpServer; +use tracedecay_contracts::retrieval::{ + RetrievalPortContext, SessionLookupRequest, TemporalRetrievalFailure, TemporalRetrievalFuture, + TemporalRetrievalPort, +}; use tracedecay_daemon_service::DaemonPrimitiveRuntimeRegistrationError; use tracedecay_domain::errors::{Result, TraceDecayError}; use tracedecay_graph_query::SourceReadContext; -use tracedecay_session_runtime::session_retrieval::DaemonSessionLookupPrimitiveV1; + +/// Session lookups on a route whose session store is held in its typed +/// reset-required state: every lookup answers that refusal. +pub(super) struct ResetRequiredSessionLookupV1; + +impl TemporalRetrievalPort for ResetRequiredSessionLookupV1 { + fn session_lookup<'a>( + &'a self, + _context: RetrievalPortContext<'a>, + _request: &'a SessionLookupRequest, + ) -> TemporalRetrievalFuture<'a> { + Box::pin(async { Err(TemporalRetrievalFailure::ResetRequired) }) + } +} #[hotpath::measure(label = "daemon.project.owners.primitive", future = true)] pub(super) async fn open_and_register_project_primitive_runtime( @@ -22,9 +40,13 @@ pub(super) async fn open_and_register_project_primitive_runtime( source: SourceReadContext, server: &McpServer, session_db: tracedecay_global_db::RegisteredGlobalDbLeaseV1, + temporal: Arc, access: ProjectSourceAccessSnapshot, - admitted_root_uri: &str, -) -> Result<()> { +) -> Result { + let admitted_root_uri = + admitted_root_uri_for_project(project_root).map_err(|error| TraceDecayError::Config { + message: format!("project-open admitted root URI denied: {error}"), + })?; let code_graph = server .code_graph_projection_read_port() .ok_or_else(|| TraceDecayError::Config { @@ -37,9 +59,6 @@ pub(super) async fn open_and_register_project_primitive_runtime( message: "project-open primitive runtime requires the mounted ignored-dependency admission authority" .to_owned(), })?; - let temporal = Arc::new(DaemonSessionLookupPrimitiveV1::new( - server.project_session_application_retrieval_service(&access.scope)?, - )); invocation .primitive_runtime_registrar() .open_and_register( @@ -56,11 +75,12 @@ pub(super) async fn open_and_register_project_primitive_runtime( session_db, temporal, access, - admitted_root_uri.to_owned(), + admitted_root_uri.clone(), ), ) .await - .map_err(primitive_runtime_registration_error) + .map_err(primitive_runtime_registration_error)?; + Ok(admitted_root_uri) } fn primitive_runtime_registration_error( diff --git a/crates/tracedecay/src/daemon/project_routing.rs b/crates/tracedecay/src/daemon/project_routing.rs index 37f56d6b50..93f87c522d 100644 --- a/crates/tracedecay/src/daemon/project_routing.rs +++ b/crates/tracedecay/src/daemon/project_routing.rs @@ -139,11 +139,9 @@ pub(super) async fn bind_authenticated_profile_identity( // A reset-required profile authority keeps its registered location. // Binding the connection to it lets every request on it answer the // typed refusal instead of closing without a frame. - Err(error) if error.store_reset_required("profile authority").is_some() => { - authority::canonical_identity_path( - &profile_root.join(tracedecay_runtime_core::config::GLOBAL_DB_FILENAME), - )? - } + Err(error) if error.is_store_reset_required() => authority::canonical_identity_path( + &profile_root.join(tracedecay_runtime_core::config::GLOBAL_DB_FILENAME), + )?, Err(error) => return Err(error), }; let supplied_global_db_path = diff --git a/crates/tracedecay/src/daemon/projectless.rs b/crates/tracedecay/src/daemon/projectless.rs index 2087dd5940..51f91fc63f 100644 --- a/crates/tracedecay/src/daemon/projectless.rs +++ b/crates/tracedecay/src/daemon/projectless.rs @@ -356,7 +356,7 @@ async fn projectless_tools_call_response_with_connection( } if let Err(error) = boxed_projectless_phase(store_administration.ensure_account_active()).await { - if error.store_reset_required("profile authority").is_some() { + if error.is_store_reset_required() { return tool_error_response(id, tool_name, &error); } return JsonRpcResponse::error(id, ErrorCode::InternalError, error.to_string()); diff --git a/crates/tracedecay/src/daemon/remote_deletion.rs b/crates/tracedecay/src/daemon/remote_deletion.rs index b2c68e05d9..cbd3a949b0 100644 --- a/crates/tracedecay/src/daemon/remote_deletion.rs +++ b/crates/tracedecay/src/daemon/remote_deletion.rs @@ -40,7 +40,7 @@ pub(super) async fn resume_remote_account_deletion_for_boot( // daemon boots to serve that typed state, and every profile read, // including the account-active guard, meets the same refusal until // the operator's reset deletes the store. - Err(error) if error.store_reset_required("profile authority").is_some() => { + Err(error) if error.is_store_reset_required() => { log_daemon_event( "remote_account_deletion_resume", &[ diff --git a/crates/tracedecay/src/doctor.rs b/crates/tracedecay/src/doctor.rs index 2365ad64e9..bff7382369 100644 --- a/crates/tracedecay/src/doctor.rs +++ b/crates/tracedecay/src/doctor.rs @@ -217,7 +217,7 @@ fn render_current_project_daemon_status( Some(Err(error)) => { if let Some((authority, reason)) = tracedecay_mcp::reset_required_context(error) { *pending_reset = true; - dc.warn(&format!( + dc.pending(&format!( "Current project is not served: {authority} requires reset ({reason}). \ Pending operator action: run `{}`", tracedecay_mcp::reset_required_command(&authority, Some(project_path)) @@ -389,7 +389,7 @@ fn check_reset_required_stores( match tracedecay_daemon_control::daemon_reset_required_stores(profile, build_version) { Ok(stores) => { for store in &stores { - dc.warn(&format!( + dc.pending(&format!( "Store {} requires reset ({}). Pending operator action: run `{}`", store.store, store.reason, store.remedy )); diff --git a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance.rs b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance.rs index 3fce18ba1d..e6093fd070 100644 --- a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance.rs +++ b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance.rs @@ -38,9 +38,12 @@ use crate::common; /// A reset-required registered store served as a typed state. mod reset_required_serving; +/// Stale session stores refuse only sessions until their scoped reset. +mod stale_sessions_store_reset; /// The HTTP, MCP-host, and Rust SDK legs of this same journey. mod transport_boundaries; +use std::io::{BufRead, BufReader}; use std::path::Path; use std::process::Stdio; use std::time::{Duration, Instant}; @@ -60,6 +63,8 @@ const MAX_MARKER_TOKEN_BYTES: usize = 36; /// live; the barrier, not this number, decides when settlement happens. const PARTIAL_EFFECT_DEADLINE: Duration = Duration::from_secs(8); const BARRIER_ARRIVAL_TIMEOUT: Duration = Duration::from_secs(60); +/// Bound on the scoped reset reaching its wait for the profile lease. +const SCOPED_RESET_TIMEOUT: Duration = Duration::from_secs(120); /// Fails loudly if a marker could be refused by memory hygiene instead of /// committing, so a future marker edit cannot silently invalidate the @@ -327,6 +332,72 @@ fn assert_reset_required(payload: &Value, context: &str) { ); } +/// Runs the scoped reset the way an operator does while the daemon that +/// reported the stale stores is still serving: the reset reads the daemon's +/// reset census, then waits for the operator to stop the unmanaged test +/// daemon (a managed service is stopped by the reset itself). `between` runs +/// after the daemon exited and before the reset can take the profile, so it +/// observes exactly the bytes the reset starts from. +fn run_scoped_reset( + home: &Path, + project: &Path, + daemon: &mut common::DaemonProcess, + between: impl FnOnce(), +) -> (std::process::ExitStatus, String) { + let profile_root = home.join(".tracedecay"); + let hold = tracedecay_runtime_core::lifecycle_lease::acquire_shared_for_profile( + &profile_root, + "stale-sessions-store-reset", + ) + .expect("hold the profile while the daemon stops"); + let mut child = tracedecay_command_with_home(home) + .args(["wipe", "--stale", "--yes"]) + .current_dir(project) + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .expect("spawn the scoped reset"); + let stdout = child.stdout.take().expect("piped stdout"); + let stdout_reader = std::thread::spawn(move || { + let mut text = String::new(); + for line in BufReader::new(stdout).lines() { + text.push_str(&line.expect("reset stdout line")); + text.push('\n'); + } + text + }); + let mut stderr = BufReader::new(child.stderr.take().expect("piped stderr")); + let mut transcript = String::new(); + let started = Instant::now(); + loop { + let mut line = String::new(); + let read = stderr.read_line(&mut line).expect("read reset stderr"); + transcript.push_str(&line); + if read == 0 { + panic!("the scoped reset exited before waiting for the profile:\n{transcript}"); + } + if line.contains("stop it and wipe --stale continues") { + break; + } + assert!( + started.elapsed() < SCOPED_RESET_TIMEOUT, + "the scoped reset never waited for the profile:\n{transcript}" + ); + } + daemon + .kill_and_wait() + .expect("stop the daemon for the scoped reset"); + between(); + drop(hold); + let mut rest = String::new(); + std::io::Read::read_to_string(&mut stderr, &mut rest).expect("read reset stderr"); + transcript.push_str(&rest); + let status = child.wait().expect("wait for the scoped reset"); + transcript.push_str(&stdout_reader.join().expect("join reset stdout")); + (status, transcript) +} + fn spawn_daemon_with_commit_barrier(home: &Path, barrier_dir: &Path) -> common::DaemonProcess { let barrier_dir = barrier_dir.to_path_buf(); spawn_tracedecay_daemon_with(home, move |command| { diff --git a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/reset_required_serving.rs b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/reset_required_serving.rs index e4bb4b618f..a0f1430a6b 100644 --- a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/reset_required_serving.rs +++ b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/reset_required_serving.rs @@ -11,14 +11,11 @@ //! still answer. After the named reset the same read serves. use std::path::Path; -use std::process::Stdio; use std::time::{Duration, Instant}; use serde_json::{Value, json}; -use crate::common::{ - canonical_existing_path, spawn_tracedecay_daemon_with, tracedecay_command_with_home, -}; +use crate::common::{canonical_existing_path, spawn_tracedecay_daemon_with}; /// Bound on waiting for the first sealed code generation of a two-line fixture. const CODE_INDEX_READY_TIMEOUT: Duration = Duration::from_secs(120); @@ -149,7 +146,7 @@ fn reset_required_profile_session_store_is_served_typed_until_its_named_reset() "required_version": 6, "reason": "git correlation profile schema 5 is incompatible with required schema 6; \ reset the profile", - "remedy": "tracedecay wipe --all --yes", + "remedy": "tracedecay wipe --stale --yes", }]) ); let refused = super::cli_problem_envelope( @@ -166,28 +163,16 @@ fn reset_required_profile_session_store_is_served_typed_until_its_named_reset() "required_version": 6, "reason": "git correlation profile schema 5 is incompatible with required schema 6; \ reset the profile", - "remedy": "tracedecay wipe --all --yes", + "remedy": "tracedecay wipe --stale --yes", }) ); wait_for_code_index_hit(&home_path, &project_path, "probe"); - // `tracedecay wipe --all --yes` takes the profile offline by stopping the - // managed service and restarts it afterwards; this unmanaged daemon is - // stopped and started around the same command. - daemon - .kill_and_wait() - .expect("stop the daemon for the profile reset"); - let wipe = tracedecay_command_with_home(&home_path) - .args(["wipe", "--all", "--yes"]) - .current_dir(&project_path) - .stdin(Stdio::null()) - .output() - .expect("run the named reset"); + let (reset_status, reset_output) = + super::run_scoped_reset(&home_path, &project_path, &mut daemon, || {}); assert!( - wipe.status.success(), - "tracedecay wipe --all --yes failed\nstdout:\n{}\nstderr:\n{}", - String::from_utf8_lossy(&wipe.stdout), - String::from_utf8_lossy(&wipe.stderr) + reset_status.success(), + "tracedecay wipe --stale --yes failed:\n{reset_output}" ); let mut daemon = spawn_tracedecay_daemon_with(&home_path, |_| {}); diff --git a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs new file mode 100644 index 0000000000..d5a30bbae7 --- /dev/null +++ b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs @@ -0,0 +1,519 @@ +//! Session stores whose observation rows predate the unified observation +//! identity are refused as typed reset states without taking code +//! intelligence down with them, and their scoped reset deletes exactly those +//! stores. +//! +//! A physically spawned `tracedecay daemon run` first writes a real profile +//! and project session store. Both are then given the shape a released binary +//! left behind: observation rows written before the unified identity, with no +//! unified-identity marker. Over that profile the project must still open, the +//! MCP host must initialize and list tools, code search must answer, session +//! reads must return the typed `reset_required` refusal naming +//! `tracedecay wipe --stale --yes`, and `tracedecay doctor` must name that +//! command instead of `tracedecay install`. The scoped reset then deletes both +//! session stores and nothing else, and the restarted daemon serves sessions +//! from empty stores. + +use std::collections::BTreeMap; +use std::io::Write; +use std::path::{Path, PathBuf}; +use std::process::Stdio; +use std::time::{Duration, Instant}; + +use serde_json::{Value, json}; +use sha2::{Digest, Sha256}; + +use crate::common::{ + TestChildProcess, canonical_existing_path, spawn_tracedecay_daemon_with, + tracedecay_command_with_home, +}; + +const STALE_STORE_RESET: &str = "tracedecay wipe --stale --yes"; +const OBSERVATIONS_RESET_REASON: &str = "observations persisted shape requires reset: \ + observation rows predate the unified observation identity and cannot be read; reset the \ + profile so ingestion can rebuild them from host transcripts"; +const PRE_UNIFIED_OBSERVATION_ID: &str = + "sha256:efd99c7fd87f4ad156b40f16d982d18511ebfb708afc140f9f67e63e0c73f5ba"; +const PRE_UNIFIED_RECEIPT_ID: &str = + "privacy.claude.v1.0000000000000000000000000000000000000000000000000000000000000000"; +/// A Claude observation as binaries before the unified identity wrote it: +/// its id derives from `tracedecay.claude.observation.v1` and its JSON still +/// carries `idempotency_key`. +const PRE_UNIFIED_OBSERVATION_JSON: &str = concat!( + r#"{"observation_id":""#, + "sha256:efd99c7fd87f4ad156b40f16d982d18511ebfb708afc140f9f67e63e0c73f5ba", + r#"","idempotency_key":""#, + "sha256:efd99c7fd87f4ad156b40f16d982d18511ebfb708afc140f9f67e63e0c73f5ba", + r#"","identity":{"source":{"provider":"claude","session_id":"session.fixture"}},"#, + r#""receipt":{"receipt_id":""#, + "privacy.claude.v1.0000000000000000000000000000000000000000000000000000000000000000", + r#""},"retention_class":"transcript.fixture","payload":{"message":"safe"}}"# +); +/// Bound on the daemon recording both refused stores after it restarts over +/// them: the project store refuses on project open, the profile store on the +/// full-route upgrade that follows. +const RESET_CENSUS_TIMEOUT: Duration = Duration::from_secs(60); +const SERVE_TIMEOUT: Duration = Duration::from_secs(90); + +/// Gives an existing session store the persisted shape a released binary +/// wrote before the unified observation identity: one observation row and no +/// unified-identity marker. +fn seed_pre_unified_observation_rows(db_path: &Path) { + let connection = rusqlite::Connection::open(db_path).expect("open the session store"); + let unmarked = connection + .execute( + "DELETE FROM global_schema_migrations \ + WHERE migration = 'observations-unified-identity-v1'", + [], + ) + .expect("remove the unified-identity marker"); + assert_eq!(unmarked, 1, "{} records the marker once", db_path.display()); + connection + .execute( + "INSERT INTO sanitization_receipts \ + (receipt_id, sanitizer_version, payload_digest, receipt_json) \ + VALUES (?1, 'privacy.claude-record.v1', ?2, '{}')", + [PRE_UNIFIED_RECEIPT_ID, PRE_UNIFIED_OBSERVATION_ID], + ) + .expect("seed the pre-unified receipt"); + connection + .execute( + "INSERT INTO observations \ + (observation_id, payload_digest, receipt_id, observation_json, \ + committed_cursor_json) VALUES (?1, ?1, ?2, ?3, '{}')", + [ + PRE_UNIFIED_OBSERVATION_ID, + PRE_UNIFIED_RECEIPT_ID, + PRE_UNIFIED_OBSERVATION_JSON, + ], + ) + .expect("seed the pre-unified observation"); +} + +/// SHA-256 of every regular file under `root`, keyed by relative path. +fn file_digests(root: &Path) -> BTreeMap { + fn walk(root: &Path, directory: &Path, digests: &mut BTreeMap) { + for entry in std::fs::read_dir(directory).expect("read profile directory") { + let path = entry.expect("profile directory entry").path(); + let file_type = std::fs::symlink_metadata(&path) + .expect("inspect profile entry") + .file_type(); + if file_type.is_dir() { + walk(root, &path, digests); + } else if file_type.is_file() { + let bytes = std::fs::read(&path).expect("read profile file"); + digests.insert( + path.strip_prefix(root) + .expect("relative path") + .to_path_buf(), + hex::encode(Sha256::digest(&bytes)), + ); + } + } + } + let mut digests = BTreeMap::new(); + walk(root, root, &mut digests); + digests +} + +/// Whether `relative` belongs to one of the two refused session stores. +fn is_session_store_member(relative: &Path, project_store: &Path) -> bool { + let name = relative.to_string_lossy(); + let in_project_sessions = relative.starts_with(project_store) + && relative + .strip_prefix(project_store) + .ok() + .and_then(|rest| rest.components().next()) + .is_some_and(|first| { + let first = first.as_os_str().to_string_lossy(); + first.starts_with("sessions.") || first == ".sessions.db.host-admission" + }); + in_project_sessions + || name.starts_with("user-sessions.") + || name.starts_with(".user-sessions.db.host-admission") +} + +fn status(home: &Path, project: &Path) -> Value { + let status = super::tool_call( + home, + project, + "tracedecay_status", + &json!({ "format": "json" }), + ); + super::typed_envelope(&status) +} + +fn find_key(value: &Value, key: &str) -> Option { + match value { + Value::Object(map) => map + .get(key) + .cloned() + .or_else(|| map.values().find_map(|child| find_key(child, key))), + Value::Array(items) => items.iter().find_map(|child| find_key(child, key)), + Value::String(text) => serde_json::from_str::(text) + .ok() + .and_then(|parsed| find_key(&parsed, key)), + _ => None, + } +} + +fn reset_required_stores(home: &Path, project: &Path) -> Vec { + let status = status(home, project); + let mut stores = find_key(&status, "reset_required_stores") + .and_then(|stores| stores.as_array().cloned()) + .unwrap_or_else(|| panic!("tracedecay_status omitted reset_required_stores: {status}")); + stores.sort_by(|left, right| left["store"].as_str().cmp(&right["store"].as_str())); + stores +} + +fn wait_for_reset_required_stores(home: &Path, project: &Path, expected: &[Value]) { + let started = Instant::now(); + loop { + let stores = reset_required_stores(home, project); + if stores == expected { + return; + } + assert!( + started.elapsed() < RESET_CENSUS_TIMEOUT, + "the daemon never recorded exactly the refused session stores within \ + {RESET_CENSUS_TIMEOUT:?}\nexpected: {expected:#?}\nobserved: {stores:#?}" + ); + std::thread::sleep(Duration::from_millis(500)); + } +} + +fn names_symbol(value: &Value, name: &str) -> bool { + match value { + Value::Object(map) => { + map.get("name").and_then(Value::as_str) == Some(name) + || map.values().any(|child| names_symbol(child, name)) + } + Value::Array(items) => items.iter().any(|child| names_symbol(child, name)), + Value::String(text) => { + serde_json::from_str::(text).is_ok_and(|parsed| names_symbol(&parsed, name)) + } + _ => false, + } +} + +fn wait_for_code_index_hit(home: &Path, project: &Path, symbol: &str) { + let started = Instant::now(); + loop { + let result = super::tool_call( + home, + project, + "tracedecay_search", + &json!({ "query": symbol, "format": "json" }), + ); + if names_symbol(&result, symbol) { + return; + } + assert!( + started.elapsed() < Duration::from_secs(120), + "code search never answered `{symbol}`: {result}" + ); + std::thread::sleep(Duration::from_millis(500)); + } +} + +/// `initialize` then `tools/list` through a real `tracedecay serve` host. +fn mcp_initialize_and_list_tools(home: &Path, project: &Path) -> (Value, Value) { + let mut command = tracedecay_command_with_home(home); + command + .arg("serve") + .arg("--path") + .arg(project) + .current_dir(project) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()); + let mut child = TestChildProcess::new(command.spawn().expect("spawn tracedecay serve")); + { + let stdin = child.stdin_mut().expect("MCP host stdin is piped"); + for request in [ + json!({ + "jsonrpc": "2.0", + "id": 1, + "method": "initialize", + "params": { + "protocolVersion": "2024-11-05", + "capabilities": {}, + "clientInfo": { "name": "stale-sessions-store-reset", "version": "0.0.0" } + } + }), + json!({ "jsonrpc": "2.0", "method": "notifications/initialized" }), + json!({ "jsonrpc": "2.0", "id": 2, "method": "tools/list" }), + ] { + writeln!(stdin, "{request}").expect("write MCP request"); + } + } + let output = child + .wait_with_output(SERVE_TIMEOUT) + .expect("the MCP host exits after stdin closes"); + let stdout = String::from_utf8_lossy(&output.stdout).into_owned(); + let response = |id: i64| { + stdout + .lines() + .filter_map(|line| serde_json::from_str::(line).ok()) + .find(|message| message.get("id") == Some(&json!(id))) + .unwrap_or_else(|| { + panic!( + "the MCP host returned no response {id}\nstdout:\n{stdout}\nstderr:\n{}", + String::from_utf8_lossy(&output.stderr) + ) + }) + }; + (response(1), response(2)) +} + +/// The outcome kind a code-read primitive answered, or the problem it +/// refused with. +fn code_read_outcome(home: &Path, project: &Path, tool: &str, args: &Value) -> Value { + let envelope = super::typed_envelope(&super::tool_call(home, project, tool, args)); + if envelope["problem"].is_object() { + return json!({ "problem": envelope["problem"]["kind"] }); + } + envelope["outcome"]["outcome"].clone() +} + +fn probe_symbol_id(home: &Path, project: &Path) -> String { + let found = super::typed_envelope(&super::tool_call( + home, + project, + "tracedecay_find_exact_symbol", + &json!({ "name": "probe", "format": "json" }), + )); + found["matches"][0]["id"] + .as_str() + .unwrap_or_else(|| panic!("find_exact_symbol did not resolve `probe`: {found}")) + .to_owned() +} + +fn session_status(home: &Path, project: &Path, storage_scope: &str) -> Value { + super::tool_call( + home, + project, + "tracedecay_lcm_status", + &json!({ "storage_scope": storage_scope, "format": "json" }), + ) +} + +#[test] +fn stale_session_stores_refuse_sessions_only_until_their_scoped_reset() { + let home = tempfile::TempDir::new().expect("isolated home"); + let home_path = canonical_existing_path(home.path()); + let project = tempfile::TempDir::new().expect("project"); + let project_path = canonical_existing_path(project.path()); + let profile_root = home_path.join(".tracedecay"); + let project_id = tracedecay_runtime_core::storage::default_profile_project_id(&project_path); + let project_store = PathBuf::from("projects").join(&project_id); + + let mut daemon = spawn_tracedecay_daemon_with(&home_path, |_| {}); + super::initialize_project(&home_path, &project_path, "stale-sessions-store-reset"); + wait_for_code_index_hit(&home_path, &project_path, "probe"); + for scope in ["project", "user"] { + let served = session_status(&home_path, &project_path, scope); + assert!( + find_key(&served, "problem").is_none_or(|problem| problem.is_null()), + "the fresh {scope} session store serves before it is aged: {served}" + ); + } + daemon + .kill_and_wait() + .expect("stop the daemon that wrote the profile"); + + seed_pre_unified_observation_rows(&profile_root.join("user-sessions.db")); + seed_pre_unified_observation_rows(&profile_root.join(&project_store).join("sessions.db")); + + let mut daemon = spawn_tracedecay_daemon_with(&home_path, |_| {}); + + let (initialize, tools) = mcp_initialize_and_list_tools(&home_path, &project_path); + assert_eq!( + initialize.get("error"), + None, + "MCP initialize must serve over stale session stores: {initialize}" + ); + let tool_names: Vec<&str> = tools["result"]["tools"] + .as_array() + .unwrap_or_else(|| panic!("tools/list returned no catalog: {tools}")) + .iter() + .filter_map(|tool| tool["name"].as_str()) + .collect(); + for code_tool in [ + "tracedecay_search", + "tracedecay_callers", + "tracedecay_status", + ] { + assert!( + tool_names.contains(&code_tool), + "tools/list omitted {code_tool}: {tools}" + ); + } + + wait_for_code_index_hit(&home_path, &project_path, "probe"); + let probe_id = probe_symbol_id(&home_path, &project_path); + assert_eq!( + code_read_outcome( + &home_path, + &project_path, + "tracedecay_callers", + &json!({ "node_id": probe_id, "format": "json" }), + ), + json!("evidence"), + "callers must serve over stale session stores" + ); + assert_eq!( + code_read_outcome( + &home_path, + &project_path, + "tracedecay_file_dependents", + &json!({ "file": "src/lib.rs", "format": "json" }), + ), + json!("evidence"), + "file dependents must serve over stale session stores" + ); + let expected_stale = vec![ + json!({ + "store": "profile sessions", + "authority": "observations", + "found_version": null, + "required_version": null, + "reason": OBSERVATIONS_RESET_REASON, + "remedy": STALE_STORE_RESET, + }), + json!({ + "store": format!("project sessions {project_id}"), + "authority": "observations", + "found_version": null, + "required_version": null, + "reason": OBSERVATIONS_RESET_REASON, + "remedy": STALE_STORE_RESET, + }), + ]; + wait_for_reset_required_stores(&home_path, &project_path, &expected_stale); + let project_open = find_key(&status(&home_path, &project_path), "project_open"); + assert!( + project_open + .as_ref() + .is_none_or(|open| open.is_null() || open["state"] == "completed"), + "project open must not stall on a session-store verdict: {project_open:?}" + ); + + for scope in ["project", "user"] { + let refused = super::cli_problem_envelope( + &session_status(&home_path, &project_path, scope), + &format!("{scope} session read over a stale store"), + ); + super::assert_reset_required(&refused, &format!("{scope} session read")); + assert_eq!( + refused["problem"]["detail"]["authority"], "observations", + "{scope} session read names the refused authority: {refused}" + ); + assert_eq!( + refused["problem"]["detail"]["remedy"], STALE_STORE_RESET, + "{scope} session read names the scoped reset: {refused}" + ); + } + + let doctor = tracedecay_command_with_home(&home_path) + .arg("doctor") + .current_dir(&project_path) + .stdin(Stdio::null()) + .output() + .expect("run doctor"); + let doctor_text = format!( + "{}{}", + String::from_utf8_lossy(&doctor.stdout), + String::from_utf8_lossy(&doctor.stderr) + ); + for store in [ + "profile sessions".to_owned(), + format!("project sessions {project_id}"), + ] { + let line = format!( + "Store {store} requires reset ({OBSERVATIONS_RESET_REASON}). Pending operator \ + action: run `{STALE_STORE_RESET}`" + ); + assert!( + doctor_text.contains(&line), + "doctor omitted `{line}`:\n{doctor_text}" + ); + } + assert!( + doctor_text.contains("pending operator action(s)") && doctor_text.contains("no issues."), + "stale session stores are pending operator actions, not issues:\n{doctor_text}" + ); + assert!( + !doctor_text.contains("to fix most issues"), + "doctor must not send a stale session store to `tracedecay install`:\n{doctor_text}" + ); + assert!( + !doctor_text.contains("Stalled"), + "doctor must not report project open stalled on a session store:\n{doctor_text}" + ); + + let mut before_reset = BTreeMap::new(); + let (reset_status, reset_output) = + super::run_scoped_reset(&home_path, &project_path, &mut daemon, || { + before_reset = file_digests(&profile_root); + }); + assert!( + reset_status.success(), + "the scoped reset failed:\n{reset_output}" + ); + let after_reset = file_digests(&profile_root); + let (stale_members, kept): (BTreeMap<_, _>, BTreeMap<_, _>) = before_reset + .into_iter() + .partition(|(relative, _)| is_session_store_member(relative, &project_store)); + assert!( + stale_members.contains_key(Path::new("user-sessions.db")) + && stale_members.contains_key(&project_store.join("sessions.db")), + "both refused stores existed before the reset: {stale_members:#?}" + ); + for relative in stale_members.keys() { + assert!( + !after_reset.contains_key(relative), + "the scoped reset left refused store member {} behind", + relative.display() + ); + } + // The lifecycle lock records whichever command holds the profile lease; + // it is coordination, not stored data. + let changed: Vec<_> = kept + .iter() + .filter(|(relative, _)| relative.as_path() != Path::new("lifecycle.lock")) + .filter(|(relative, digest)| after_reset.get(*relative) != Some(digest)) + .map(|(relative, _)| relative.clone()) + .collect(); + assert_eq!( + changed, + Vec::::new(), + "the scoped reset touched files outside the refused stores:\n{reset_output}" + ); + for store in [ + "profile sessions".to_owned(), + format!("project sessions {project_id}"), + ] { + assert!( + reset_output.contains(&format!("reset {store}")), + "the scoped reset did not report resetting {store}:\n{reset_output}" + ); + } + + let mut daemon = spawn_tracedecay_daemon_with(&home_path, |_| {}); + wait_for_code_index_hit(&home_path, &project_path, "probe"); + for scope in ["project", "user"] { + let served = session_status(&home_path, &project_path, scope); + assert!( + find_key(&served, "problem").is_none_or(|problem| problem.is_null()), + "the {scope} session store must serve after the scoped reset: {served}" + ); + } + assert_eq!( + reset_required_stores(&home_path, &project_path), + Vec::::new(), + "no store stays refused after the scoped reset" + ); + + let _ = daemon.kill_and_wait(); +}