Skip to content
Merged
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
5 changes: 3 additions & 2 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -359,8 +359,9 @@ This simplifies the codebase while maintaining semantic clarity.

```
.atomic/
├── pristine/ # Graph database (redb)
│ └── data.mdb # Single database file
├── atomic.redb # Repository database (redb): graph, views,
│ # sessions, vault, provenance journal
├── legacy/<timestamp>/ # pristine.redb + changes.redb replaced by a merge
├── changes/ # Content-addressed change files
│ └── AB/CDEF... # Two-level directory structure
├── config.toml # Repository configuration
Expand Down
10 changes: 6 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -43,13 +43,13 @@ atomic agent enable

### Repository owner locks and connection recovery

Agent hooks are short-lived processes, while redb permits only one writable process to open a database. Atomic therefore starts one long-lived **database owner per canonical repository**. Hooks send committed provenance requests to that owner instead of opening the repository's redb change store directly.
Agent hooks are short-lived processes, while redb permits only one process at a time to open a database file. Atomic therefore starts one long-lived **database owner per canonical repository**. Hooks send committed provenance requests to that owner instead of opening the repository's provenance journal directly. The journal lives in the repository database, `.atomic/atomic.redb`, alongside the graph, so the owner opens that file only while it serves a request and leaves it free for other `atomic` commands in between.

Each repository has an independent ownership namespace:

| Resource | Location | Purpose |
|---|---|---|
| redb change store | `.atomic/changes.redb` | Mutable provenance journal and checkpoint state |
| repository database | `.atomic/atomic.redb` | Graph, views, sessions, vault, and the provenance journal with its checkpoint state |
| owner election lock | `.atomic/changes-owner.lock` | OS-backed authority deciding which process may own that database |
| Unix endpoint | `/tmp/atomic-owner-<repository-digest>.sock` | Local IPC transport on macOS and Linux |
| Windows endpoint | `\\.\pipe\atomic-owner-<repository-digest>` | Local IPC transport on Windows |
Expand All @@ -63,17 +63,19 @@ Agent A ─┐ Agent C ─┐
Agent B ─┴─> Project 1 owner Agent D ─┴─> Project 2 owner
│ │
v v
Project 1 changes.redb Project 2 changes.redb
Project 1 atomic.redb Project 2 atomic.redb
```

The **lock is the authority; the socket or named pipe is only transport**. The lock file may remain on disk after normal operation, but an unlocked file does not block a new owner. If an owner crashes, the OS releases its lock automatically. Recovery then proceeds as follows:

1. A hook cannot reach the old endpoint and starts or reconnects to an owner.
2. Owner candidates race for that repository's `changes-owner.lock`.
3. Only the lock winner may open `changes.redb`.
3. Only the lock winner serves provenance requests against `atomic.redb`.
4. On Unix, the winner removes any stale socket left by the dead owner and binds a fresh endpoint.
5. The hook retries with the same request/event ID, so a request committed before the crash is acknowledged once rather than applied twice.

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.

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
2 changes: 1 addition & 1 deletion atomic-agent/src/turn/orchestrator/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ mod provenance;
mod session_end;
mod session_start;
mod turn;
mod wait_budget;
pub mod wait_budget;

#[cfg(test)]
mod tests;
Expand Down
24 changes: 14 additions & 10 deletions atomic-agent/src/turn/orchestrator/provenance.rs
Original file line number Diff line number Diff line change
Expand Up @@ -822,20 +822,24 @@ impl TurnOrchestrator {
}
};

let repository = atomic_repository::Repository::open_existing_wait(
&self.repo_root,
super::wait_budget::database_wait(),
)
.map_err(|error| AgentError::ProvenanceJournalFailed {
session_id: session_id.to_string(),
reason: error.to_string(),
})?;
let publication = repository
.publish_provenance_checkpoint(&graph, session_turn)
// The owner opens the same database file to record the acknowledgement,
// so this handle must be closed before calling it.
let publication = {
let repository = atomic_repository::Repository::open_existing_wait(
&self.repo_root,
super::wait_budget::database_wait(),
)
.map_err(|error| AgentError::ProvenanceJournalFailed {
session_id: session_id.to_string(),
reason: error.to_string(),
})?;
repository
.publish_provenance_checkpoint(&graph, session_turn)
.map_err(|error| AgentError::ProvenanceJournalFailed {
session_id: session_id.to_string(),
reason: error.to_string(),
})?
};
sink.acknowledge_checkpoint(
&checkpoint,
publication.manifest_hash,
Expand Down
198 changes: 165 additions & 33 deletions atomic-cli/src/commands/agent/owner.rs
Original file line number Diff line number Diff line change
@@ -1,14 +1,19 @@
//! Per-repository owner for the redb-native change and provenance store.
//!
//! redb permits one writable process to open a database. This service elects
//! that process with an OS-backed file lock and exposes a small, versioned local
//! protocol so short-lived agent hooks never open `changes.redb` themselves.
//! redb permits one process at a time to open a database file. This service
//! elects one owner per repository with an OS-backed file lock and exposes a
//! small, versioned local protocol so concurrent agent hooks never race for
//! the provenance journal. The journal lives in the repository database,
//! `atomic.redb`, so the owner opens that file only while it serves a request
//! (see [`StoreLease`]); between requests other atomic processes can open the
//! repository.

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;
use std::time::Duration;
use std::sync::{Arc, Mutex, TryLockError, Weak};
use std::time::{Duration, Instant};

use anyhow::{anyhow, Context};
use atomic_agent::{
Expand All @@ -23,7 +28,7 @@ use atomic_repository::redb_change_store::{
ProvenanceCheckpointSource, ProvenanceId, ProvenanceTurnState, RedbChangeStore, StopCause,
StopState, StoredProvenanceTurn, MAX_FROZEN_PAGE_FRAGMENTS,
};
use atomic_repository::Repository;
use atomic_repository::{ensure_database, Repository, RepositoryError};
use clap::{Args, Subcommand};
use fs2::FileExt;
use serde::{Deserialize, Serialize};
Expand All @@ -45,6 +50,7 @@ const FROZEN_PAGE_METADATA_BYTES: usize = 256 * MAX_FROZEN_PAGE_FRAGMENTS + 256;
const OWNER_LOCK_FILE: &str = "changes-owner.lock";
const START_ATTEMPTS: usize = 80;
const START_RETRY_DELAY: Duration = Duration::from_millis(25);
const DATABASE_RETRY_DELAY: Duration = Duration::from_millis(10);

#[derive(Debug, Args)]
#[command(arg_required_else_help = true)]
Expand Down Expand Up @@ -895,17 +901,105 @@ fn request(repository: &Path, request_body: OwnerRequest) -> anyhow::Result<Owne
fn serve(repository: &Path) -> CliResult<()> {
let dot_dir = Repository::canonical_dot_dir(repository)?;
let owner_lock = acquire_owner_lock(&dot_dir)?;
let store_path = Repository::canonical_change_store_path(repository)?;
let store = Arc::new(
RedbChangeStore::open(&store_path)
.with_context(|| format!("failed to open {}", store_path.display()))?,
);
// Bind before touching the database: clients give a starting owner only
// a short window to answer pings, while opening (or merging a legacy
// layout) may wait for another process. The first store request opens it.
let lease = Arc::new(StoreLease::new(dot_dir.clone()));
let endpoint = endpoint_name(&dot_dir);
runtime()?.block_on(run_server(&endpoint, store))?;
runtime()?.block_on(run_server(&endpoint, lease))?;
FileExt::unlock(&owner_lock).context("failed to release database-owner lock")?;
Ok(())
}

/// Opens the repository database only while store requests are in flight.
///
/// Graph state and the provenance journal share `atomic.redb`, and redb locks
/// the whole file per handle, so an owner that stayed open would lock every
/// other atomic process out of the repository. Concurrent requests share one
/// handle, which closes as soon as the last of them finishes.
struct StoreLease {
dot_dir: PathBuf,
current: Mutex<Weak<RedbChangeStore>>,
}

impl StoreLease {
fn new(dot_dir: PathBuf) -> Self {
Self {
dot_dir,
current: Mutex::new(Weak::new()),
}
}

fn acquire(&self, started: Instant) -> anyhow::Result<Arc<RedbChangeStore>> {
let timeout = atomic_agent::turn::orchestrator::wait_budget::database_wait();
let mut current = loop {
match self.current.try_lock() {
Ok(guard) => break guard,
Err(TryLockError::Poisoned(_)) => {
return Err(anyhow!("database-owner store lease is poisoned"));
}
Err(TryLockError::WouldBlock) if started.elapsed() < timeout => {
std::thread::sleep(
DATABASE_RETRY_DELAY.min(timeout.saturating_sub(started.elapsed())),
);
}
Err(TryLockError::WouldBlock) => {
return Err(anyhow!(
"database-owner store lease stayed busy for {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 open_waiting(&self, started: Instant, timeout: Duration) -> anyhow::Result<RedbChangeStore> {
loop {
match self.open_once() {
Err(OpenFailure::Busy) if started.elapsed() < timeout => {
std::thread::sleep(DATABASE_RETRY_DELAY);
}
Err(OpenFailure::Busy) => {
return Err(anyhow!(
"repository database in {} stayed busy for {:?}",
self.dot_dir.display(),
timeout
))
}
Err(OpenFailure::Other(error)) => return Err(error),
Ok(store) => return Ok(store),
}
}
}

fn open_once(&self) -> Result<RedbChangeStore, OpenFailure> {
let path = match ensure_database(&self.dot_dir) {
Ok(path) => path,
Err(RepositoryError::DatabaseBusy) => return Err(OpenFailure::Busy),
Err(error) => return Err(OpenFailure::Other(error.into())),
};
RedbChangeStore::open_existing(&path).map_err(|error| {
if error.is_database_busy() {
OpenFailure::Busy
} else {
OpenFailure::Other(
anyhow::Error::new(error).context(format!("failed to open {}", path.display())),
)
}
})
}
}

enum OpenFailure {
Busy,
Other(anyhow::Error),
}

fn acquire_owner_lock(dot_dir: &Path) -> anyhow::Result<File> {
let path = dot_dir.join(OWNER_LOCK_FILE);
let file = OpenOptions::new()
Expand Down Expand Up @@ -942,6 +1036,18 @@ fn owner_failpoint(name: &str) {
}

fn handle_request(store: &RedbChangeStore, frame: RequestFrame) -> (ResponseFrame, bool) {
handle_request_with(frame, || Ok(store))
}

/// Answer `Ping` and `Shutdown` without the database: a busy repository must
/// not make a live owner look dead, which would spawn a competing owner.
fn handle_request_with<S>(
frame: RequestFrame,
store: impl FnOnce() -> anyhow::Result<S>,
) -> (ResponseFrame, bool)
where
S: Deref<Target = RedbChangeStore>,
{
let request_id = frame.request_id;
if frame.version != PROTOCOL_VERSION {
return (
Expand All @@ -960,14 +1066,22 @@ fn handle_request(store: &RedbChangeStore, frame: RequestFrame) -> (ResponseFram
);
}

let store = match frame.request {
OwnerRequest::Ping => return respond(request_id, pong(), false),
OwnerRequest::Shutdown => return respond(request_id, shutting_down(), true),
_ => match store() {
Ok(store) => store,
Err(error) => {
let response = OwnerResponse::Error {
code: "database-unavailable".to_string(),
message: format!("{error:#}"),
};
return respond(request_id, response, false);
}
},
};
let (response, shutdown) = match frame.request {
OwnerRequest::Ping => (
OwnerResponse::Pong {
pid: std::process::id(),
frozen_envelope_paging: true,
},
false,
),
OwnerRequest::Ping => (pong(), false),
OwnerRequest::ReserveProvenanceTurn {
session_id,
turn_number,
Expand Down Expand Up @@ -1292,14 +1406,12 @@ fn handle_request(store: &RedbChangeStore, frame: RequestFrame) -> (ResponseFram
false,
),
},
OwnerRequest::Shutdown => (
OwnerResponse::ShuttingDown {
pid: std::process::id(),
},
true,
),
OwnerRequest::Shutdown => (shutting_down(), true),
};
respond(request_id, response, shutdown)
}

fn respond(request_id: String, response: OwnerResponse, shutdown: bool) -> (ResponseFrame, bool) {
(
ResponseFrame {
version: PROTOCOL_VERSION,
Expand All @@ -1310,16 +1422,36 @@ fn handle_request(store: &RedbChangeStore, frame: RequestFrame) -> (ResponseFram
)
}

fn pong() -> OwnerResponse {
OwnerResponse::Pong {
pid: std::process::id(),
frozen_envelope_paging: true,
}
}

fn shutting_down() -> OwnerResponse {
OwnerResponse::ShuttingDown {
pid: std::process::id(),
}
}

async fn handle_connection<S>(
mut stream: S,
store: Arc<RedbChangeStore>,
lease: Arc<StoreLease>,
shutdown: Arc<Notify>,
) -> anyhow::Result<()>
where
S: AsyncRead + AsyncWrite + Unpin,
{
let request: RequestFrame = read_frame(&mut stream).await?;
let (response, should_shutdown) = handle_request(&store, request);
let started = Instant::now();
// 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))
})
.await
.context("database-owner request handler panicked")?;
write_frame(&mut stream, &response).await?;
stream.shutdown().await?;
if should_shutdown {
Expand Down Expand Up @@ -1459,7 +1591,7 @@ async fn connect_owner_pipe(
}

#[cfg(unix)]
async fn run_server(endpoint: &str, store: Arc<RedbChangeStore>) -> anyhow::Result<()> {
async fn run_server(endpoint: &str, lease: Arc<StoreLease>) -> anyhow::Result<()> {
use tokio::net::UnixListener;

let path = Path::new(endpoint);
Expand All @@ -1475,10 +1607,10 @@ async fn run_server(endpoint: &str, store: Arc<RedbChangeStore>) -> anyhow::Resu
tokio::select! {
accepted = listener.accept() => {
let (stream, _) = accepted?;
let store = Arc::clone(&store);
let lease = Arc::clone(&lease);
let shutdown = Arc::clone(&shutdown);
tokio::spawn(async move {
if let Err(error) = handle_connection(stream, store, shutdown).await {
if let Err(error) = handle_connection(stream, lease, shutdown).await {
log::warn!("database-owner connection failed: {error}");
}
});
Expand All @@ -1496,7 +1628,7 @@ async fn run_server(endpoint: &str, store: Arc<RedbChangeStore>) -> anyhow::Resu
}

#[cfg(windows)]
async fn run_server(endpoint: &str, store: Arc<RedbChangeStore>) -> anyhow::Result<()> {
async fn run_server(endpoint: &str, lease: Arc<StoreLease>) -> anyhow::Result<()> {
use tokio::net::windows::named_pipe::ServerOptions;

let shutdown = Arc::new(Notify::new());
Expand All @@ -1515,10 +1647,10 @@ async fn run_server(endpoint: &str, store: Arc<RedbChangeStore>) -> anyhow::Resu
.create(endpoint)
.with_context(|| format!("failed to create database owner pipe {endpoint}"))?;
let connected = std::mem::replace(&mut server, next);
let store = Arc::clone(&store);
let lease = Arc::clone(&lease);
let shutdown = Arc::clone(&shutdown);
tokio::spawn(async move {
if let Err(error) = handle_connection(connected, store, shutdown).await {
if let Err(error) = handle_connection(connected, lease, shutdown).await {
log::warn!("database-owner connection failed: {error}");
}
});
Expand Down
Loading
Loading