Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

25 changes: 25 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,31 @@ The **lock is the authority; the socket or named pipe is only transport**. The l

Repositories created before `atomic.redb` kept graph state in `.atomic/pristine.redb` and the journal in `.atomic/changes.redb`. The first open by a newer `atomic` merges both into `atomic.redb`, verifies every table against its source, and moves the old files to `.atomic/legacy/<timestamp>/`. If an owner started by an older `atomic` still holds `changes.redb`, run `atomic agent database-owner shutdown` in that repository first.

### Reclaiming database space

Run `atomic compact` to reclaim unused space in an existing `atomic.redb`:

```sh
atomic compact
atomic compact --repository /path/to/repo --json
```

The command starts or reconnects to the repository's database owner, waits for
in-flight database requests to finish, and compacts with exclusive access.
It reports `database`, `before_bytes`, `after_bytes`, and `reclaimed_bytes`.
Running it inside a sandbox compacts that sandbox's canonical repository.
Changes, provenance, views, and working-copy files are preserved; this does not
delete history or `.change` files. The amount reclaimed depends on unused space,
and may be zero. Compaction runs only when explicitly requested.

Use a quiet period for large repositories: database operations wait while
compaction runs and can exhaust their normal lock-wait budgets. The owner's
`ATOMIC_DB_LOCK_WAIT_MS` setting bounds acquisition waits, not the duration of
an already-started compaction. Owners started by an older executable must be
shut down once before retrying. For a legacy repository, first open it with the
current Atomic to complete its database migration; `compact` does not migrate
or create a database.

Within one repository, redb write transactions are serialized by design, while reads and independent repositories can proceed concurrently. Startup, reconnect, and retry loops are bounded so a broken owner fails instead of hanging hooks indefinitely. Checkpoint attempts, event cutoffs, and fencing generations let interrupted turns resume without rewriting completed session turns.

Current operational considerations:
Expand Down
3 changes: 3 additions & 0 deletions atomic-cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,9 @@ console = { workspace = true }
# Async runtime
tokio = { workspace = true }

# Exclusive maintenance inside the database owner
redb = { workspace = true }

# Error handling
anyhow = { workspace = true }
thiserror = { workspace = true }
Expand Down
2 changes: 2 additions & 0 deletions atomic-cli/src/commands/agent/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@ mod lifecycle;
mod owner;
mod status;

pub(crate) use owner::compact_repository;

use clap::{Args, Subcommand};

use crate::commands::Command;
Expand Down
173 changes: 161 additions & 12 deletions atomic-cli/src/commands/agent/owner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ use std::fs::{File, OpenOptions};
use std::ops::Deref;
use std::path::{Path, PathBuf};
use std::process::{Command as ProcessCommand, Stdio};
use std::sync::{Arc, Mutex, TryLockError, Weak};
use std::sync::{Arc, Mutex, MutexGuard, TryLockError, Weak};
use std::time::{Duration, Instant};

use anyhow::{anyhow, Context};
Expand All @@ -28,7 +28,7 @@ use atomic_repository::redb_change_store::{
ProvenanceCheckpointSource, ProvenanceId, ProvenanceTurnState, RedbChangeStore, StopCause,
StopState, StoredProvenanceTurn, MAX_FROZEN_PAGE_FRAGMENTS,
};
use atomic_repository::{ensure_database, Repository, RepositoryError};
use atomic_repository::{ensure_database, Repository, RepositoryError, DATABASE_FILE};
use clap::{Args, Subcommand};
use fs2::FileExt;
use serde::{Deserialize, Serialize};
Expand All @@ -39,6 +39,9 @@ use uuid::Uuid;
use crate::commands::Command;
use crate::error::CliResult;

mod maintenance;
use maintenance::{compact_database, is_database_busy, DatabaseCompaction};

const PROTOCOL_VERSION: u16 = 1;
// Local IPC only; a single oversized envelope (e.g. a huge recovered reasoning
// block) must still fit one frame even after client-side chunking.
Expand Down Expand Up @@ -118,6 +121,7 @@ struct RequestFrame {
#[derive(Clone, Debug, Serialize, Deserialize)]
enum OwnerRequest {
Ping,
Compact,
ReserveProvenanceTurn {
session_id: String,
turn_number: u32,
Expand Down Expand Up @@ -203,6 +207,11 @@ pub(crate) enum OwnerResponse {
// Additive v1 capability: old clients ignore it, old owners omit it.
#[serde(default)]
frozen_envelope_paging: bool,
#[serde(default)]
database_compaction: bool,
},
Compacted {
report: DatabaseCompaction,
},
ProvenanceTurn {
turn: StoredProvenanceTurn,
Expand Down Expand Up @@ -838,6 +847,29 @@ pub(crate) fn start_or_reconnect(repository: &Path) -> anyhow::Result<OwnerRespo
))
}

/// Run explicit maintenance through the owner, never an ordinary CLI opener.
pub(crate) fn compact_repository(repository: &Path) -> anyhow::Result<DatabaseCompaction> {
require_compaction(start_or_reconnect(repository)?)?;
match request(repository, OwnerRequest::Compact)? {
OwnerResponse::Compacted { report } => Ok(report),
other => Err(unexpected_response("compact", other)),
}
}

fn require_compaction(response: OwnerResponse) -> anyhow::Result<()> {
match response {
OwnerResponse::Pong {
database_compaction: true,
..
} => Ok(()),
OwnerResponse::Pong { .. } => Err(anyhow!(
"the running database owner does not support compaction; run \
`atomic agent database-owner shutdown --repository <path>` and retry"
)),
other => Err(unexpected_response("ping", other)),
}
}

fn request_with_reconnect(
repository: &Path,
request_body: OwnerRequest,
Expand Down Expand Up @@ -932,9 +964,23 @@ impl StoreLease {

fn acquire(&self, started: Instant) -> anyhow::Result<Arc<RedbChangeStore>> {
let timeout = atomic_agent::turn::orchestrator::wait_budget::database_wait();
let mut current = loop {
let mut current = self.lock_current(started, timeout)?;
if let Some(store) = current.upgrade() {
return Ok(store);
}
let store = Arc::new(self.open_waiting(started, timeout)?);
*current = Arc::downgrade(&store);
Ok(store)
}

fn lock_current(
&self,
started: Instant,
timeout: Duration,
) -> anyhow::Result<MutexGuard<'_, Weak<RedbChangeStore>>> {
loop {
match self.current.try_lock() {
Ok(guard) => break guard,
Ok(guard) => return Ok(guard),
Err(TryLockError::Poisoned(_)) => {
return Err(anyhow!("database-owner store lease is poisoned"));
}
Expand All @@ -949,13 +995,45 @@ impl StoreLease {
));
}
}
};
if let Some(store) = current.upgrade() {
return Ok(store);
}
let store = Arc::new(self.open_waiting(started, timeout)?);
*current = Arc::downgrade(&store);
Ok(store)
}

fn compact(&self, started: Instant) -> anyhow::Result<DatabaseCompaction> {
let timeout = atomic_agent::turn::orchestrator::wait_budget::database_wait();
// Keep this gate locked through compaction. Existing requests can
// finish and release their Arcs, but no new store lease may open.
let current = self.lock_current(started, timeout)?;
while current.strong_count() != 0 {
if started.elapsed() >= timeout {
return Err(anyhow!(
"database compaction waited {timeout:?} for active requests"
));
}
std::thread::sleep(DATABASE_RETRY_DELAY.min(timeout.saturating_sub(started.elapsed())));
}
// The owner may have started with `--repository .` in a different
// working directory from this client (for example a sandbox).
let path = self.dot_dir.join(DATABASE_FILE);
let path = std::fs::canonicalize(&path)
.with_context(|| format!("failed to locate database {}", path.display()))?;
loop {
match compact_database(&path) {
Err(error) if is_database_busy(&error) && started.elapsed() < timeout => {
std::thread::sleep(
DATABASE_RETRY_DELAY.min(timeout.saturating_sub(started.elapsed())),
);
}
Err(error) if is_database_busy(&error) => {
return Err(anyhow!(
"database compaction waited {timeout:?} for {} to become available",
path.display()
));
}
result => {
return result.with_context(|| format!("failed to compact {}", path.display()))
}
}
}
}

fn open_waiting(&self, started: Instant, timeout: Duration) -> anyhow::Result<RedbChangeStore> {
Expand Down Expand Up @@ -1082,6 +1160,14 @@ where
};
let (response, shutdown) = match frame.request {
OwnerRequest::Ping => (pong(), false),
// Maintenance needs StoreLease's exclusive gate, not a shared handle.
OwnerRequest::Compact => (
OwnerResponse::Error {
code: "database-compaction".into(),
message: "compaction requires an exclusive owner lease".into(),
},
false,
),
OwnerRequest::ReserveProvenanceTurn {
session_id,
turn_number,
Expand Down Expand Up @@ -1426,6 +1512,7 @@ fn pong() -> OwnerResponse {
OwnerResponse::Pong {
pid: std::process::id(),
frozen_envelope_paging: true,
database_compaction: true,
}
}

Expand All @@ -1448,7 +1535,18 @@ where
// Opening the database may wait for another process; keep that off the
// runtime threads that accept connections and answer pings.
let (response, should_shutdown) = tokio::task::spawn_blocking(move || {
handle_request_with(request, || lease.acquire(started))
if request.version == PROTOCOL_VERSION && matches!(request.request, OwnerRequest::Compact) {
let response = match lease.compact(started) {
Ok(report) => OwnerResponse::Compacted { report },
Err(error) => OwnerResponse::Error {
code: "database-compaction".into(),
message: format!("{error:#}"),
},
};
respond(request.request_id, response, false)
} else {
handle_request_with(request, || lease.acquire(started))
}
})
.await
.context("database-owner request handler panicked")?;
Expand Down Expand Up @@ -1704,6 +1802,56 @@ fn configure_detached(command: &mut ProcessCommand) -> anyhow::Result<()> {
mod tests {
use super::*;

#[test]
fn legacy_owner_requires_an_explicit_restart_for_compaction() {
let old: OwnerResponse = serde_json::from_str(r#"{"Pong":{"pid":42}}"#).unwrap();
assert!(require_compaction(old)
.unwrap_err()
.to_string()
.contains("database-owner shutdown"));
assert!(require_compaction(pong()).is_ok());
}

#[test]
fn compaction_drains_active_leases_before_opening_the_database() {
let dir = tempfile::tempdir().unwrap();
let repo = Repository::init(dir.path()).unwrap();
let lease = Arc::new(StoreLease::new(repo.dot_dir().to_path_buf()));
drop(repo);
let held = lease.acquire(Instant::now()).unwrap();
let worker_lease = Arc::clone(&lease);
let (sender, receiver) = std::sync::mpsc::channel();
let worker = std::thread::spawn(move || {
sender.send(worker_lease.compact(Instant::now())).unwrap();
});
let deadline = Instant::now() + Duration::from_secs(5);
// Observe the maintenance gate deterministically before releasing
// the old lease, rather than relying on a sleep racing the worker.
loop {
if matches!(lease.current.try_lock(), Err(TryLockError::WouldBlock)) {
break;
}
assert!(
Instant::now() < deadline,
"compaction did not acquire its gate"
);
std::thread::yield_now();
}
assert!(receiver.try_recv().is_err());
drop(held);
let report = receiver
.recv_timeout(Duration::from_secs(10))
.unwrap()
.unwrap();
assert_eq!(
report.database,
std::fs::canonicalize(dir.path().join(".atomic").join(DATABASE_FILE)).unwrap()
);
worker.join().unwrap();
// Maintenance releases the gate and its redb handle for new requests.
assert!(lease.acquire(Instant::now()).is_ok());
}

#[cfg(windows)]
#[tokio::test]
async fn busy_owner_pipe_waits_for_a_new_listening_instance() {
Expand Down Expand Up @@ -1869,7 +2017,8 @@ mod tests {
.contains("database-owner shutdown"));
assert!(require_frozen_paging(OwnerResponse::Pong {
pid: 42,
frozen_envelope_paging: true
frozen_envelope_paging: true,
database_compaction: false,
})
.is_ok());
}
Expand Down
Loading