diff --git a/Cargo.lock b/Cargo.lock index 358b0f79..786558ad 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1876,6 +1876,7 @@ checksum = "20870f649af7073d53e38067b2a84312175d56ea15217e1b15bc83506ec50afb" name = "libatomic" version = "0.19.1" dependencies = [ + "anyhow", "atomic-agent", "atomic-canonical", "atomic-core", @@ -1892,13 +1893,16 @@ dependencies = [ "prost", "prost-types", "protox", + "redb 4.2.0", "serde_json", "similar", + "tempfile", "tokio", "tokio-stream", "tonic", "tonic-build", "ulid", + "uuid", ] [[package]] diff --git a/README.md b/README.md index 951065c1..1e5f1a6f 100644 --- a/README.md +++ b/README.md @@ -84,6 +84,36 @@ Current operational considerations: - Unix endpoints use `/tmp` to stay below macOS Unix-socket path limits. A future hardening step should move them into a user-private runtime directory where available, enforce restrictive socket permissions, and validate peer credentials. - A 96-bit endpoint digest makes accidental cross-project collisions extraordinarily unlikely, but the repository-local lock remains the final ownership check. +### Reclaiming unused database space + +Run explicit maintenance when the repository is quiet: + +```bash +atomic compact +atomic compact --repository /path/to/repo --json +``` + +The command compacts the existing `.atomic/atomic.redb` and reports its +before/after sizes and reclaimed bytes. It preserves history, views, +provenance, change files and unrecorded working files. It does not delete +live data or deduplicate change objects; zero bytes reclaimed is a valid result. + +Compaction runs through libatomic, in-process by default. With +`ATOMIC_SERVICE=reactor`, it uses the same handler through the configured +Reactor service; that service must advertise `CompactDatabase`. An older or +unavailable service produces an error, with no local fallback. + +Other database readers and writers must release their handles. Lock waiting +uses `ATOMIC_DB_LOCK_WAIT_MS` (30 seconds by default); a busy database produces +an error that can be retried. Once blocking maintenance starts, it finishes +independently of the caller, including the bounded lock wait. Use a quiet period: +other commands may time out while the database is being compacted. Persistent +savepoints block compaction and are never deleted automatically. + +This is manual maintenance of an existing combined database, with no implicit +initialization or migration. Hosted Storage's open-handle cache requires separate +maintenance coordination before online compaction can be supported. + ### Provenance Graphs Every agent session builds a causal decision DAG. Not just *what* changed, but *why*: diff --git a/atomic-cli/Cargo.toml b/atomic-cli/Cargo.toml index 26ebfb9a..60900cda 100644 --- a/atomic-cli/Cargo.toml +++ b/atomic-cli/Cargo.toml @@ -95,5 +95,6 @@ git2 = { workspace = true } windows-sys = { version = "0.61", features = ["Win32_Foundation"] } [dev-dependencies] +tokio-stream = { workspace = true, features = ["net"] } tempfile = { workspace = true } serial_test = { workspace = true } diff --git a/atomic-cli/src/commands/compact.rs b/atomic-cli/src/commands/compact.rs new file mode 100644 index 00000000..f315add1 --- /dev/null +++ b/atomic-cli/src/commands/compact.rs @@ -0,0 +1,52 @@ +//! Explicit database maintenance through the shared service layer. + +use crate::commands::Command; +use crate::error::{CliError, CliResult}; +use crate::service::Service; +use atomic_client::proto::{CompactDatabaseRequest, RequestMeta}; +use clap::Args; +use std::path::PathBuf; + +/// Reclaim unused database space without deleting repository data. +#[derive(Debug, Args)] +pub struct Compact { + /// Path inside the repository or one of its sandboxes. + #[arg(long, default_value = ".")] + repository: PathBuf, + /// Emit the database path and before/after byte counts as JSON. + #[arg(long)] + json: bool, +} + +impl Command for Compact { + fn run(&self) -> CliResult<()> { + let service = + Service::open_root(&self.repository)?.ok_or_else(|| CliError::RepositoryNotFound { + searched_path: self.repository.clone(), + })?; + let report = service.compact_database(CompactDatabaseRequest { + repository: Some(service.reference.clone()), + meta: Some(RequestMeta { + request_id: uuid::Uuid::new_v4().to_string(), + observed_at: None, + }), + })?; + if self.json { + println!( + "{}", + serde_json::json!({ + "database": report.database, + "before_bytes": report.before_bytes, + "after_bytes": report.after_bytes, + "reclaimed_bytes": report.reclaimed_bytes, + }) + ); + } else { + println!( + "Compacted {}: {} -> {} bytes ({} bytes reclaimed)", + report.database, report.before_bytes, report.after_bytes, report.reclaimed_bytes + ); + } + Ok(()) + } +} diff --git a/atomic-cli/src/commands/mod.rs b/atomic-cli/src/commands/mod.rs index 4e758131..ae5f6dbd 100644 --- a/atomic-cli/src/commands/mod.rs +++ b/atomic-cli/src/commands/mod.rs @@ -69,6 +69,7 @@ use crate::error::{CliError, CliResult}; // Phase 2: Core Local Commands pub mod add; pub mod change; +pub mod compact; pub mod complete; pub mod completions; pub mod conflicts; @@ -137,6 +138,7 @@ pub use add::Add; pub use agent::Agent; pub use change::ChangeCmd; pub use clone::Clone; +pub use compact::Compact; pub use completions::Completions; pub use conflicts::Conflicts; pub use diff::Diff; diff --git a/atomic-cli/src/main.rs b/atomic-cli/src/main.rs index b15eb341..c9174172 100644 --- a/atomic-cli/src/main.rs +++ b/atomic-cli/src/main.rs @@ -62,6 +62,7 @@ use commands::{ ChangeCmd, Clone, Command, + Compact, Completions, Conflicts, Diff, @@ -447,6 +448,9 @@ enum Commands { /// such as backfilling the dependency index for legacy repositories. Doctor(Doctor), + /// Reclaim unused space in the repository database. + Compact(Compact), + /// Git interoperability commands. /// /// Import Git repositories into Atomic, preserving history, authorship, @@ -996,6 +1000,7 @@ fn main() { Commands::Diff(diff) => diff.run(), Commands::Doctor(doctor) => doctor.run(), + Commands::Compact(compact) => compact.run(), Commands::Git(git) => git.run(), diff --git a/atomic-cli/src/service/compact.rs b/atomic-cli/src/service/compact.rs new file mode 100644 index 00000000..ea979b98 --- /dev/null +++ b/atomic-cli/src/service/compact.rs @@ -0,0 +1,49 @@ +use super::{pb, services_maintenance, Backend, Service}; +use crate::error::CliResult; +use libatomic::atomic::maintenance_service_server::MaintenanceService as _; +use tonic::{Request, Status}; + +impl Service { + pub fn compact_database( + &self, + request: pb::CompactDatabaseRequest, + ) -> CliResult { + self.call(move |backend| async move { + let mut request = Request::new(request); + request + .metadata_mut() + .insert("x-atomic-contract-version", "2".parse().unwrap()); + match backend { + Backend::Local(state) => services_maintenance::MaintenanceImpl { state } + .compact_database(request) + .await + .map(|response| response.into_inner()), + Backend::Reactor(channel) => { + let mut daemon = + pb::daemon_service_client::DaemonServiceClient::new(channel.clone()); + let capabilities = daemon + .get_capabilities(pb::GetCapabilitiesRequest {}) + .await? + .into_inner(); + if capabilities.protocol_version != 2 { + return Err(Status::failed_precondition( + "the running Reactor uses an incompatible service contract; update/restart it and retry", + )); + } + let supported = capabilities.methods.iter().any(|method| { + method.service == "MaintenanceService" && method.method == "CompactDatabase" + }); + if !supported { + return Err(Status::failed_precondition( + "the running Reactor does not support database compaction; update/restart it and retry", + )); + } + pb::maintenance_service_client::MaintenanceServiceClient::new(channel) + .compact_database(request) + .await + .map(|response| response.into_inner()) + } + } + }) + } +} diff --git a/atomic-cli/src/service/mod.rs b/atomic-cli/src/service/mod.rs index a4296065..8b4c1130 100644 --- a/atomic-cli/src/service/mod.rs +++ b/atomic-cli/src/service/mod.rs @@ -56,6 +56,8 @@ use tonic::{Request, Status}; use crate::error::{CliError, CliResult}; +mod compact; + /// The service-area selector. `ATOMIC_SERVICE` wins; `ATOMIC_RPC` is the /// legacy alias (1 → reactor, anything else → local); unset routes local. pub const ENV_SERVICE: &str = "ATOMIC_SERVICE"; diff --git a/atomic-cli/tests/compact_service_integration_test.rs b/atomic-cli/tests/compact_service_integration_test.rs new file mode 100644 index 00000000..8907e39d --- /dev/null +++ b/atomic-cli/tests/compact_service_integration_test.rs @@ -0,0 +1,198 @@ +//! Exercise the real CLI through the shared in-process and served handlers. +use std::path::Path; +use std::process::{Command, Output}; + +use atomic_repository::{ChangeStore, Repository, DEFAULT_CACHE_CAPACITY}; +use serde_json::Value; + +fn command(path: &Path) -> Command { + let mut command = Command::new(env!("CARGO_BIN_EXE_atomic")); + command + .args(["compact", "--repository"]) + .arg(path) + .arg("--json") + .env("ATOMIC_SERVICE", "local") + .env_remove("ATOMIC_RPC") + .env("ATOMIC_DB_LOCK_WAIT_MS", "100") + .env("ATOMIC_DAEMON_BIN", path.join("no-daemon-binary")) + .env("ATOMIC_DAEMON_SOCKET", path.join("no-daemon.socket")); + command +} + +fn report(output: Output) -> Value { + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + serde_json::from_slice(&output.stdout).unwrap() +} + +#[test] +fn local_compaction_preserves_views_journal_changes_and_dirty_sandbox() { + let dir = tempfile::tempdir().unwrap(); + let root = dir.path().join("repo"); + let sandbox = dir.path().join("sandbox"); + let mut repo = Repository::init(&root).unwrap(); + let contents = b"recorded contents\n"; + std::fs::write(root.join("file.txt"), contents).unwrap(); + repo.add("file.txt", Default::default()).unwrap(); + let hash = *repo.record_all("compaction fixture").unwrap().hash(); + let change_path = ChangeStore::new(repo.changes_dir(), DEFAULT_CACHE_CAPACITY) + .unwrap() + .change_path(&hash); + let change_bytes = std::fs::read(&change_path).unwrap(); + let view = repo.current_view().to_owned(); + repo.create_view("compact-child").unwrap(); + let views = repo.list_views().unwrap(); + repo.provision_sandbox(&sandbox, &view).unwrap(); + let store = repo.redb_change_store().unwrap(); + let turn = store.reserve_provenance_turn("compact", 1, 1).unwrap(); + let envelope = serde_json::to_vec(&serde_json::json!({ + "schema_version":1,"event_id":"pending","session_id":"compact", + "turn_number":1,"generation":turn.generation,"timestamp_ms":1700000000000_i64, + "event":{"type":"tool","phase":"after","tool_name":"Read", + "tool_call_id":"pending","input":{"path":"file.txt"}, + "output":"recorded contents","status":"completed"} + })) + .unwrap(); + store + .append_provenance_envelope(turn.provenance_id, turn.generation, "pending", &envelope, 2) + .unwrap(); + let pending = store.load_provenance_envelopes(turn.provenance_id).unwrap(); + let turn_before = store.get_provenance_turn(turn.provenance_id).unwrap(); + drop(store); + drop(repo); + std::fs::write(root.join("file.txt"), b"dirty main\n").unwrap(); + std::fs::write(sandbox.join("file.txt"), b"dirty sandbox\n").unwrap(); + let database = root.join(".atomic/atomic.redb").canonicalize().unwrap(); + let before = database.metadata().unwrap().len(); + let log = dir.path().join("service.jsonl"); + let result = report( + command(&sandbox) + .env("ATOMIC_DAEMON_LOG_REQUESTS", &log) + .output() + .unwrap(), + ); + assert_eq!(result["database"], database.to_str().unwrap()); + assert_eq!(result["before_bytes"], before); + let after = database.metadata().unwrap().len(); + assert_eq!(result["after_bytes"], after); + assert_eq!(result["reclaimed_bytes"], before.saturating_sub(after)); + assert!(std::fs::read_to_string(log) + .unwrap() + .contains("CompactDatabase")); + report(command(&root).output().unwrap()); // Repeat maintenance is valid. + let repo = Repository::open_existing(&root).unwrap(); + assert_eq!(repo.list_views().unwrap(), views); + assert_eq!(repo.log(Default::default()).unwrap()[0].hash, hash); + for view in [&view, "compact-child"] { + assert_eq!( + repo.get_file_content_on_view("file.txt", view) + .unwrap() + .unwrap(), + contents + ); + } + assert_eq!(std::fs::read(change_path).unwrap(), change_bytes); + assert_eq!( + std::fs::read(root.join("file.txt")).unwrap(), + b"dirty main\n" + ); + assert_eq!( + std::fs::read(sandbox.join("file.txt")).unwrap(), + b"dirty sandbox\n" + ); + let store = repo.redb_change_store().unwrap(); + assert_eq!( + store.load_provenance_envelopes(turn.provenance_id).unwrap(), + pending + ); + assert_eq!( + store.get_provenance_turn(turn.provenance_id).unwrap(), + turn_before + ); + assert!(!sandbox.join(".atomic/atomic.redb").exists()); +} + +#[test] +fn busy_database_reports_actionable_error_and_can_retry() { + let dir = tempfile::tempdir().unwrap(); + let repo = Repository::init(dir.path()).unwrap(); + let output = command(dir.path()).output().unwrap(); + assert!(!output.status.success()); + assert!(String::from_utf8_lossy(&output.stderr).contains("timed out")); + drop(repo); + report(command(dir.path()).output().unwrap()); +} + +#[test] +fn missing_repository_is_not_created() { + let dir = tempfile::tempdir().unwrap(); + assert!(!command(dir.path()).output().unwrap().status.success()); + assert!(!dir.path().join(".atomic").exists()); +} + +#[cfg(unix)] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn reactor_cli_uses_same_handler_over_socket_without_fallback() { + use libatomic::atomic::daemon_service_server::DaemonServiceServer; + use libatomic::atomic::maintenance_service_server::MaintenanceServiceServer; + use libatomic::daemon::{ + services::DaemonImpl, services_maintenance::MaintenanceImpl, state::DaemonState, + }; + use std::sync::Arc; + let dir = tempfile::tempdir().unwrap(); + let root = dir.path().join("repo"); + drop(Repository::init(&root).unwrap()); + // Short socket path for macOS. The test server serves the real library; + // this does not require or modify the private Reactor checkout. + let socket_dir = tempfile::Builder::new() + .prefix("compact-") + .tempdir_in("/tmp") + .unwrap(); + let socket = socket_dir.path().join("rpc.sock"); + let listener = tokio::net::UnixListener::bind(&socket).unwrap(); + let state = Arc::new(DaemonState::new()); + let (stop, stopped) = tokio::sync::oneshot::channel(); + let server = tokio::spawn( + tonic::transport::Server::builder() + .add_service(DaemonServiceServer::new(DaemonImpl { + state: state.clone(), + })) + .add_service(MaintenanceServiceServer::new(MaintenanceImpl { state })) + .serve_with_incoming_shutdown( + tokio_stream::wrappers::UnixListenerStream::new(listener), + async { + let _ = stopped.await; + }, + ), + ); + let mut cli = command(&root); + cli.env("ATOMIC_SERVICE", "reactor") + .env("ATOMIC_DAEMON_SOCKET", &socket); + let output = tokio::task::spawn_blocking(move || cli.output().unwrap()) + .await + .unwrap(); + let result = report(output); + assert_eq!( + result["database"], + root.join(".atomic/atomic.redb") + .canonicalize() + .unwrap() + .to_str() + .unwrap() + ); + stop.send(()).unwrap(); + server.await.unwrap().unwrap(); + let mut cli = command(&root); + cli.env("ATOMIC_SERVICE", "reactor") + .env("ATOMIC_DAEMON_SOCKET", &socket); + let output = tokio::task::spawn_blocking(move || cli.output().unwrap()) + .await + .unwrap(); + assert!( + !output.status.success(), + "reactor mode must never fall back to local compaction" + ); +} diff --git a/libatomic/Cargo.toml b/libatomic/Cargo.toml index c553692c..02d987dc 100644 --- a/libatomic/Cargo.toml +++ b/libatomic/Cargo.toml @@ -37,6 +37,12 @@ blake3 = { workspace = true } data-encoding = { workspace = true } similar = { workspace = true } ulid = { workspace = true } +redb = { workspace = true } +anyhow = { workspace = true } +uuid = { workspace = true } + +[dev-dependencies] +tempfile = { workspace = true } [build-dependencies] tonic-build = { workspace = true } diff --git a/libatomic/proto/CONTRACT.md b/libatomic/proto/CONTRACT.md index 6c04e4ab..c0394801 100644 --- a/libatomic/proto/CONTRACT.md +++ b/libatomic/proto/CONTRACT.md @@ -70,6 +70,32 @@ that would incorrectly exclude the serving local atomicd. stream/message/queue limits apply per principal; slow streaming cannot keep a writer open. Negotiate limits before allocating opaque payloads. +## Repeatable physical maintenance + +`MaintenanceService.CompactDatabase` is an explicit exception to logical-mutation +replay/publication semantics. It changes the physical page layout of an existing +database without changing its logical records. redb performs its own internal +transactions; there is no Atomic replay receipt written into the database. +`RequestMeta.request_id` correlates an attempt and is echoed, with `replayed=false`. +A repeated request may run again and report different byte counts, including zero +reclaimed bytes. The CLI does not automatically retry a lost response. + +Only local administrators with `maintenance.admin` may invoke it. The serving +transport must enforce caller identity and the descriptor allowlist before +dispatch; a sandbox working-directory path does not confer sandbox-token +authority. Compaction resolves that path to the canonical database without +opening a Repository or migrating it. + +Acquiring the repository gate and exclusive redb file lock is bounded by +`ATOMIC_DB_LOCK_WAIT_MS`; a busy database returns UNAVAILABLE. Read paths which +do not use the repository gate remain protected by redb's file lock and may +fail/wait during maintenance. Cancellation while queued for the gate starts no +blocking work. Once blocking maintenance starts, that task retains the gate and +finishes independently of the RPC caller, including its bounded file-lock wait. +Persistent savepoints are preserved and reported as a precondition failure. +Successful byte counts are sampled before compaction and after closing redb; +other processes should remain idle until the command returns. + ## Repository and workspace routing RepositoryRef.authority selects an explicit configured authority; it is not an diff --git a/libatomic/proto/COVERAGE.md b/libatomic/proto/COVERAGE.md index 5d40e213..c7153d68 100644 --- a/libatomic/proto/COVERAGE.md +++ b/libatomic/proto/COVERAGE.md @@ -40,7 +40,7 @@ completion engine are load-bearing surfaces. | SyncService | PushChanges, PullChanges, ListRemotes, ManageRemotes, BindRemoteProject | all LOCAL | | GitInteropService | ImportFromGit, PreviewGitImport, PushToGit | LOCAL | | TriageService | ListTriageCandidates, GenerateTriageReview | BOTH | -| MaintenanceService | Repair, CheckRepository, ReindexWorkspace | LOCAL, BOTH, LOCAL | +| MaintenanceService | Repair, CheckRepository, ReindexWorkspace, CompactDatabase | LOCAL, BOTH, LOCAL, LOCAL | | ReactorService | Publish, PublishStream, Subscribe, Replay, CheckpointConsumer | all BOTH (contract draft) | ## 2. Legacy agent journal protocol (13 ops, BEING DISMANTLED) → ProvenanceService @@ -151,6 +151,7 @@ daemon, never the write queue. | `split [--switch]` | direct-RW | CreateView (explicit base); --switch composes SwitchView | | `stash push/pop/apply/list/show/drop/clear` | direct-RW | CreateStash / ApplyStash / ListStashes / DropStash; pop composes apply then drop | | `doctor repair-dependency-index/materialize-crdt/check` | direct-RW (check: read semantics) | Repair + CheckRepository (maintenance.read, BOTH) | +| `compact [--repository PATH] [--json]` | shared service (local or Reactor) | CompactDatabase (maintenance.admin, LOCAL); repeatable physical maintenance | | `session show/fork/rebuild` | direct-RW | ProvenanceService.GetSession / ForkSession / RebuildSessionIndex | | `sandbox create/stage/seal` | direct-RW | SandboxService.CreateSandboxTree / StageSandboxImage / SealSandboxImage (SHAPE-CONFIRM) | | `tag create/delete/list/show` | direct-RW | TagService (create/delete/list/get) | diff --git a/libatomic/proto/atomic/services/maintenance/maintenance_messages.proto b/libatomic/proto/atomic/services/maintenance/maintenance_messages.proto index 7270676a..b67c9f4c 100644 --- a/libatomic/proto/atomic/services/maintenance/maintenance_messages.proto +++ b/libatomic/proto/atomic/services/maintenance/maintenance_messages.proto @@ -12,6 +12,21 @@ package atomic; import "atomic/common/common.proto"; +message CompactDatabaseRequest { + RepositoryRef repository = 1; + // Correlates this physical maintenance attempt. Unlike logical mutations, + // compaction is repeatable, without a durable replay receipt (CONTRACT.md). + RequestMeta meta = 2; +} +message CompactDatabaseResponse { + // Canonical physical path, available only to local administrative callers. + string database = 1; + uint64 before_bytes = 2; + uint64 after_bytes = 3; + uint64 reclaimed_bytes = 4; + ResponseMeta meta = 100; +} + enum RepairAction { diff --git a/libatomic/proto/atomic/services/maintenance/maintenance_service.proto b/libatomic/proto/atomic/services/maintenance/maintenance_service.proto index 4600e3fe..28d5c2ad 100644 --- a/libatomic/proto/atomic/services/maintenance/maintenance_service.proto +++ b/libatomic/proto/atomic/services/maintenance/maintenance_service.proto @@ -10,6 +10,14 @@ import "atomic/common/options.proto"; import "atomic/services/maintenance/maintenance_messages.proto"; service MaintenanceService { + // Explicit physical maintenance of an existing atomic.redb. Busy handles + // cause bounded waiting, then UNAVAILABLE. Not logical garbage collection. + rpc CompactDatabase(CompactDatabaseRequest) returns (CompactDatabaseResponse) { + option (atomic.execution_scope) = EXECUTION_SCOPE_LOCAL; + option (atomic.required_capability) = "maintenance.admin"; + option (atomic.allowed_caller) = CALLER_CLASS_LOCAL; + option (atomic.effect) = RPC_EFFECT_REPOSITORY_WRITE; + } // `atomic doctor repair-dependency-index | materialize-crdt` — mutating // repairs only; the read-only consistency check lives in CheckRepository. rpc Repair(RepairRequest) returns (RepairResponse) { diff --git a/libatomic/src/daemon/services.rs b/libatomic/src/daemon/services.rs index 55d7fa02..ff315beb 100644 --- a/libatomic/src/daemon/services.rs +++ b/libatomic/src/daemon/services.rs @@ -78,6 +78,7 @@ const METHOD_MATRIX: &[(&str, &str)] = &[ ("ProvenanceService", "GetSession"), ("ProvenanceService", "ListSessions"), ("MaintenanceService", "Repair"), + ("MaintenanceService", "CompactDatabase"), ("MaintenanceService", "CheckRepository"), ("ViewService", "ListViews"), ("ViewService", "CreateView"), @@ -109,13 +110,26 @@ impl daemon_service_server::DaemonService for DaemonImpl { ) -> Result, Status> { let methods = METHOD_MATRIX .iter() - .map(|(service, method)| MethodDescriptorInfo { - service: service.to_string(), - method: method.to_string(), - scope: ExecutionScope::Both as i32, - required_capabilities: Vec::new(), - allowed_callers: Vec::new(), - effects: Vec::new(), + .map(|(service, method)| { + if *service == "MaintenanceService" && *method == "CompactDatabase" { + MethodDescriptorInfo { + service: service.to_string(), + method: method.to_string(), + scope: ExecutionScope::Local as i32, + required_capabilities: vec!["maintenance.admin".to_string()], + allowed_callers: vec![CallerClass::Local as i32], + effects: vec![RpcEffect::RepositoryWrite as i32], + } + } else { + MethodDescriptorInfo { + service: service.to_string(), + method: method.to_string(), + scope: ExecutionScope::Both as i32, + required_capabilities: Vec::new(), + allowed_callers: Vec::new(), + effects: Vec::new(), + } + } }) .collect(); Ok(Response::new(GetCapabilitiesResponse { diff --git a/libatomic/src/daemon/services_maintenance.rs b/libatomic/src/daemon/services_maintenance.rs index 49569f37..4b47e4a9 100644 --- a/libatomic/src/daemon/services_maintenance.rs +++ b/libatomic/src/daemon/services_maintenance.rs @@ -3,11 +3,11 @@ //! `atomic doctor` folds into Repair (the mutating actions) and //! CheckRepository (the read-only consistency check). The read/write split //! the state gate defers lands here for the read side: CheckRepository -//! opens the repository READ-ONLY (redb read-only handles coexist with -//! the writer's exclusive lock) and never acquires the per-repo gate — -//! a consistency check never queues behind writers, exactly as the -//! contract's "EXECUTES AS A READ TRANSACTION" demands. Repair is a -//! repository write: gated like every other mutating handler. +//! opens the repository READ-ONLY without acquiring the per-repo gate. +//! Separate read-only handles share redb's file lock with other readers; +//! an exclusive writable handle still prevents these opens. Repair uses the +//! mutation gate. Compaction additionally requires exclusive access to the +//! physical database, including exclusion of independent reader handles. use std::sync::Arc; @@ -19,6 +19,8 @@ use tonic::{Request, Response, Status}; use super::services::default_ref; use super::state::{domain_status, repository_error, DaemonState}; +mod compact; + pub struct MaintenanceImpl { pub state: Arc, } @@ -41,6 +43,19 @@ fn repair_action(action: i32) -> Result { #[tonic::async_trait] impl MaintenanceService for MaintenanceImpl { + async fn compact_database( + &self, + request: Request, + ) -> Result, Status> { + compact::compact( + &self.state, + request, + atomic_agent::turn::orchestrator::wait_budget::database_wait(), + ) + .await + .map(Response::new) + } + async fn repair( &self, request: Request, @@ -139,9 +154,8 @@ impl MaintenanceService for MaintenanceImpl { .state .resolve(request.repository.as_ref().unwrap_or(&default_ref()))?; self.state.log_rpc("CheckRepository", Some(&handle)); - // No gate, no writer queue: read-only opens coexist with the - // writer's exclusive handle, so the check runs concurrent with - // in-flight writes ("EXECUTES AS A READ TRANSACTION"). + // No mutation gate. redb's shared file lock still refuses this + // open while a writable handle (including compaction) is live. let root = handle.root.clone(); let result = tokio::task::spawn_blocking(move || { let repo = atomic_repository::Repository::open_readonly(&root).map_err(|error| { diff --git a/libatomic/src/daemon/services_maintenance/compact.rs b/libatomic/src/daemon/services_maintenance/compact.rs new file mode 100644 index 00000000..0b8e35ec --- /dev/null +++ b/libatomic/src/daemon/services_maintenance/compact.rs @@ -0,0 +1,178 @@ +//! Physical maintenance. Only this private module opens redb for compaction. + +use std::path::Path; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use atomic_repository::Repository; +use prost::Message; +use redb::ReadableDatabase; +use tonic::{Request, Status}; + +use crate::atomic::{ + CompactDatabaseRequest, CompactDatabaseResponse, ErrorCode, ErrorInfo, ResponseMeta, +}; +use crate::daemon::state::DaemonState; + +pub(super) async fn compact( + state: &Arc, + request: Request, + wait: Duration, +) -> Result { + // Serving transports must enforce LOCAL + maintenance.admin before calling + // this handler. Restricted credentials must never fall through to the local + // administrator path; invoking from a sandbox directory is a separate case. + if request + .metadata() + .contains_key("x-atomic-sandbox-token-bin") + || request.metadata().contains_key("authorization") + { + return Err(Status::permission_denied( + "database compaction requires a local administrator", + )); + } + let request = request.into_inner(); + let request_id = request + .meta + .as_ref() + .map(|meta| meta.request_id.as_str()) + .unwrap_or(""); + if uuid::Uuid::parse_str(request_id).is_err() { + return Err(Status::invalid_argument( + "compaction requires a UUID request_id", + )); + } + compact_inner(state, &request, wait) + .await + .map_err(|status| with_request_id(status, request_id)) +} + +async fn compact_inner( + state: &Arc, + request: &CompactDatabaseRequest, + wait: Duration, +) -> Result { + let started = Instant::now(); + let reference = request + .repository + .as_ref() + .ok_or_else(|| Status::invalid_argument("repository is required"))?; + let resolved = state.resolve(reference)?; + // Resolve sandbox pointers without opening, migrating or recovering a repo. + let dot_dir = Repository::canonical_dot_dir(&resolved.root) + .map_err(|error| Status::failed_precondition(error.to_string()))? + .canonicalize() + .map_err(|error| Status::failed_precondition(error.to_string()))?; + let path = dot_dir + .join(atomic_repository::DATABASE_FILE) + .canonicalize() + .map_err(|error| { + Status::failed_precondition(format!("existing atomic.redb required: {error}")) + })?; + // The sandbox and canonical repository must share maintenance exclusion. + let handle = state.register( + dot_dir + .parent() + .ok_or_else(|| Status::failed_precondition("canonical .atomic has no parent"))? + .to_path_buf(), + ); + let guard = tokio::time::timeout( + wait.saturating_sub(started.elapsed()), + handle.exclusive_owned(), + ) + .await + .map_err(|_| { + Status::unavailable("database compaction timed out waiting for active requests") + })?; + state.log_rpc("CompactDatabase", Some(&handle)); + let report = tokio::task::spawn_blocking(move || { + // The blocking task, not the RPC future, owns this until redb closes. + let _guard = guard; + compact_waiting(&path, started, wait) + }) + .await + .map_err(|error| Status::internal(format!("compaction task failed: {error}")))??; + Ok(CompactDatabaseResponse { + meta: Some(ResponseMeta { + request_id: request + .meta + .as_ref() + .expect("validated metadata") + .request_id + .clone(), + replayed: false, + snapshot: None, + }), + ..report + }) +} + +fn compact_waiting( + path: &Path, + started: Instant, + wait: Duration, +) -> Result { + loop { + match compact_file(path) { + Ok(report) => return Ok(report), + Err(error) + if matches!( + error.downcast_ref::(), + Some(redb::DatabaseError::DatabaseAlreadyOpen) + ) => + { + let remaining = wait.saturating_sub(started.elapsed()); + if remaining.is_zero() { + return Err(Status::unavailable(format!( + "database compaction timed out waiting for {}; close other database readers/writers and retry", + path.display() + ))); + } + std::thread::sleep(remaining.min(Duration::from_millis(25))); + } + Err(error) => { + return Err(Status::failed_precondition(format!( + "cannot compact {}: {error:#}", + path.display() + ))) + } + } + } +} + +fn compact_file(path: &Path) -> anyhow::Result { + let mut database = redb::Builder::new().open(path)?; + atomic_core::pristine::schema::check_schema_version(&database.begin_read()?)?; + let before_bytes = std::fs::metadata(path)?.len(); + database.compact()?; + // Closing persists allocator metadata; include that in the final size. + drop(database); + let after_bytes = std::fs::metadata(path)?.len(); + Ok(CompactDatabaseResponse { + database: path.to_string_lossy().into_owned(), + before_bytes, + after_bytes, + reclaimed_bytes: before_bytes.saturating_sub(after_bytes), + meta: None, + }) +} + +fn with_request_id(status: Status, request_id: &str) -> Status { + let code = match status.code() { + tonic::Code::Unavailable => ErrorCode::ResourceExhausted, + tonic::Code::InvalidArgument => ErrorCode::InvalidArgument, + tonic::Code::NotFound => ErrorCode::RepositoryNotFound, + tonic::Code::FailedPrecondition => ErrorCode::PreconditionFailed, + _ => ErrorCode::Internal, + }; + let info = ErrorInfo { + code: code as i32, + message: status.message().to_string(), + request_id: Some(request_id.to_string()), + ..Default::default() + }; + Status::with_details(status.code(), status.message(), info.encode_to_vec().into()) +} + +#[cfg(test)] +mod tests; diff --git a/libatomic/src/daemon/services_maintenance/compact/tests.rs b/libatomic/src/daemon/services_maintenance/compact/tests.rs new file mode 100644 index 00000000..ade5cf95 --- /dev/null +++ b/libatomic/src/daemon/services_maintenance/compact/tests.rs @@ -0,0 +1,323 @@ +use super::*; +use redb::{Database, ReadableTableMetadata, TableDefinition}; + +const DATA: TableDefinition = TableDefinition::new("compaction_data"); + +fn request(handle: &crate::daemon::state::RepoHandle) -> Request { + Request::new(CompactDatabaseRequest { + repository: Some(handle.repository_ref()), + meta: Some(crate::atomic::RequestMeta { + request_id: uuid::Uuid::new_v4().to_string(), + observed_at: None, + }), + }) +} + +#[tokio::test] +async fn busy_reader_and_writer_time_out_then_retry_succeeds() { + let dir = tempfile::tempdir().unwrap(); + let repo = Repository::init(dir.path()).unwrap(); + let state = Arc::new(DaemonState::new()); + let handle = state.register(dir.path().to_path_buf()); + let call = request(&handle); + let id = call.get_ref().meta.as_ref().unwrap().request_id.clone(); + let error = compact(&state, call, Duration::from_millis(20)) + .await + .unwrap_err(); + assert_eq!(error.code(), tonic::Code::Unavailable); + assert_eq!( + ErrorInfo::decode(error.details()).unwrap().request_id, + Some(id) + ); + drop(repo); + let reader = Repository::open_readonly(dir.path()).unwrap(); + assert_eq!( + compact(&state, request(&handle), Duration::from_millis(20)) + .await + .unwrap_err() + .code(), + tonic::Code::Unavailable + ); + drop(reader); + let call = request(&handle); + let id = call.get_ref().meta.as_ref().unwrap().request_id.clone(); + let report = compact(&state, call, Duration::from_secs(1)).await.unwrap(); + assert_eq!(report.meta.unwrap().request_id, id); + assert!(Repository::open_existing(dir.path()).is_ok()); +} + +#[tokio::test] +async fn refuses_invalid_requests_before_opening_a_database() { + let dir = tempfile::tempdir().unwrap(); + let state = Arc::new(DaemonState::new()); + let handle = state.register(dir.path().to_path_buf()); + let mut call = request(&handle); + call.get_mut().meta = None; + assert_eq!( + compact(&state, call, Duration::ZERO) + .await + .unwrap_err() + .code(), + tonic::Code::InvalidArgument + ); + for header in ["authorization", "x-atomic-sandbox-token-bin"] { + let mut call = request(&handle); + if header.ends_with("-bin") { + call.metadata_mut() + .insert_bin(header, tonic::metadata::MetadataValue::from_bytes(b"token")); + } else { + call.metadata_mut() + .insert(header, "Bearer token".parse().unwrap()); + } + assert_eq!( + compact(&state, call, Duration::ZERO) + .await + .unwrap_err() + .code(), + tonic::Code::PermissionDenied + ); + } + assert!(compact(&state, request(&handle), Duration::ZERO) + .await + .is_err()); + assert!(!dir.path().join(".atomic").exists()); + let mut unknown = request(&handle); + unknown.get_mut().repository.as_mut().unwrap().repository_id = vec![7; 32]; + assert_eq!( + compact(&state, unknown, Duration::ZERO) + .await + .unwrap_err() + .code(), + tonic::Code::NotFound + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn cancellation_keeps_gate_until_blocking_compaction_finishes() { + use std::future::Future; + use std::task::Poll; + let dir = tempfile::tempdir().unwrap(); + let held = Repository::init(dir.path()).unwrap(); + let state = Arc::new(DaemonState::new()); + let handle = state.register(dir.path().to_path_buf()); + let held_gate = handle.exclusive().await; + let mut call = Box::pin(compact(&state, request(&handle), Duration::from_secs(3))); + // Poll to queue behind the held gate, without depending on scheduling sleeps. + assert!(std::future::poll_fn(|cx| Poll::Ready(call.as_mut().poll(cx).is_pending())).await); + drop(held_gate); + // This poll obtains the gate and starts blocking work; the independent + // repository handle ensures compaction cannot have finished yet. + assert!(std::future::poll_fn(|cx| Poll::Ready(call.as_mut().poll(cx).is_pending())).await); + drop(call); // RPC cancellation/disconnection drops its future. + assert!( + tokio::time::timeout(Duration::from_millis(20), handle.exclusive()) + .await + .is_err() + ); + drop(held); + let _gate = tokio::time::timeout(Duration::from_secs(2), handle.exclusive()) + .await + .unwrap(); + assert!(Repository::open_existing(dir.path()).is_ok()); +} + +#[tokio::test] +async fn concurrent_compactions_finish_and_release_the_gate() { + let dir = tempfile::tempdir().unwrap(); + drop(Repository::init(dir.path()).unwrap()); + let state = Arc::new(DaemonState::new()); + let handle = state.register(dir.path().to_path_buf()); + let (first, second) = tokio::join!( + compact(&state, request(&handle), Duration::from_secs(2)), + compact(&state, request(&handle), Duration::from_secs(2)), + ); + assert_eq!(first.unwrap().database, second.unwrap().database); + let _gate = tokio::time::timeout(Duration::from_secs(1), handle.exclusive()) + .await + .unwrap(); + assert!(Repository::open_existing(dir.path()).is_ok()); +} + +#[tokio::test] +async fn sandbox_uses_canonical_database_and_gate() { + let dir = tempfile::tempdir().unwrap(); + let canonical = dir.path().join("repo"); + let sandbox = dir.path().join("sandbox"); + let repo = Repository::init(&canonical).unwrap(); + repo.provision_sandbox(&sandbox, repo.current_view()) + .unwrap(); + drop(repo); + let state = Arc::new(DaemonState::new()); + let main_handle = state.register(canonical.clone()); + let sandbox_handle = state.register(sandbox.clone()); + let gate = main_handle.exclusive().await; + assert_eq!( + compact(&state, request(&sandbox_handle), Duration::from_millis(20)) + .await + .unwrap_err() + .code(), + tonic::Code::Unavailable + ); + drop(gate); + let report = compact(&state, request(&sandbox_handle), Duration::from_secs(1)) + .await + .unwrap(); + assert_eq!( + Path::new(&report.database), + canonical + .join(".atomic/atomic.redb") + .canonicalize() + .unwrap() + ); + assert!(!sandbox.join(".atomic/atomic.redb").exists()); +} + +#[tokio::test] +async fn capability_advertises_local_administrative_write() { + use crate::atomic::daemon_service_server::DaemonService; + use crate::atomic::{CallerClass, ExecutionScope, GetCapabilitiesRequest, RpcEffect}; + let service = crate::daemon::services::DaemonImpl { + state: Arc::new(DaemonState::new()), + }; + let capabilities = service + .get_capabilities(Request::new(GetCapabilitiesRequest {})) + .await + .unwrap() + .into_inner(); + let method = capabilities + .methods + .iter() + .find(|m| m.method == "CompactDatabase") + .unwrap(); + assert_eq!(method.scope, ExecutionScope::Local as i32); + assert_eq!(method.required_capabilities, ["maintenance.admin"]); + assert_eq!(method.allowed_callers, [CallerClass::Local as i32]); + assert_eq!(method.effects, [RpcEffect::RepositoryWrite as i32]); +} + +#[test] +fn compaction_reclaims_space_and_preserves_surviving_values() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("atomic.redb"); + let db = Database::create(&path).unwrap(); + let value = vec![0x5a; 4096]; + let write = db.begin_write().unwrap(); + { + let mut table = write.open_table(DATA).unwrap(); + for key in 0..1024 { + table.insert(key, value.as_slice()).unwrap(); + } + } + write.commit().unwrap(); + // Keep the old pages pinned through deletion so the fixture has + // reclaimable space without depending on the allocator's layout. + let reader = db.begin_read().unwrap(); + let write = db.begin_write().unwrap(); + { + let mut table = write.open_table(DATA).unwrap(); + for key in 1..1024 { + table.remove(key).unwrap(); + } + } + write.commit().unwrap(); + drop(reader); + drop(db); + + let before = std::fs::metadata(&path).unwrap().len(); + let report = compact_file(&path).unwrap(); + assert_eq!(report.before_bytes, before); + assert!(report.after_bytes < report.before_bytes, "{report:?}"); + assert_eq!(report.after_bytes, std::fs::metadata(&path).unwrap().len()); + assert_eq!(report.reclaimed_bytes, before - report.after_bytes); + let db = redb::Builder::new().open_read_only(&path).unwrap(); + let read = db.begin_read().unwrap(); + assert_eq!(read.list_tables().unwrap().count(), 1); + let table = read.open_table(DATA).unwrap(); + assert_eq!(table.len().unwrap(), 1); + assert_eq!(table.get(0).unwrap().unwrap().value(), value); + drop(table); + drop(read); + drop(db); + // A repeated maintenance request remains valid. + compact_file(&path).unwrap(); +} + +#[test] +fn compaction_does_not_create_a_missing_database() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("atomic.redb"); + assert!(compact_file(&path).is_err()); + assert!(!path.exists()); +} + +#[test] +fn compaction_respects_existing_writer_and_reader_handles() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("atomic.redb"); + let db = Database::create(&path).unwrap(); + assert!(matches!( + compact_file(&path) + .unwrap_err() + .downcast_ref::(), + Some(redb::DatabaseError::DatabaseAlreadyOpen) + )); + drop(db); + let db = redb::Builder::new().open_read_only(&path).unwrap(); + assert!(matches!( + compact_file(&path) + .unwrap_err() + .downcast_ref::(), + Some(redb::DatabaseError::DatabaseAlreadyOpen) + )); + drop(db); + compact_file(&path).unwrap(); +} + +#[test] +fn compaction_rejects_a_future_repository_schema() { + use atomic_core::pristine::{schema::SCHEMA_VERSION, tables::ATOMIC_META}; + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("atomic.redb"); + let db = Database::create(&path).unwrap(); + let write = db.begin_write().unwrap(); + { + let mut meta = write.open_table(ATOMIC_META).unwrap(); + meta.insert( + atomic_core::pristine::schema::SCHEMA_VERSION_KEY, + (SCHEMA_VERSION + 1).to_le_bytes().as_slice(), + ) + .unwrap(); + } + write.commit().unwrap(); + drop(db); + let error = compact_file(&path).unwrap_err().to_string(); + assert!(error.contains("schema"), "{error}"); +} + +#[test] +fn compaction_preserves_persistent_savepoints_and_reports_the_blocker() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("atomic.redb"); + let db = Database::create(&path).unwrap(); + let write = db.begin_write().unwrap(); + let savepoint = write.persistent_savepoint().unwrap(); + write.commit().unwrap(); + drop(db); + + assert!(matches!( + compact_file(&path) + .unwrap_err() + .downcast::() + .unwrap(), + redb::CompactionError::PersistentSavepointExists + )); + let db = redb::Builder::new().open(&path).unwrap(); + let write = db.begin_write().unwrap(); + assert_eq!( + write + .list_persistent_savepoints() + .unwrap() + .collect::>(), + vec![savepoint] + ); +} diff --git a/libatomic/src/daemon/state.rs b/libatomic/src/daemon/state.rs index bbde5255..c7e8e3cf 100644 --- a/libatomic/src/daemon/state.rs +++ b/libatomic/src/daemon/state.rs @@ -36,14 +36,14 @@ pub const READ_OPEN_WAIT: std::time::Duration = std::time::Duration::from_secs(1 /// serialization for this slice — the read/write split is future work.) pub struct RepoHandle { pub root: PathBuf, - gate: tokio::sync::Mutex<()>, + gate: std::sync::Arc>, } impl RepoHandle { pub fn new(root: PathBuf) -> Self { Self { root, - gate: tokio::sync::Mutex::new(()), + gate: std::sync::Arc::new(tokio::sync::Mutex::new(())), } } @@ -53,6 +53,12 @@ impl RepoHandle { self.gate.lock().await } + /// Move maintenance exclusion into blocking work so cancelling its caller + /// cannot release the gate while that work is still using the database. + pub(super) async fn exclusive_owned(&self) -> tokio::sync::OwnedMutexGuard<()> { + self.gate.clone().lock_owned().await + } + /// The persistent repository ID for this daemon: a blake3 digest of the /// canonical root. Slice limitation: RFC D3 wants IDs assigned at init /// and recorded in `.atomic` so they survive daemon restarts and follow