Merge PR #823 into release/2.5

feat: add a persistent lifecycle relay companion (#822)

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01MDbhmszrjG9s5MrPrTuNtm

# Conflicts:
#	.github/workflows/ci.yml
#	CHANGELOG.md
#	crates/ai-memory-cli/tests/suite/packaging.rs
#	scripts/install-git-hooks.sh
This commit is contained in:
AkitaOnRails
2026-09-24 17:56:52 -03:00
20 changed files with 6428 additions and 11 deletions
+30 -10
View File
@@ -75,6 +75,18 @@ jobs:
# block that fails to compile (e.g. an indented shell example rustdoc
# reads as Rust) shipped unnoticed. Run them explicitly.
- run: cargo test --workspace --doc
- uses: astral-sh/setup-uv@bec219d24cd3e171d82865faccec33120bb574f4 # v10.1.0
with:
enable-cache: false
- name: Build the external lifecycle relay
env:
CARGO_TARGET_DIR: ${{ github.workspace }}/target
run: cargo build --locked --manifest-path companions/ai-memory-relay/Cargo.toml
- name: External lifecycle delivery against the real server
run: >
uv run --no-project python tests/e2e/external_relay_smoke.py
--ai-memory-bin "${{ github.workspace }}/target/debug/ai-memory"
--relay-bin "${{ github.workspace }}/target/debug/ai-memory-relay"
- name: Vendored tailwind.css is current
if: matrix.os == 'ubuntu-latest'
run: git diff --exit-code -- crates/ai-memory-web/static/tailwind.css
@@ -125,13 +137,13 @@ jobs:
# is here to catch before it reaches windows.yml or a release.
- run: cargo xwin build --workspace --all-targets --target x86_64-pc-windows-msvc
# The companion importer is deliberately outside the root workspace
# (its own [workspace] in companions/ai-memory-importer/Cargo.toml), so
# none of the jobs above compile it. Gate it here with the same
# fmt/clippy/test trio docs/companion-crates.md prescribes, or a
# toolchain bump / convention change rots it silently.
# Companions have their own workspaces and need separate checks.
companions:
name: companions (ai-memory-importer)
name: companions (${{ matrix.companion }})
strategy:
fail-fast: false
matrix:
companion: [ai-memory-importer, ai-memory-relay]
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@3d3c42e5aac5ba805825da76410c181273ba90b1 # v7.0.1
@@ -140,10 +152,10 @@ jobs:
components: rustfmt, clippy
- uses: Swatinem/rust-cache@f0d9c3887740aee45f6153b24b3a6b815192ec16 # v2
with:
workspaces: companions/ai-memory-importer
- run: cargo fmt --check --manifest-path companions/ai-memory-importer/Cargo.toml
- run: cargo clippy --manifest-path companions/ai-memory-importer/Cargo.toml --all-targets -- -D warnings
- run: cargo test --manifest-path companions/ai-memory-importer/Cargo.toml
workspaces: companions/${{ matrix.companion }}
- run: cargo fmt --check --manifest-path companions/${{ matrix.companion }}/Cargo.toml
- run: cargo clippy --locked --manifest-path companions/${{ matrix.companion }}/Cargo.toml --all-targets -- -D warnings
- run: cargo test --locked --manifest-path companions/${{ matrix.companion }}/Cargo.toml
# Build the release binary on its own so a release-only failure
# (LTO crash, codegen issue, dead-code-with-debug-assertions) doesn't
@@ -308,6 +320,13 @@ jobs:
log-level: warn
command: check
arguments: --all-features
- name: Relay dependency policy
uses: EmbarkStudios/cargo-deny-action@3c6349835b2b7b196a839186cb8b78e02f7b5f25 # v2
with:
log-level: warn
command: check
manifest-path: companions/ai-memory-relay/Cargo.toml
arguments: --all-features
audit:
name: cargo-audit
@@ -323,6 +342,7 @@ jobs:
--ignore RUSTSEC-2024-0320
--ignore RUSTSEC-2026-0194
--ignore RUSTSEC-2026-0195
- run: cargo audit --file companions/ai-memory-relay/Cargo.lock
gitleaks:
name: gitleaks (secret scan)
+21
View File
@@ -70,6 +70,27 @@ jobs:
- uses: dtolnay/rust-toolchain@4cda84d5c5c54efe2404f9d843567869ab1699d4 # stable
- uses: Swatinem/rust-cache@f0d9c3887740aee45f6153b24b3a6b815192ec16 # v2
- run: cargo test --workspace --all-targets
- name: Check the external lifecycle relay
env:
CARGO_TARGET_DIR: ${{ github.workspace }}/target
run: |
cargo fmt --check --manifest-path companions/ai-memory-relay/Cargo.toml
if ($LASTEXITCODE -ne 0) { exit $LASTEXITCODE }
cargo clippy --locked --manifest-path companions/ai-memory-relay/Cargo.toml --all-targets -- -D warnings
if ($LASTEXITCODE -ne 0) { exit $LASTEXITCODE }
cargo test --locked --manifest-path companions/ai-memory-relay/Cargo.toml
if ($LASTEXITCODE -ne 0) { exit $LASTEXITCODE }
cargo build --locked --manifest-path companions/ai-memory-relay/Cargo.toml
if ($LASTEXITCODE -ne 0) { exit $LASTEXITCODE }
- uses: astral-sh/setup-uv@bec219d24cd3e171d82865faccec33120bb574f4 # v10.1.0
with:
enable-cache: false
- name: External lifecycle delivery against the real server
run: >
uv run --no-project python tests/e2e/external_relay_smoke.py
--ai-memory-bin "${{ github.workspace }}/target/debug/ai-memory.exe"
--relay-bin "${{ github.workspace }}/target/debug/ai-memory-relay.exe"
--timeout-seconds 120
# `ci.yml`'s `hooks-shell` job covers the POSIX bundle across four awks, all
# on Linux. The suite also drives `hooks/lib/ai-memory-hook.ps1`, and that
+4
View File
@@ -64,6 +64,10 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
pre-existing identity, so upgrading installs need no migration, and a
query-only prefix change never needs a rebuild. See
`docs/llm-providers.md`. (#859)
- Added the optional `ai-memory-relay` companion for external lifecycle events.
It queues events locally and sends ordered batches through the public hook API,
using stable event identities for retries. The package has its own workspace
and does not change the server's capture or storage defaults. (#823)
- `auto_improve.patchable_page_prefixes` makes the folders whose page bodies the
reviewer reads configurable, defaulting to the historical `_rules/` and
`procedures/`. Only those two folders' contents were ever sent; every other
+11
View File
@@ -93,6 +93,17 @@ the block's shell options stay inside it; a failing test run still fails the
hook even when your own commands follow the block without `set -e`.
Run the installer again to update an existing installation.
The managed test block clears Git's repository environment and disables global
and system Git configuration for Cargo and its children. Fixture commands can
then use their own repositories without inheriting the checkout being pushed.
The publishing Git process and other hook code retain their configuration.
Run the installer again to update an existing installation.
Companions have separate Cargo workspaces. Check each changed companion with
`cargo fmt`, `cargo clippy --all-targets -- -D warnings`, and `cargo test`, passing
its `--manifest-path`. Changes to the lifecycle relay also need the real-server
test documented in [its README](companions/ai-memory-relay/README.md#validation).
Integration tests live in `tests/suite/` per crate and compile into the
crate's own test harness (declare a new file with `mod name;` in
`tests/suite/mod.rs`); only the CLI keeps a separate test binary, because its
+1 -1
View File
@@ -483,7 +483,7 @@ diagram, crate breakdown, schema notes, and invariants.
| [`docs/ARCHITECTURE.md`](docs/ARCHITECTURE.md) | Operational summary: data flow, crate layout, cross-cutting invariants, schema. |
| [`docs/design-decisions.md`](docs/design-decisions.md) | The full v1 spec. |
| [`docs/managed-harness-contributions.md`](docs/managed-harness-contributions.md) | Protocol and acceptance bar for adding managed resume, transcript import, and startup context delivery to another harness. |
| [`docs/companion-crates.md`](docs/companion-crates.md) | Boundary and plan for optional companion projects, including the standalone importer at [`companions/ai-memory-importer`](companions/ai-memory-importer). |
| [`docs/companion-crates.md`](docs/companion-crates.md) | Optional companion projects: the [importer](companions/ai-memory-importer) and [external lifecycle relay](companions/ai-memory-relay). |
| [`docs/external-lifecycle.md`](docs/external-lifecycle.md) | External lifecycle producers: per-execution native capture suppression, preserved handoffs, batch ingestion and stable retry identity. |
| [`docs/auto-improvement-loop.md`](docs/auto-improvement-loop.md) | Auto-improvement design notes: scheduled review, auto-approval default, manual review opt-in, pending proposal storage, and curator work. |
File diff suppressed because it is too large Load Diff
+38
View File
@@ -0,0 +1,38 @@
# Separate workspace, as described in docs/companion-crates.md.
[workspace]
[package]
name = "ai-memory-relay"
version = "0.1.0"
edition = "2024"
publish = false
license = "MIT"
description = "Persistent queue for external ai-memory lifecycle events"
[[bin]]
name = "ai-memory-relay"
path = "src/main.rs"
[lints.rust]
unsafe_code = "forbid"
[dependencies]
anyhow = "1"
# Read the bearer token separately to keep it out of clap's Debug output.
clap = { version = "4", features = ["derive"] }
fs2 = "0.4"
# Native roots match the workspace (#492).
reqwest = { version = "0.12", default-features = false, features = [
"blocking",
"json",
"rustls-tls-native-roots",
] }
rusqlite = { version = "0.32", features = ["bundled"] }
serde = { version = "1", features = ["derive"] }
# Default map ordering makes content comparisons independent of input key order.
serde_json = "1"
sha2 = "0.10"
url = "2"
[dev-dependencies]
tempfile = "3"
+215
View File
@@ -0,0 +1,215 @@
# ai-memory-relay
A companion CLI for sending externally captured lifecycle events to ai-memory.
It records events in a local SQLite queue and sends them through `POST /hook/batch`.
Use it when an orchestrator needs delivery to survive an unavailable server or a
process restart.
The relay has its own Cargo workspace. It does not open ai-memory's database or
wiki, launch agents, or claim handoffs. See the
[external lifecycle contract](../../docs/external-lifecycle.md) for capture
ownership and the [companion policy](../../docs/companion-crates.md) for the core
boundary.
## Build and use
From the repository root:
```sh
cargo build --locked --manifest-path companions/ai-memory-relay/Cargo.toml
```
The executable is `companions/ai-memory-relay/target/debug/ai-memory-relay`, unless
`CARGO_TARGET_DIR` selects another target directory. The commands below assume
the executable is on `PATH` and an ai-memory server is running.
Initialize a queue. Its parent directory must exist; the relay creates the queue
directory itself.
```sh
ai-memory-relay init \
--queue-dir "$HOME/.ai-memory-relay" \
--server-url http://127.0.0.1:49374 \
--producer my-orchestrator \
--actor developer \
--workspace work \
--project example
```
This binds the queue to one destination and scope. Repeating an identical `init`
is safe. A different binding is rejected. `actor` is a stable namespace used to
derive retry keys; authentication still comes from the server's bearer token.
It does not select or impersonate a server user.
When the orchestrator launches a harness whose lifecycle it captures, pass
`AI_MEMORY_CAPTURE_OWNER=my-orchestrator` to that process. Native capture is then
suppressed as described in the lifecycle contract. Supported handoff delivery
remains active, and the agent can keep using MCP for retrieval and deliberate
memory writes. A standalone harness without that context keeps normal hooks.
The orchestrator supplies a JSON array, for example `events.json`:
```json
[
{
"event_id": "event-001",
"agent": "claude-code",
"event": "session-start",
"body": {
"session_id": "native-session-123",
"cwd": "/workspace/example"
}
},
{
"event_id": "event-002",
"agent": "claude-code",
"event": "user-prompt-submit",
"body": {
"session_id": "native-session-123",
"cwd": "/workspace/example",
"prompt": "Check the failing test."
}
}
]
```
```sh
ai-memory-relay enqueue --queue-dir "$HOME/.ai-memory-relay" --file events.json
ai-memory-relay flush --queue-dir "$HOME/.ai-memory-relay"
ai-memory-relay status --queue-dir "$HOME/.ai-memory-relay"
```
`enqueue` validates the whole array and commits it in one transaction. Each body
must be an object with explicit `session_id` and `cwd` strings. Use the harness's
native identity and canonical hook event names. Other body fields keep their
values, including an `_ai_memory_capture` block. The relay does not translate
provider transcripts or infer parent agents and workflows.
For authenticated servers, supply `AI_MEMORY_AUTH_TOKEN` through the flush
process's environment. The relay does not store or print it. The destination
must use HTTP or HTTPS, with no credentials, query string or fragment in the
base URL. Redirects are refused.
## Delivery and retries
Each event gets a deterministic `ingest_key` from the producer, actor, native
agent/session identity, event name and producer-assigned `event_id`, following
the lifecycle contract. `extension` carries the producer and `source_event`
carries the event name.
Re-enqueueing the same identity and body is recognized while the pending event
or receipt remains in the queue. JSON object key order does not matter. Reusing
that identity with different content rejects the entire input array. Two equal
bodies with different event IDs remain separate events.
Only the oldest pending event from each session enters a batch. The next event
for that session is eligible after acknowledgement. This keeps `session-end`
behind earlier events while allowing other sessions to proceed. Order is the
queue's committed enqueue order; producers must serialize events within a
session before enqueueing them. Separate queue directories do not coordinate
session order.
One process can flush a queue at a time. Other processes can enqueue during a
flush. HTTP runs outside the SQLite write transaction. If a response is lost
after the server commits, the event stays pending and a retry uses the same key.
The relay validates the complete batch acknowledgement before changing its
queue. It honors both contiguous and noncontiguous accepted indexes, including
partial success in HTTP 200 or 429 responses. Malformed acknowledgements,
authentication failures and transport errors retain unacknowledged events.
A failed head blocks its own session. HTTP 429 ends the flush; the caller should
back off before scheduling another one.
An acknowledgement can also mean the server deliberately dropped an event, for
example because of capture policy or a session identity collision. A drained
queue does not prove that every event became an observation. The relay catches
a native session changing agents within one queue, but server ownership checks
still apply across queues and producers.
The server's ingest keys expire after 30 days. The relay records the first
attempt before sending and retains expired pending events without replaying
them automatically. Unsent events can remain offline longer. Keep the host
clock accurate; the retry deadline uses wall-clock timestamps. Recreating a
queue discards its attempt history and receipts, so it is not a safe way to
recover an ambiguous delivery after the retry window.
`flush` attempts at most 64 batches by default; `--max-batches` changes that
count. Requests, transport retries and the total flush duration also have finite
bounds. Unfinished work stays queued for the next invocation.
Exit codes are:
| Code | Meaning |
| --- | --- |
| `0` | Command succeeded; a flush has no pending events. |
| `2` | Command failed. Inspect stderr; unacknowledged events remain queued. |
| `3` | Flush ended with events still pending. |
`status` emits a JSON object with counters, including `pending_items`,
`pending_bytes`, `pending_sessions`, `expired_items` and `receipts`. It does not
print event bodies. Receipts count acknowledgements, including policy drops.
## Local data and recovery
The queue contains producer-supplied event bodies before server sanitization.
The producer must apply its capture exclusions before enqueueing. The relay
does not load the project's `[capture] ignore_paths` settings or inspect tool
payloads for sensitive paths. Server sanitization cannot protect a local queue
that already contains those payloads.
On Unix, new queue directories use mode `0700` and files use `0600`. Existing
queue paths must already be private. Queue symlinks and Windows reparse points
are rejected. On Windows, the relay does not set ACLs; use a directory restricted
to the account that runs the orchestrator.
SQLite uses WAL mode and `synchronous=FULL`. Keep the database and its sidecar
files together. To move or back up a queue, stop producers and flush processes
first. Reopen the same queue after a restart; do not edit its SQLite tables to
mark events as delivered. An unknown queue schema is rejected.
The queue uses these limits:
| Resource | Limit |
| --- | --- |
| Pending events | 50,000 |
| Pending body bytes | 64 MiB |
| Pending events plus retained receipts | 200,000 |
| Retained session identities | 50,000 |
| One event body | 256 KiB |
| One input file | 32 MiB and 10,000 events |
| One HTTP batch | 8 MiB and 256 events |
| Acknowledgement body | 64 KiB |
Reaching a limit rejects new input instead of dropping older events. Receipts
and unused session records can expire after the 30-day retry window. Pending
events are retained. Review expired entries and producer failures before
choosing how to archive a queue; the CLI has no command that discards them
automatically. These are logical record limits, not a fixed SQLite file size.
## Validation
The package needs separate checks because it is outside the root workspace:
```sh
cargo fmt --check --manifest-path companions/ai-memory-relay/Cargo.toml
cargo clippy --locked --manifest-path companions/ai-memory-relay/Cargo.toml --all-targets -- -D warnings
cargo test --locked --manifest-path companions/ai-memory-relay/Cargo.toml
```
Run the integration test against locally built binaries:
```sh
cargo build --workspace
cargo build --locked --manifest-path companions/ai-memory-relay/Cargo.toml
uv run --no-project python tests/e2e/external_relay_smoke.py \
--ai-memory-bin "$PWD/target/debug/ai-memory" \
--relay-bin "$PWD/companions/ai-memory-relay/target/debug/ai-memory-relay"
```
Adjust the binary paths if `CARGO_TARGET_DIR` is set; Windows executables end in
`.exe`. The test starts an isolated real ai-memory server and invokes its native
hooks. It checks persistent observations for offline recovery, retries after a
lost response, distinct event identities and 15 concurrent sessions. It also
checks handoff delivery with native capture suppressed, rejected authentication
and acknowledged session collisions. Temporary data and logs are retained, and
the test prints their directory.
+98
View File
@@ -0,0 +1,98 @@
//! Acknowledgement parsing for `POST /hook/batch`.
//!
//! Validate the entire response before releasing any queued event. An
//! inconsistent acknowledgement leaves the whole batch pending for retry.
use serde::Deserialize;
/// The server's ack. Unknown fields are accepted (the server may add some);
/// a missing `accepted` is a malformed ack and preserves the batch.
#[derive(Debug, Clone, Deserialize)]
pub struct BatchAck {
/// Contiguous leading prefix committed, oldest-first.
pub accepted: usize,
/// Non-contiguous committed indexes, when per-source rate limiting skipped items.
#[serde(default)]
pub accepted_indices: Option<Vec<usize>>,
/// Item that failed processing after earlier skips.
#[serde(default)]
pub failed_index: Option<usize>,
}
/// Why an ack was refused. The batch stays pending in every case.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AckRejected(pub String);
impl std::fmt::Display for AckRejected {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
/// Validate an ack against the batch it answers, returning the indexes that may
/// be released.
///
/// `accepted` is the contiguous leading prefix *even when* `accepted_indices` is
/// present, so the two must agree. Duplicated or out-of-range indexes, a
/// disagreeing `accepted`, a `failed_index` outside the batch, and a
/// `failed_index` that also claims to be accepted are all refusals.
pub fn validate(batch_len: usize, ack: &BatchAck) -> Result<Vec<usize>, AckRejected> {
let reject = |detail: String| Err(AckRejected(detail));
let accepted: Vec<usize> = match &ack.accepted_indices {
Some(indices) => {
if let Some(bad) = indices.iter().find(|idx| **idx >= batch_len) {
return reject(format!(
"accepted_indices contains {bad}, outside a {batch_len}-item batch"
));
}
if indices.windows(2).any(|pair| pair[0] >= pair[1]) {
return reject(
"accepted_indices must be strictly ascending and duplicate-free".into(),
);
}
let prefix = indices
.iter()
.enumerate()
.take_while(|(pos, idx)| pos == *idx)
.count();
if prefix != ack.accepted {
return reject(format!(
"accepted={} disagrees with the {prefix}-item contiguous prefix of accepted_indices",
ack.accepted
));
}
indices.clone()
}
None => {
if ack.accepted > batch_len {
return reject(format!(
"accepted={} exceeds the {batch_len} items sent",
ack.accepted
));
}
(0..ack.accepted).collect()
}
};
if let Some(failed) = ack.failed_index {
if failed >= batch_len {
return reject(format!(
"failed_index={failed} is outside a {batch_len}-item batch"
));
}
if accepted.contains(&failed) {
return reject(format!(
"failed_index={failed} is also reported as accepted"
));
}
// The server fails fast: it stops at `failed_index` and processes
// nothing after it. An ack claiming a later item committed describes a
// run that cannot have happened, so the batch is preserved whole.
if let Some(after) = accepted.iter().find(|idx| **idx > failed) {
return reject(format!(
"accepted index {after} comes after failed_index={failed}, which the server \
never processes past"
));
}
}
Ok(accepted)
}
+220
View File
@@ -0,0 +1,220 @@
//! Filesystem guards for the queue directory.
//!
//! Event bodies reach this directory before server sanitization. Unix paths
//! must be private; existing permissions are checked without changing them.
//!
//! Windows: only symlink/reparse-point rejection is implemented. No ACL
//! hardening is performed, so a queue directory there is exactly as private as
//! the operator made it.
use std::path::{Component, Path, PathBuf};
use anyhow::{Context, Result, bail};
/// SQLite sidecars that must be checked alongside the database itself.
pub const SIDECARS: &[&str] = &[
"relay.sqlite-wal",
"relay.sqlite-shm",
"relay.sqlite-journal",
];
/// The queue database file name.
pub const DB_FILE: &str = "relay.sqlite";
/// Reject a path that is a symlink or a Windows reparse point.
///
/// A missing path is fine: it is about to be created inside a directory that
/// was itself checked.
pub fn reject_symlink(path: &Path) -> Result<()> {
let meta = match std::fs::symlink_metadata(path) {
Ok(meta) => meta,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(e) => return Err(e).with_context(|| format!("inspect {}", path.display())),
};
if meta.file_type().is_symlink() {
bail!(
"{} is a symlink; the relay refuses to follow one into a queue directory",
path.display()
);
}
#[cfg(windows)]
{
use std::os::windows::fs::MetadataExt;
const FILE_ATTRIBUTE_REPARSE_POINT: u32 = 0x400;
if meta.file_attributes() & FILE_ATTRIBUTE_REPARSE_POINT != 0 {
bail!(
"{} is a reparse point; the relay refuses to follow one into a queue directory",
path.display()
);
}
}
Ok(())
}
/// Check one queue file: not a symlink/reparse point, a regular file, not
/// hardlinked elsewhere, and owner-only on Unix. A missing file passes.
pub fn check_queue_file(path: &Path) -> Result<()> {
reject_symlink(path)?;
let meta = match std::fs::symlink_metadata(path) {
Ok(meta) => meta,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(()),
Err(e) => return Err(e).with_context(|| format!("inspect {}", path.display())),
};
if !meta.is_file() {
bail!(
"{} exists but is not a regular file; refusing to use it as queue state",
path.display()
);
}
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
use std::os::unix::fs::PermissionsExt;
if meta.nlink() > 1 {
bail!(
"{} has {} hard links; refusing to write queue state through a shared inode",
path.display(),
meta.nlink()
);
}
let mode = meta.permissions().mode();
if mode & 0o077 != 0 {
bail!(
"{} is mode {:o}; queue state must be owner-only (chmod 600) \
(existing permissions are unchanged)",
path.display(),
mode & 0o7777
);
}
}
Ok(())
}
/// Prepare `--queue-dir`.
///
/// A directory the relay creates is created owner-only. An existing directory
/// must already be private and must be either empty or an existing relay queue;
/// anything else is refused untouched, so pointing the relay at a shared
/// directory can never re-permission it.
pub fn prepare_queue_dir(dir: &Path) -> Result<PathBuf> {
if dir.components().any(|c| c == Component::ParentDir) {
bail!(
"--queue-dir must not contain a `..` component: {}",
dir.display()
);
}
reject_symlink(dir)?;
match std::fs::symlink_metadata(dir) {
Ok(meta) if !meta.is_dir() => bail!("--queue-dir {} is not a directory", dir.display()),
Ok(_) => check_existing_dir(dir)?,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => create_private_dir(dir)?,
Err(e) => return Err(e).with_context(|| format!("inspect {}", dir.display())),
}
for name in std::iter::once(DB_FILE).chain(SIDECARS.iter().copied()) {
check_queue_file(&dir.join(name))?;
}
std::fs::canonicalize(dir).with_context(|| format!("resolve {}", dir.display()))
}
/// An existing directory must be private and dedicated to this queue.
fn check_existing_dir(dir: &Path) -> Result<()> {
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mode = std::fs::metadata(dir)?.permissions().mode();
if mode & 0o077 != 0 {
bail!(
"--queue-dir {} is mode {:o}; it must be owner-only (chmod 700) before the relay \
will store event bodies in it. The relay does not re-permission a directory it \
did not create",
dir.display(),
mode & 0o7777
);
}
}
let mut foreign = Vec::new();
for entry in std::fs::read_dir(dir).with_context(|| format!("read {}", dir.display()))? {
let name = entry?.file_name();
let name = name.to_string_lossy().to_string();
let known = name == DB_FILE || name == "flush.lock" || SIDECARS.contains(&name.as_str());
if !known {
foreign.push(name);
}
}
if !foreign.is_empty() {
foreign.sort();
foreign.truncate(5);
bail!(
"--queue-dir {} already holds unrelated files ({}); point the relay at an empty or \
existing queue directory instead of sharing one",
dir.display(),
foreign.join(", ")
);
}
Ok(())
}
#[cfg(unix)]
fn create_private_dir(dir: &Path) -> Result<()> {
use std::os::unix::fs::DirBuilderExt;
if let Some(parent) = dir.parent()
&& !parent.as_os_str().is_empty()
&& !parent.exists()
{
bail!(
"parent of --queue-dir {} does not exist; create it deliberately first",
dir.display()
);
}
std::fs::DirBuilder::new()
.mode(0o700)
.create(dir)
.with_context(|| format!("create {}", dir.display()))
}
#[cfg(not(unix))]
fn create_private_dir(dir: &Path) -> Result<()> {
if let Some(parent) = dir.parent()
&& !parent.as_os_str().is_empty()
&& !parent.exists()
{
bail!(
"parent of --queue-dir {} does not exist; create it deliberately first",
dir.display()
);
}
std::fs::create_dir(dir).with_context(|| format!("create {}", dir.display()))
}
/// Create a queue file owner-only, or leave an existing one alone.
///
/// SQLite copies the database mode onto its sidecars. Set 0600 at creation so
/// WAL and SHM files are private from their first write.
#[cfg(unix)]
pub fn create_private_file(path: &Path) -> Result<()> {
use std::os::unix::fs::OpenOptionsExt;
match std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.mode(0o600)
.open(path)
{
Ok(_) => Ok(()),
// Another process won the race and made it; its own guard applies.
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
}
}
/// No ACL hardening on non-Unix targets; documented in the README.
#[cfg(not(unix))]
pub fn create_private_file(path: &Path) -> Result<()> {
match std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(path)
{
Ok(_) => Ok(()),
Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => Ok(()),
Err(e) => Err(e).with_context(|| format!("create {}", path.display())),
}
}
+200
View File
@@ -0,0 +1,200 @@
//! The only network surface: `POST <server>/hook/batch`.
//!
//! Per-item URLs are always constructed here from the bound destination and the
//! queue's own fields. An input file can never steer a request anywhere.
use std::io::Read;
use std::time::Duration;
use anyhow::{Context, Result, bail};
use url::Url;
use crate::queue::{Binding, PendingItem};
/// Bearer material. Never stored on disk, never printed, never in a Debug line.
pub struct Token(String);
impl std::fmt::Debug for Token {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("Token(<redacted>)")
}
}
impl Token {
/// Read `AI_MEMORY_AUTH_TOKEN` from the environment of *this* flush.
///
/// The value is never echoed, not even when it is rejected.
pub fn from_env() -> Result<Option<Self>> {
let Ok(raw) = std::env::var("AI_MEMORY_AUTH_TOKEN") else {
return Ok(None);
};
let trimmed = raw.trim();
if trimmed.is_empty() {
return Ok(None);
}
if trimmed.chars().any(|c| c.is_control()) {
bail!("AI_MEMORY_AUTH_TOKEN contains control characters; refusing to build a request");
}
Ok(Some(Self(trimmed.to_owned())))
}
fn expose(&self) -> &str {
&self.0
}
}
/// Normalize and harden `--server-url`.
///
/// http/https only, a host, no userinfo, no query, no fragment. No host
/// allowlist: a remote HTTPS server is a legitimate deployment.
pub fn normalize_server_url(raw: &str) -> Result<String> {
let url = Url::parse(raw).context("invalid --server-url")?;
if !matches!(url.scheme(), "http" | "https") {
bail!("--server-url must be http or https, got {:?}", url.scheme());
}
if !url.username().is_empty() || url.password().is_some() {
bail!("--server-url must not carry userinfo; pass credentials in AI_MEMORY_AUTH_TOKEN");
}
if url.query().is_some() {
bail!("--server-url must not carry a query string");
}
if url.fragment().is_some() {
bail!("--server-url must not carry a fragment");
}
let host = url.host_str().context("--server-url needs a host")?;
let authority = match url.port() {
Some(port) => format!("{host}:{port}"),
None => host.to_owned(),
};
let path = url.path().trim_end_matches('/');
Ok(format!("{}://{}{}", url.scheme(), authority, path))
}
/// The `{url, body}` pair the batch endpoint expects.
#[derive(Debug, serde::Serialize)]
pub struct BatchItem {
pub url: String,
pub body: serde_json::Value,
}
/// Build one item's hook URL from the binding plus the queued event.
///
/// `extension` carries the producer namespace, `source_event` repeats the
/// canonical event explicitly (both are required to preserve provenance), and
/// `ingest_key` is the stable key persisted with the event at enqueue time.
pub fn item_url(base: &str, binding: &Binding, item: &PendingItem) -> String {
let query = url::form_urlencoded::Serializer::new(String::new())
.append_pair("event", &item.event)
.append_pair("agent", &item.agent)
.append_pair("workspace", &binding.workspace)
.append_pair("project", &binding.project)
.append_pair("extension", &binding.producer)
.append_pair("source_event", &item.event)
.append_pair("ingest_key", &item.ingest_key)
.finish();
format!("{base}/hook?{query}")
}
/// Largest ack body the relay will read. A 256-index ack is under 2 KiB; 64 KiB
/// leaves room for future fields while refusing to stream an unbounded response
/// from a server that is not behaving like ai-memory.
pub const MAX_ACK_BYTES: usize = 64 * 1024;
/// One completed HTTP exchange. The body is bounded and never logged.
pub struct Delivery {
pub status: u16,
pub body: Vec<u8>,
}
impl std::fmt::Debug for Delivery {
// Length only: an ack body is server output, and output never becomes a log
// line here by accident.
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Delivery")
.field("status", &self.status)
.field("body_bytes", &self.body.len())
.finish()
}
}
/// Blocking batch sender.
pub struct Sender {
client: reqwest::blocking::Client,
base: String,
token: Option<Token>,
}
impl std::fmt::Debug for Sender {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Sender")
.field("base", &self.base)
.field("token", &self.token)
.finish()
}
}
impl Sender {
/// Build the client. Redirects are disabled: a 3xx from the ingestion
/// endpoint must never resend the bearer header, or the batch, elsewhere.
/// The read timeout is per request, because a batch's budget scales with the
/// number of items it carries.
pub fn new(base: String, token: Option<Token>, connect_timeout: Duration) -> Result<Self> {
let client = reqwest::blocking::Client::builder()
.redirect(reqwest::redirect::Policy::none())
.connect_timeout(connect_timeout)
.build()
.context("build HTTP client")?;
Ok(Self {
client,
base,
token,
})
}
/// POST one batch under an explicit timeout.
///
/// Transport failures carry a short class label only: never the request, the
/// response, or the bearer header.
pub fn post_batch(&self, items: &[BatchItem], timeout: Duration) -> Result<Delivery> {
let mut request = self
.client
.post(format!("{}/hook/batch", self.base))
.timeout(timeout)
.json(items);
if let Some(token) = &self.token {
request = request.bearer_auth(token.expose());
}
let response = request.send().map_err(|e| {
anyhow::anyhow!("hook batch transport failure: {}", transport_class(&e))
})?;
let status = response.status().as_u16();
// Bounded read: `bytes()` would accept whatever the peer sends.
let mut body = Vec::new();
let mut limited = response.take(MAX_ACK_BYTES as u64 + 1);
if limited.read_to_end(&mut body).is_err() {
body.clear();
}
if body.len() > MAX_ACK_BYTES {
bail!(
"hook batch ack exceeded {MAX_ACK_BYTES} bytes; refusing to parse it \
(the batch stays pending)"
);
}
Ok(Delivery { status, body })
}
}
/// A short, content-free label for a transport error.
fn transport_class(error: &reqwest::Error) -> &'static str {
if error.is_timeout() {
"timeout"
} else if error.is_connect() {
"connect"
} else if error.is_redirect() {
"redirect refused"
} else if error.is_request() {
"request"
} else {
"io"
}
}
+224
View File
@@ -0,0 +1,224 @@
//! Envelope validation and the stable retry identity.
//!
//! Validate input before enqueueing so malformed events cannot block a batch.
//! Body values, including `_ai_memory_capture`, are preserved for the server.
//! These checks do not sanitize locally queued content.
use serde::Deserialize;
use sha2::{Digest, Sha256};
/// `extension` accepts up to 64 ASCII token characters (docs/external-lifecycle.md).
pub const MAX_PRODUCER_LEN: usize = 64;
/// A stable adapter-side operator namespace. Never a bearer token.
pub const MAX_ACTOR_LEN: usize = 64;
/// `source_event` accepts 128 (docs/external-lifecycle.md).
pub const MAX_EVENT_LEN: usize = 128;
/// Producer-assigned event id. It only feeds the ingest-key hash, so it is
/// bounded free-form text rather than a
/// token: an orchestrator numbering events `run/123/event/2` is normal.
pub const MAX_EVENT_ID_LEN: usize = 128;
/// Wire `agent` such as `claude-code` or `codex`.
pub const MAX_AGENT_LEN: usize = 64;
/// Workspace / project names, URL-encoded into the item query.
pub const MAX_SCOPE_LEN: usize = 128;
/// Native session id, preserved exactly as the harness minted it.
pub const MAX_SESSION_ID_LEN: usize = 256;
/// Native cwd, preserved exactly as the harness reported it.
pub const MAX_CWD_LEN: usize = 4096;
/// Canonicalized body bytes accepted for one event.
pub const MAX_BODY_BYTES: usize = 256 * 1024;
/// Items in one `POST /hook/batch` (server: `MAX_HOOK_BATCH_ITEMS`).
pub const MAX_BATCH_ITEMS: usize = 256;
/// Server ingest keys expire after 30 days; nothing older may be auto-resent.
pub const RETRY_WINDOW_MS: i64 = 30 * 24 * 60 * 60 * 1000;
/// Lifecycle events that end an execution. A terminal event may only ship once
/// it is its session's oldest pending item, so it can never pass an earlier one.
pub const TERMINAL_EVENTS: &[&str] = &["session-end"];
/// One item of the `enqueue --file` input array.
///
/// `deny_unknown_fields` applies to the envelope only. `body` is an opaque
/// object: unknown keys inside it (a harness payload field, the
/// `_ai_memory_capture` protocol block) are preserved untouched.
#[derive(Debug, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct InputEvent {
/// Producer-assigned id, unique within the producer/actor/agent/session tuple.
pub event_id: String,
/// Wire agent of the live harness (`claude-code`, `codex`, ...).
pub agent: String,
/// Canonical lifecycle event name, used for both `event` and `source_event`.
pub event: String,
/// Harness payload, delivered verbatim.
pub body: serde_json::Value,
}
/// A validated event, ready to be persisted and later replayed unchanged.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ValidEvent {
pub event_id: String,
pub agent: String,
pub event: String,
pub session_id: String,
pub cwd: String,
/// Canonical JSON of the body exactly as supplied (sorted keys, no reformatting
/// of values). The bytes delivered later are these bytes.
pub body_json: String,
pub body_sha256: String,
pub ingest_key: String,
}
impl ValidEvent {
/// True when this event ends the execution for its session.
pub fn is_terminal(&self) -> bool {
TERMINAL_EVENTS.contains(&self.event.as_str())
}
}
/// `extension`, `source_event` and `ingest_key` share the server's token
/// alphabet: letters, digits, `.`, `_`, `-` and `:`.
fn is_token(value: &str, max: usize) -> bool {
!value.is_empty()
&& value.len() <= max
&& value
.bytes()
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'.' | b'_' | b'-' | b':'))
}
/// Free-form identity text (session id, cwd, scope names): bounded, non-blank,
/// and free of control characters that would corrupt a URL or a log line.
fn is_plain(value: &str, max: usize) -> bool {
!value.trim().is_empty()
&& value.chars().count() <= max
&& !value.chars().any(|c| c.is_control())
}
/// Validate the fields that identify a producer namespace.
pub fn check_producer(producer: &str) -> Result<(), String> {
is_token(producer, MAX_PRODUCER_LEN)
.then_some(())
.ok_or_else(|| format!("--producer must be 1..={MAX_PRODUCER_LEN} token characters"))
}
/// Validate the adapter-side operator namespace.
pub fn check_actor(actor: &str) -> Result<(), String> {
is_token(actor, MAX_ACTOR_LEN)
.then_some(())
.ok_or_else(|| format!("--actor must be 1..={MAX_ACTOR_LEN} token characters"))
}
/// Validate a workspace or project name.
pub fn check_scope(label: &str, value: &str) -> Result<(), String> {
is_plain(value, MAX_SCOPE_LEN)
.then_some(())
.ok_or_else(|| format!("--{label} must be 1..={MAX_SCOPE_LEN} printable characters"))
}
/// The retry identity from `docs/external-lifecycle.md`, byte-for-byte.
///
/// The tuple defines the namespace of `event_id`. Payload content does not
/// contribute to the key, so a restarted producer can derive the same identity.
pub fn ingest_key(
producer: &str,
actor: &str,
agent: &str,
session_id: &str,
source_event: &str,
event_id: &str,
) -> String {
let identity = serde_json::to_string(&[
"external-capture-v1",
producer,
actor,
agent,
session_id,
source_event,
event_id,
])
.expect("a fixed-size array of strings always serializes");
let digest = Sha256::digest(identity.as_bytes());
let mut key = String::with_capacity(64);
for byte in digest {
use std::fmt::Write as _;
let _ = write!(key, "{byte:02x}");
}
key
}
/// Validate one input item against every bound above.
///
/// Errors name the array index and the producer's own `event_id`. They never
/// quote body content: an invalid payload must not become a log line.
pub fn validate(
index: usize,
input: InputEvent,
producer: &str,
actor: &str,
) -> Result<ValidEvent, String> {
let at = |detail: &str| format!("item {index}: {detail}");
if !is_plain(&input.event_id, MAX_EVENT_ID_LEN) {
return Err(at(&format!(
"event_id must be 1..={MAX_EVENT_ID_LEN} printable characters"
)));
}
let tag = format!("item {index} (event_id {}): ", input.event_id);
let at = |detail: &str| format!("{tag}{detail}");
if !is_token(&input.agent, MAX_AGENT_LEN) {
return Err(at(&format!(
"agent must be 1..={MAX_AGENT_LEN} token characters"
)));
}
if !is_token(&input.event, MAX_EVENT_LEN) {
return Err(at(&format!(
"event must be 1..={MAX_EVENT_LEN} token characters"
)));
}
let serde_json::Value::Object(body) = &input.body else {
return Err(at("body must be a JSON object"));
};
let Some(serde_json::Value::String(session_id)) = body.get("session_id") else {
return Err(at("body.session_id must be an explicit string"));
};
if !is_plain(session_id, MAX_SESSION_ID_LEN) {
return Err(at(&format!(
"body.session_id must be 1..={MAX_SESSION_ID_LEN} printable characters"
)));
}
let Some(serde_json::Value::String(cwd)) = body.get("cwd") else {
return Err(at("body.cwd must be an explicit string"));
};
if !is_plain(cwd, MAX_CWD_LEN) {
return Err(at(&format!(
"body.cwd must be 1..={MAX_CWD_LEN} printable characters"
)));
}
let (session_id, cwd) = (session_id.clone(), cwd.clone());
let body_json = serde_json::to_string(&input.body).map_err(|_| at("body is not encodable"))?;
if body_json.len() > MAX_BODY_BYTES {
return Err(at(&format!(
"body is {} bytes, over the {MAX_BODY_BYTES}-byte limit",
body_json.len()
)));
}
let digest = Sha256::digest(body_json.as_bytes());
let body_sha256 = digest.iter().map(|b| format!("{b:02x}")).collect();
let ingest_key = ingest_key(
producer,
actor,
&input.agent,
&session_id,
&input.event,
&input.event_id,
);
Ok(ValidEvent {
event_id: input.event_id,
agent: input.agent,
event: input.event,
session_id,
cwd,
body_json,
body_sha256,
ingest_key,
})
}
+28
View File
@@ -0,0 +1,28 @@
//! Persistent queue for external lifecycle events sent through `POST /hook/batch`.
//!
//! The queue is separate from ai-memory's database and wiki. Events leave it
//! after a validated acknowledgement, which may include a server policy drop.
//! Pending events remain on disk when delivery fails or their retry window
//! expires; reaching a capacity limit rejects new input.
//!
//! Bodies retain their input values before server sanitization. Producers must
//! apply capture exclusions before enqueueing and protect the queue directory.
pub mod ack;
pub mod fsguard;
pub mod http;
pub mod identity;
pub mod queue;
pub mod relay;
/// Wall-clock milliseconds since the Unix epoch.
///
/// Used for the retry window, which is why it is clamped at 0 rather than
/// panicking on a pre-epoch clock: a nonsensical clock must not take the process
/// down mid-flush.
pub fn now_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| i64::try_from(d.as_millis()).unwrap_or(i64::MAX))
.unwrap_or(0)
}
+113
View File
@@ -0,0 +1,113 @@
//! Thin CLI over [`ai_memory_relay::relay`]. Parse, call, print, exit.
use std::path::PathBuf;
use ai_memory_relay::relay::{self, FlushOptions, Report};
use clap::{Parser, Subcommand};
/// Queue external lifecycle events for ai-memory.
///
/// Exit codes: 0 success, 2 failure, 3 flush ended with events pending.
#[derive(Debug, Parser)]
#[command(name = "ai-memory-relay", version, about, long_about = None)]
struct Cli {
#[command(subcommand)]
command: Command,
}
#[derive(Debug, Subcommand)]
enum Command {
/// Bind a queue directory to one destination, producer, actor and scope.
Init {
#[arg(long)]
queue_dir: PathBuf,
/// http(s) base URL of the ai-memory server. No userinfo, query or fragment.
#[arg(long)]
server_url: String,
/// Producer namespace, sent as `extension` (also your `AI_MEMORY_CAPTURE_OWNER`).
#[arg(long)]
producer: String,
/// Stable adapter-side operator namespace. Never a bearer token.
#[arg(long)]
actor: String,
#[arg(long)]
workspace: String,
#[arg(long)]
project: String,
},
/// Add events from a JSON array file. All-or-nothing.
Enqueue {
#[arg(long)]
queue_dir: PathBuf,
/// `[{"event_id","agent","event","body"}]`; body needs explicit `session_id` and `cwd`.
#[arg(long)]
file: PathBuf,
},
/// Deliver pending events through POST /hook/batch.
Flush {
#[arg(long)]
queue_dir: PathBuf,
/// Batches attempted in one flush. Finite by default: a flush never loops
/// forever, and what it does not deliver stays queued for the next run.
#[arg(long, default_value_t = 64)]
max_batches: usize,
},
/// Payload-free counters for the queue, as JSON.
Status {
#[arg(long)]
queue_dir: PathBuf,
},
}
fn main() {
let cli = Cli::parse();
let outcome = match cli.command {
Command::Init {
queue_dir,
server_url,
producer,
actor,
workspace,
project,
} => relay::init(
&queue_dir,
&server_url,
&producer,
&actor,
&workspace,
&project,
),
Command::Enqueue { queue_dir, file } => relay::enqueue(&queue_dir, &file),
Command::Flush {
queue_dir,
max_batches,
} => relay::flush(
&queue_dir,
&FlushOptions {
max_batches: max_batches.max(1),
},
),
Command::Status { queue_dir } => relay::status(&queue_dir),
};
std::process::exit(finish(outcome));
}
fn finish(outcome: anyhow::Result<Report>) -> i32 {
match outcome {
Ok(report) => {
for line in &report.summary {
println!("{line}");
}
if let Some(failure) = &report.failure {
eprintln!("error: {failure}");
}
report.exit_code()
}
Err(error) => {
// `{error:#}` prints the context chain, which is built from paths,
// counts and classes only.
eprintln!("error: {error:#}");
2
}
}
}
+727
View File
@@ -0,0 +1,727 @@
//! The companion's own durable queue: one private SQLite database per queue
//! directory. It never opens ai-memory's database or wiki.
//!
//! The binding fixes the destination, producer identity and scope. Pending
//! events keep their body values until acknowledgement. Receipts retain keys
//! and content hashes for duplicate recognition during the retry window.
//! Session records reject a native session changing agents within one queue.
use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::time::Duration;
use anyhow::{Context, Result, bail};
use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
use crate::fsguard;
use crate::identity::{MAX_BODY_BYTES, RETRY_WINDOW_MS, TERMINAL_EVENTS, ValidEvent};
/// Undelivered events allowed in one queue.
pub const MAX_PENDING_ITEMS: i64 = 50_000;
/// Undelivered body bytes allowed in one queue.
pub const MAX_PENDING_BYTES: i64 = 64 * 1024 * 1024;
/// Pending events plus retained receipts. Admission reserves one receipt for
/// each new event so acknowledgement cannot exceed the limit.
pub const MAX_RECEIPTS: i64 = 200_000;
/// Session-to-agent pins retained at once, enforced the same way.
pub const MAX_SESSION_PINS: i64 = 50_000;
const SCHEMA_VERSION: &str = "1";
const IDENTITY: &str = "ai-memory-relay-queue";
const SCHEMA: &str = "
CREATE TABLE IF NOT EXISTS meta(
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS binding(
id INTEGER PRIMARY KEY CHECK(id = 1),
server_url TEXT NOT NULL,
producer TEXT NOT NULL,
actor TEXT NOT NULL,
workspace TEXT NOT NULL,
project TEXT NOT NULL,
created_at_ms INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS pending(
seq INTEGER PRIMARY KEY AUTOINCREMENT,
ingest_key TEXT NOT NULL UNIQUE,
event_id TEXT NOT NULL,
agent TEXT NOT NULL,
event TEXT NOT NULL,
session_id TEXT NOT NULL,
cwd TEXT NOT NULL,
body_json TEXT NOT NULL,
body_sha256 TEXT NOT NULL,
body_bytes INTEGER NOT NULL,
first_seen_ms INTEGER NOT NULL,
first_attempt_ms INTEGER,
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT
);
CREATE INDEX IF NOT EXISTS pending_session ON pending(agent, session_id, seq);
CREATE TABLE IF NOT EXISTS receipt(
ingest_key TEXT PRIMARY KEY,
body_sha256 TEXT NOT NULL,
first_attempt_ms INTEGER NOT NULL,
delivered_at_ms INTEGER NOT NULL
);
CREATE INDEX IF NOT EXISTS receipt_first_attempt ON receipt(first_attempt_ms);
CREATE TABLE IF NOT EXISTS session_agent(
session_id TEXT PRIMARY KEY,
agent TEXT NOT NULL,
last_seen_ms INTEGER NOT NULL
);
";
/// What `init` bound this queue directory to.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Binding {
pub server_url: String,
pub producer: String,
pub actor: String,
pub workspace: String,
pub project: String,
}
/// Whether `init` created the binding or recognized an identical one.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum BindOutcome {
Created,
Recognized,
}
/// Per-file result of `enqueue`.
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct EnqueueReport {
pub accepted: usize,
/// Same id, byte-identical body: already pending or already delivered.
pub recognized: usize,
}
/// One item selected for delivery.
#[derive(Debug, Clone)]
pub struct PendingItem {
pub seq: i64,
pub ingest_key: String,
pub event_id: String,
pub agent: String,
pub event: String,
pub session_id: String,
pub body_json: String,
pub body_bytes: i64,
pub first_seen_ms: i64,
pub first_attempt_ms: Option<i64>,
}
impl PendingItem {
/// Session coordinate: one head per `(agent, session_id)` per batch.
pub fn session(&self) -> (String, String) {
(self.agent.clone(), self.session_id.clone())
}
}
/// A batch plus what was held back while building it.
#[derive(Debug, Default, Clone)]
pub struct Batch {
pub items: Vec<PendingItem>,
/// Sessions whose head is past the retry window. The whole session is held:
/// skipping its head and sending the next event would reorder it.
pub blocked_sessions: usize,
/// Sessions skipped because the byte budget filled first.
pub budget_deferred: usize,
}
/// Counters for `status`. Deliberately payload-free.
#[derive(Debug, Default, Clone, PartialEq, Eq)]
pub struct Stats {
pub pending_items: i64,
pub pending_bytes: i64,
pub pending_sessions: i64,
pub attempted_items: i64,
pub expired_items: i64,
pub oldest_pending_age_ms: i64,
pub max_attempts: i64,
pub receipts: i64,
pub known_sessions: i64,
}
/// The companion's queue database.
pub struct Queue {
conn: Connection,
path: PathBuf,
}
impl std::fmt::Debug for Queue {
// Path only: never connection internals, never row data.
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Queue").field("path", &self.path).finish()
}
}
impl Queue {
/// Open (or create) the queue database inside an already-prepared directory.
///
/// Check an existing file's identity before schema writes. Schema and
/// metadata commit together, so an interrupted initialization can reopen
/// the empty database left by rollback.
pub fn open(dir: &Path) -> Result<Self> {
let path = dir.join(fsguard::DB_FILE);
fsguard::check_queue_file(&path)?;
for name in fsguard::SIDECARS {
fsguard::check_queue_file(&dir.join(name))?;
}
if !path.exists() {
// Own the mode before SQLite ever opens the file: the `-wal` and
// `-shm` sidecars inherit the database's permissions, so creating it
// under the ambient umask would make them world-readable.
fsguard::create_private_file(&path)?;
fsguard::check_queue_file(&path)?;
}
let mut conn = Connection::open(&path)
.with_context(|| format!("open relay queue at {}", path.display()))?;
conn.busy_timeout(Duration::from_secs(10))?;
let state = classify(&conn, &path)?;
conn.execute_batch(
"PRAGMA journal_mode=WAL;\nPRAGMA synchronous=FULL;\nPRAGMA foreign_keys=ON;",
)
.with_context(|| format!("configure relay queue at {}", path.display()))?;
if state != DbState::Ready {
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
tx.execute_batch(SCHEMA)
.with_context(|| format!("initialize relay queue schema at {}", path.display()))?;
tx.execute(
"INSERT OR REPLACE INTO meta(key, value) VALUES('identity', ?1), ('schema_version', ?2)",
params![IDENTITY, SCHEMA_VERSION],
)?;
tx.commit()?;
}
Ok(Self { conn, path })
}
/// Path of the database, for operator-facing messages.
pub fn path(&self) -> &Path {
&self.path
}
/// Bind the queue to a destination, producer, actor and scope.
///
/// Repeated initialization must preserve the destination and identity of
/// queued events, including events whose acknowledgement was lost.
pub fn bind(&mut self, binding: &Binding, now_ms: i64) -> Result<BindOutcome> {
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
let outcome = match read_binding(&tx)? {
Some(existing) if existing == *binding => BindOutcome::Recognized,
Some(existing) => {
let mut changed = Vec::new();
let mut note = |label: &str, old: &str, new: &str| {
if old != new {
changed.push(format!("{label}: {old:?} -> {new:?}"));
}
};
note("server-url", &existing.server_url, &binding.server_url);
note("producer", &existing.producer, &binding.producer);
note("actor", &existing.actor, &binding.actor);
note("workspace", &existing.workspace, &binding.workspace);
note("project", &existing.project, &binding.project);
bail!(
"binding mismatch: this queue is already bound ({}). \
Use a separate --queue-dir for a different destination or identity",
changed.join(", ")
);
}
None => {
tx.execute(
"INSERT INTO binding(id, server_url, producer, actor, workspace, project, created_at_ms)
VALUES(1, ?1, ?2, ?3, ?4, ?5, ?6)",
params![
binding.server_url,
binding.producer,
binding.actor,
binding.workspace,
binding.project,
now_ms,
],
)?;
BindOutcome::Created
}
};
tx.commit()?;
Ok(outcome)
}
/// The binding, or a clear error telling the operator to run `init` first.
pub fn binding(&self) -> Result<Binding> {
read_binding(&self.conn)?.ok_or_else(|| {
anyhow::anyhow!("queue is not bound yet; run `ai-memory-relay init` first")
})
}
/// Persist a whole validated file, all-or-nothing.
///
/// `first_seen_ms` is stamped here for reporting. It is *not* the retry
/// clock: that starts at the first durable attempt, recorded before the
/// first HTTP request (see [`Queue::stamp_attempt`]).
pub fn enqueue(&mut self, events: &[ValidEvent], now_ms: i64) -> Result<EnqueueReport> {
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
let (mut items, mut bytes): (i64, i64) = tx.query_row(
"SELECT COUNT(*), COALESCE(SUM(body_bytes), 0) FROM pending",
[],
|r| Ok((r.get(0)?, r.get(1)?)),
)?;
// Replay protection is bounded, and the bound covers pending *and*
// receipts: every pending item becomes a receipt when it is
// acknowledged, so counting receipts alone would let the queue commit to
// more protection than the cap allows. Pruning stays reserved for
// entries past the 30-day window, so a full queue refuses new ids rather
// than forgetting a recently delivered one. A duplicate is still
// recognized while full, because recognizing it adds nothing.
let receipts: i64 = tx.query_row("SELECT COUNT(*) FROM receipt", [], |r| r.get(0))?;
let mut pins: i64 = tx.query_row("SELECT COUNT(*) FROM session_agent", [], |r| r.get(0))?;
let mut report = EnqueueReport::default();
let mut seen_keys: HashMap<&str, &str> = HashMap::new();
let mut seen_sessions: HashMap<&str, &str> = HashMap::new();
for event in events {
// A session id belongs to exactly one wire agent. Two agents on one
// native session is the producer bug the server can only answer with
// an acknowledged SessionCollision drop.
let bound: Option<String> = tx
.query_row(
"SELECT agent FROM session_agent WHERE session_id = ?1",
params![event.session_id],
|r| r.get(0),
)
.optional()?;
let bound = bound.or_else(|| {
seen_sessions
.get(event.session_id.as_str())
.map(|a| (*a).to_owned())
});
if bound.is_none() {
pins += 1;
if pins > MAX_SESSION_PINS {
bail!(
"session pin table is full: {MAX_SESSION_PINS} session(s) already pinned \
to an agent. Nothing was enqueued and no pin was forgotten; pins are \
released as they pass the 30-day retry window"
);
}
}
if let Some(bound) = bound
&& bound != event.agent
{
bail!(
"session identity conflict: session {:?} is already bound to agent {:?} in \
this queue, but item with event_id {} claims agent {:?}. Nothing was \
enqueued; a native session belongs to one harness",
event.session_id,
bound,
event.event_id,
event.agent
);
}
seen_sessions.insert(&event.session_id, &event.agent);
if let Some(previous) = seen_keys.insert(&event.ingest_key, &event.body_sha256)
&& previous != event.body_sha256
{
bail!(collision_message(
&event.ingest_key,
"another item in this file"
));
}
let known: Option<String> = tx
.query_row(
"SELECT body_sha256 FROM pending WHERE ingest_key = ?1
UNION ALL
SELECT body_sha256 FROM receipt WHERE ingest_key = ?1",
params![event.ingest_key],
|r| r.get(0),
)
.optional()?;
if let Some(known) = known {
if known != event.body_sha256 {
bail!(collision_message(
&event.ingest_key,
"an item already recorded in this queue"
));
}
report.recognized += 1;
continue;
}
let body_bytes = event.body_json.len() as i64;
items += 1;
bytes += body_bytes;
if items > MAX_PENDING_ITEMS || bytes > MAX_PENDING_BYTES {
bail!(
"queue is full: {items} items / {bytes} bytes would exceed the \
{MAX_PENDING_ITEMS}-item / {MAX_PENDING_BYTES}-byte limit. Nothing was \
enqueued and nothing was dropped; flush the queue first"
);
}
// Reserve the receipt this item will need once it is acknowledged.
let protected = receipts + items;
if protected > MAX_RECEIPTS {
bail!(
"replay protection is full: {receipts} acknowledged key(s) plus {items} \
pending would need {protected} slots, over the {MAX_RECEIPTS} limit. \
Nothing was enqueued and nothing was forgotten; slots are released as \
entries pass the 30-day retry window, or start a fresh --queue-dir"
);
}
tx.execute(
"INSERT INTO pending(
ingest_key, event_id, agent, event, session_id, cwd,
body_json, body_sha256, body_bytes, first_seen_ms)
VALUES(?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
params![
event.ingest_key,
event.event_id,
event.agent,
event.event,
event.session_id,
event.cwd,
event.body_json,
event.body_sha256,
body_bytes,
now_ms,
],
)?;
report.accepted += 1;
}
for (session_id, agent) in seen_sessions {
tx.execute(
"INSERT INTO session_agent(session_id, agent, last_seen_ms) VALUES(?1, ?2, ?3)
ON CONFLICT(session_id) DO UPDATE SET last_seen_ms = excluded.last_seen_ms",
params![session_id, agent, now_ms],
)?;
}
tx.commit()?;
Ok(report)
}
/// Pick at most one pending head per `(agent, session_id)`, oldest first,
/// within both an item count and a wire-byte budget.
///
/// One head per session per batch is what makes a partial acknowledgement
/// safe: two events of one session are never in flight together, so no ack
/// subset can commit event 2 while event 1 is still pending, and a terminal
/// event can never overtake an earlier one. `deferred` holds sessions this
/// flush already stopped on; they are skipped without touching their order.
pub fn select_batch(
&self,
limit: usize,
byte_budget: usize,
now_ms: i64,
deferred: &HashSet<(String, String)>,
cost_of: impl Fn(&PendingItem) -> usize,
) -> Result<Batch> {
let limit = limit.clamp(1, crate::identity::MAX_BATCH_ITEMS);
let mut stmt = self.conn.prepare(
"SELECT seq, ingest_key, event_id, agent, event, session_id, body_json,
body_bytes, first_seen_ms, first_attempt_ms
FROM pending ORDER BY seq ASC",
)?;
let mut rows = stmt.query([])?;
let mut seen: HashSet<(String, String)> = HashSet::new();
let mut batch = Batch::default();
let mut spent = 0usize;
while let Some(row) = rows.next()? {
if batch.items.len() >= limit {
break;
}
let item = PendingItem {
seq: row.get(0)?,
ingest_key: row.get(1)?,
event_id: row.get(2)?,
agent: row.get(3)?,
event: row.get(4)?,
session_id: row.get(5)?,
body_json: row.get(6)?,
body_bytes: row.get(7)?,
first_seen_ms: row.get(8)?,
first_attempt_ms: row.get(9)?,
};
let session = item.session();
if !seen.insert(session.clone()) {
continue;
}
if deferred.contains(&session) {
continue;
}
if is_expired(item.first_attempt_ms, now_ms) {
batch.blocked_sessions += 1;
continue;
}
let cost = cost_of(&item);
if !batch.items.is_empty() && spent + cost > byte_budget {
batch.budget_deferred += 1;
continue;
}
debug_assert!(
!TERMINAL_EVENTS.contains(&item.event.as_str()) || self.session_head(&item)?,
"a terminal event was selected while its session had an older pending item"
);
spent += cost;
batch.items.push(item);
}
Ok(batch)
}
fn session_head(&self, item: &PendingItem) -> Result<bool> {
let older: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM pending WHERE agent = ?1 AND session_id = ?2 AND seq < ?3",
params![item.agent, item.session_id, item.seq],
|r| r.get(0),
)?;
Ok(older == 0)
}
/// Record the durable start of the retry window, **before** the request.
///
/// `COALESCE` is the whole point: a restart, a lost ack, or a later retry
/// never moves the clock forward, so the 30-day guard measures from the
/// first time this event was actually put on the wire.
pub fn stamp_attempt(&mut self, keys: &[String], now_ms: i64) -> Result<()> {
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
for key in keys {
tx.execute(
"UPDATE pending
SET first_attempt_ms = COALESCE(first_attempt_ms, ?2),
attempts = attempts + 1
WHERE ingest_key = ?1",
params![key, now_ms],
)?;
}
tx.commit()?;
Ok(())
}
/// Move acknowledged items to compact receipts, in one transaction.
///
/// An acknowledgement means delivered *or* deliberately dropped by server
/// policy (a capture-protocol drop, a subagent drop, a session-collision
/// drop). It never promises an observation was written.
pub fn confirm(&mut self, keys: &[String], now_ms: i64) -> Result<usize> {
if keys.is_empty() {
return Ok(0);
}
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
let mut moved = 0;
for key in keys {
let row: Option<(String, Option<i64>, i64)> = tx
.query_row(
"SELECT body_sha256, first_attempt_ms, first_seen_ms
FROM pending WHERE ingest_key = ?1",
params![key],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()?;
let Some((body_sha256, first_attempt_ms, first_seen_ms)) = row else {
continue;
};
tx.execute(
"INSERT INTO receipt(ingest_key, body_sha256, first_attempt_ms, delivered_at_ms)
VALUES(?1, ?2, ?3, ?4)
ON CONFLICT(ingest_key) DO UPDATE SET delivered_at_ms = excluded.delivered_at_ms",
params![
key,
body_sha256,
first_attempt_ms.unwrap_or(first_seen_ms),
now_ms
],
)?;
tx.execute("DELETE FROM pending WHERE ingest_key = ?1", params![key])?;
moved += 1;
}
tx.commit()?;
Ok(moved)
}
/// Charge one item with a delivery failure. The class is a short label,
/// never a server body and never payload.
pub fn record_failure(&mut self, key: &str, class: &str) -> Result<()> {
let class: String = class.chars().filter(|c| !c.is_control()).take(80).collect();
self.conn.execute(
"UPDATE pending SET last_error = ?2 WHERE ingest_key = ?1",
params![key, class],
)?;
Ok(())
}
/// Release expired receipts and session records without pending events.
pub fn prune(&mut self, now_ms: i64) -> Result<usize> {
let cutoff = now_ms.saturating_sub(RETRY_WINDOW_MS);
let tx = self
.conn
.transaction_with_behavior(TransactionBehavior::Immediate)?;
let receipts = tx.execute(
"DELETE FROM receipt WHERE first_attempt_ms < ?1",
params![cutoff],
)?;
let sessions = tx.execute(
"DELETE FROM session_agent WHERE last_seen_ms < ?1
AND session_id NOT IN (SELECT session_id FROM pending)",
params![cutoff],
)?;
tx.commit()?;
Ok(receipts + sessions)
}
/// Payload-free counters.
pub fn stats(&self, now_ms: i64) -> Result<Stats> {
let cutoff = now_ms.saturating_sub(RETRY_WINDOW_MS);
let (pending_items, pending_bytes, oldest, max_attempts, attempted_items) =
self.conn.query_row(
"SELECT COUNT(*), COALESCE(SUM(body_bytes), 0), COALESCE(MIN(first_seen_ms), 0),
COALESCE(MAX(attempts), 0), COUNT(first_attempt_ms)
FROM pending",
[],
|r| {
Ok((
r.get::<_, i64>(0)?,
r.get::<_, i64>(1)?,
r.get::<_, i64>(2)?,
r.get::<_, i64>(3)?,
r.get::<_, i64>(4)?,
))
},
)?;
let pending_sessions: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM (SELECT DISTINCT agent, session_id FROM pending)",
[],
|r| r.get(0),
)?;
let expired_items: i64 = self.conn.query_row(
"SELECT COUNT(*) FROM pending WHERE first_attempt_ms IS NOT NULL
AND first_attempt_ms <= ?1",
params![cutoff],
|r| r.get(0),
)?;
let receipts: i64 = self
.conn
.query_row("SELECT COUNT(*) FROM receipt", [], |r| r.get(0))?;
let known_sessions: i64 =
self.conn
.query_row("SELECT COUNT(*) FROM session_agent", [], |r| r.get(0))?;
Ok(Stats {
pending_items,
pending_bytes,
pending_sessions,
attempted_items,
expired_items,
oldest_pending_age_ms: if pending_items == 0 {
0
} else {
now_ms.saturating_sub(oldest)
},
max_attempts,
receipts,
known_sessions,
})
}
}
/// Past the retry window an event may not be auto-resent.
///
/// Refuse the boundary itself. This comparison depends on an accurate host
/// clock; a backwards clock change can extend retries past server key expiry.
pub fn is_expired(first_attempt_ms: Option<i64>, now_ms: i64) -> bool {
match first_attempt_ms {
Some(started) => now_ms.saturating_sub(started) >= RETRY_WINDOW_MS,
None => false,
}
}
fn collision_message(key: &str, other: &str) -> String {
format!(
"identity collision: ingest key {key} was already used by {other} with different content. \
The same producer event id must always carry the same body; nothing was enqueued. \
Body content is deliberately not shown. Per-body limit: {MAX_BODY_BYTES} bytes"
)
}
/// What an existing file turned out to be.
///
/// Schema and identity commit together. Recovery therefore accepts an empty
/// database or an identified queue; generic table names do not establish ownership.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DbState {
/// A SQLite file with no user objects: empty, or rolled back to empty.
Fresh,
/// A relay queue of this exact schema version.
Ready,
}
/// Classify an existing file before touching it. Never mutates.
fn classify(conn: &Connection, path: &Path) -> Result<DbState> {
let refuse = |detail: String| -> anyhow::Error {
anyhow::anyhow!(
"{} is not a usable relay queue ({detail}); nothing was modified. \
Point --queue-dir at an empty directory or restore the original queue",
path.display()
)
};
// Tables, views, triggers and standalone indexes all count: any user object
// means the file belongs to something, and that something is not this relay
// unless it also carries the identity row.
let objects: Result<i64, rusqlite::Error> = conn.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE name NOT LIKE 'sqlite_%'",
[],
|r| r.get(0),
);
// A file that is not a SQLite database fails right here, before any write.
let objects = objects.map_err(|e| refuse(format!("cannot read its object list: {e}")))?;
if objects == 0 {
return Ok(DbState::Fresh);
}
let rows: Result<Vec<(String, String)>, rusqlite::Error> = (|| {
let mut stmt =
conn.prepare("SELECT key, value FROM meta WHERE key IN ('identity','schema_version')")?;
let mapped = stmt.query_map([], |r| Ok((r.get(0)?, r.get(1)?)))?;
mapped.collect()
})();
let meta: HashMap<String, String> = rows
.map_err(|e| refuse(format!("it carries no readable relay metadata: {e}")))?
.into_iter()
.collect();
match (
meta.get("identity").map(String::as_str),
meta.get("schema_version").map(String::as_str),
) {
(Some(IDENTITY), Some(SCHEMA_VERSION)) => Ok(DbState::Ready),
(Some(IDENTITY), Some(other)) => Err(refuse(format!(
"schema version {other} but this build speaks {SCHEMA_VERSION}"
))),
(Some(other), _) => Err(refuse(format!("identity is {other:?}"))),
_ => Err(refuse(
"it holds user objects but no relay identity row".into(),
)),
}
}
fn read_binding(conn: &Connection) -> Result<Option<Binding>> {
Ok(conn
.query_row(
"SELECT server_url, producer, actor, workspace, project FROM binding WHERE id = 1",
[],
|r| {
Ok(Binding {
server_url: r.get(0)?,
producer: r.get(1)?,
actor: r.get(2)?,
workspace: r.get(3)?,
project: r.get(4)?,
})
},
)
.optional()?)
}
+556
View File
@@ -0,0 +1,556 @@
//! CLI command operations and reports.
//!
//! Reports include queue metadata and counts. Event bodies and bearer tokens
//! are kept out of diagnostics.
use std::collections::HashSet;
use std::io::Read;
use std::path::Path;
use std::time::{Duration, Instant};
use anyhow::{Context, Result, bail};
use crate::ack::{self, BatchAck};
use crate::fsguard;
use crate::http::{self, BatchItem, Sender, Token};
use crate::identity::{self, InputEvent, MAX_BATCH_ITEMS, ValidEvent};
use crate::queue::{self, Batch, BindOutcome, Binding, PendingItem, Queue};
/// Largest `enqueue --file` input accepted, before parsing.
pub const MAX_INPUT_BYTES: u64 = 32 * 1024 * 1024;
/// Most events accepted from one input file.
pub const MAX_INPUT_ITEMS: usize = 10_000;
/// Wire-byte budget for one batch. The server's `/hook` body limit is 10 MiB
/// (`serve.rs`); core's own spool drain uses 8 MiB, and so does this.
pub const MAX_BATCH_BYTES: usize = 8 * 1024 * 1024;
/// JSON framing charged per item on top of its URL and body: `{"url":"","body":},`.
const ITEM_FRAMING_BYTES: usize = 32;
/// Transport retries per batch, on top of the first attempt.
const TRANSPORT_RETRIES: usize = 2;
/// Request timeout budget per item; a batch gets this times its item count.
const PER_EVENT_TIMEOUT: Duration = Duration::from_secs(2);
/// Absolute ceiling for one batch request, however many items it carries.
const MAX_BATCH_TIMEOUT: Duration = Duration::from_secs(120);
/// Wall-clock ceiling for one whole flush.
const TOTAL_BUDGET: Duration = Duration::from_secs(900);
/// Ceiling for one backoff sleep between transport retries.
const MAX_BACKOFF: Duration = Duration::from_secs(2);
/// TCP connect timeout.
const CONNECT_TIMEOUT: Duration = Duration::from_secs(10);
/// The one knob a `flush` exposes. Everything else is a constant above, so the
/// public surface cannot be tuned into an unbounded run.
#[derive(Debug, Clone)]
pub struct FlushOptions {
/// Batches attempted in one flush. Finite by construction.
pub max_batches: usize,
}
impl Default for FlushOptions {
fn default() -> Self {
Self { max_batches: 64 }
}
}
/// What a command wants the process to exit with.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Report {
pub summary: Vec<String>,
/// Work remains: unacknowledged items are still queued.
pub pending: bool,
/// Something went wrong. The queue is intact either way.
pub failure: Option<String>,
}
impl Report {
fn ok(summary: Vec<String>) -> Self {
Self {
summary,
pending: false,
failure: None,
}
}
/// 0 = drained, 2 = failure, 3 = nothing lost but work remains.
pub fn exit_code(&self) -> i32 {
if self.failure.is_some() {
2
} else if self.pending {
3
} else {
0
}
}
}
/// `init`: bind a queue directory to one destination, producer, actor and scope.
pub fn init(
dir: &Path,
server_url: &str,
producer: &str,
actor: &str,
workspace: &str,
project: &str,
) -> Result<Report> {
identity::check_producer(producer).map_err(anyhow::Error::msg)?;
identity::check_actor(actor).map_err(anyhow::Error::msg)?;
identity::check_scope("workspace", workspace).map_err(anyhow::Error::msg)?;
identity::check_scope("project", project).map_err(anyhow::Error::msg)?;
let binding = Binding {
server_url: http::normalize_server_url(server_url)?,
producer: producer.to_owned(),
actor: actor.to_owned(),
workspace: workspace.to_owned(),
project: project.to_owned(),
};
let dir = fsguard::prepare_queue_dir(dir)?;
let mut queue = Queue::open(&dir)?;
let outcome = queue.bind(&binding, crate::now_ms())?;
let verb = match outcome {
BindOutcome::Created => "bound",
BindOutcome::Recognized => "already bound (identical)",
};
Ok(Report::ok(vec![
format!("queue {verb}: {}", queue.path().display()),
format!("server: {}", binding.server_url),
format!("producer: {}", binding.producer),
format!("actor: {}", binding.actor),
format!("scope: {}/{}", binding.workspace, binding.project),
]))
}
/// `enqueue`: validate a whole file, then persist it all-or-nothing.
pub fn enqueue(dir: &Path, file: &Path) -> Result<Report> {
let dir = fsguard::prepare_queue_dir(dir)?;
let mut queue = Queue::open(&dir)?;
let binding = queue.binding()?;
let events = read_input(file, &binding)?;
let report = queue.enqueue(&events, crate::now_ms())?;
let stats = queue.stats(crate::now_ms())?;
Ok(Report::ok(vec![
format!(
"enqueued {} event(s), recognized {} duplicate(s)",
report.accepted, report.recognized
),
format!(
"pending now: {} event(s) across {} session(s), {} byte(s)",
stats.pending_items, stats.pending_sessions, stats.pending_bytes
),
]))
}
/// Read and validate an input file without ever quoting its content back.
fn read_input(file: &Path, binding: &Binding) -> Result<Vec<ValidEvent>> {
let meta =
std::fs::metadata(file).with_context(|| format!("read --file {}", file.display()))?;
if meta.len() > MAX_INPUT_BYTES {
bail!(
"--file {} is {} bytes, over the {MAX_INPUT_BYTES}-byte input limit",
file.display(),
meta.len()
);
}
let mut raw = Vec::new();
std::fs::File::open(file)
.with_context(|| format!("open --file {}", file.display()))?
.take(MAX_INPUT_BYTES + 1)
.read_to_end(&mut raw)
.with_context(|| format!("read --file {}", file.display()))?;
// The metadata check above is an early out; this is the one that counts,
// because the file can grow between `stat` and `read`.
if raw.len() as u64 > MAX_INPUT_BYTES {
bail!(
"--file {} is over the {MAX_INPUT_BYTES}-byte input limit",
file.display()
);
}
let input: Vec<InputEvent> = serde_json::from_slice(&raw).map_err(|e| {
// Category plus position only. `Display` on a serde error can quote the
// offending value ("invalid type: string \"...\""), and an event body is
// exactly the thing that must not reach a terminal or a log.
anyhow::anyhow!(
"--file {} is not a valid event array ({} error at line {}, column {}): expected \
[{{\"event_id\",\"agent\",\"event\",\"body\"}}, ...]",
file.display(),
match e.classify() {
serde_json::error::Category::Io => "io",
serde_json::error::Category::Syntax => "syntax",
serde_json::error::Category::Data => "schema",
serde_json::error::Category::Eof => "truncated input",
},
e.line(),
e.column()
)
})?;
if input.is_empty() {
bail!("--file {} contains no events", file.display());
}
if input.len() > MAX_INPUT_ITEMS {
bail!(
"--file {} carries {} events, over the {MAX_INPUT_ITEMS}-event input limit",
file.display(),
input.len()
);
}
input
.into_iter()
.enumerate()
.map(|(index, event)| {
identity::validate(index, event, &binding.producer, &binding.actor)
.map_err(anyhow::Error::msg)
})
.collect()
}
/// `status`: payload-free counters for the queue, as one stable JSON object.
///
/// Machine-readable on purpose: an operator script (and the end-to-end smoke
/// test) reads `pending_items` from it. Every field here is a count, a name the
/// operator chose, or a path. No event body, no cwd, no token.
pub fn status(dir: &Path) -> Result<Report> {
let dir = fsguard::prepare_queue_dir(dir)?;
let queue = Queue::open(&dir)?;
let binding = queue.binding()?;
let stats = queue.stats(crate::now_ms())?;
let document = serde_json::json!({
"queue": queue.path().display().to_string(),
"server_url": binding.server_url,
"producer": binding.producer,
"actor": binding.actor,
"workspace": binding.workspace,
"project": binding.project,
"pending_items": stats.pending_items,
"pending_bytes": stats.pending_bytes,
"pending_sessions": stats.pending_sessions,
"attempted_items": stats.attempted_items,
"expired_items": stats.expired_items,
"oldest_pending_age_ms": stats.oldest_pending_age_ms,
"max_attempts": stats.max_attempts,
"receipts": stats.receipts,
"known_sessions": stats.known_sessions,
"limits": {
"pending_items": queue::MAX_PENDING_ITEMS,
"pending_bytes": queue::MAX_PENDING_BYTES,
"receipts": queue::MAX_RECEIPTS,
"known_sessions": queue::MAX_SESSION_PINS,
"body_bytes": crate::identity::MAX_BODY_BYTES,
"batch_items": MAX_BATCH_ITEMS,
"batch_bytes": MAX_BATCH_BYTES,
},
"retry_window_ms": crate::identity::RETRY_WINDOW_MS,
"ack_means": "delivered or dropped by server policy, not that an observation was written",
});
Ok(Report::ok(vec![serde_json::to_string_pretty(&document)?]))
}
/// `flush`: deliver pending events, oldest first, one head per session per batch.
pub fn flush(dir: &Path, options: &FlushOptions) -> Result<Report> {
let dir = fsguard::prepare_queue_dir(dir)?;
let lock_path = dir.join("flush.lock");
fsguard::check_queue_file(&lock_path)?;
if !lock_path.exists() {
fsguard::create_private_file(&lock_path)?;
fsguard::check_queue_file(&lock_path)?;
}
let lock = std::fs::OpenOptions::new()
.truncate(false)
.write(true)
.open(&lock_path)
.with_context(|| format!("open {}", lock_path.display()))?;
// Serialize flushes across processes. The lock file itself is never removed:
// releasing the lock is dropping the handle.
if fs2::FileExt::try_lock_exclusive(&lock).is_err() {
bail!(
"another flush already holds {}; run one flush at a time",
lock_path.display()
);
}
let result = flush_locked(&dir, options);
let _ = fs2::FileExt::unlock(&lock);
result
}
fn flush_locked(dir: &Path, options: &FlushOptions) -> Result<Report> {
let mut queue = Queue::open(dir)?;
let binding = queue.binding()?;
queue.prune(crate::now_ms())?;
let token = Token::from_env()?;
let authenticated = token.is_some();
let sender = Sender::new(binding.server_url.clone(), token, CONNECT_TIMEOUT)?;
let started = Instant::now();
let mut deferred: HashSet<(String, String)> = HashSet::new();
let mut delivered = 0usize;
let mut batches = 0usize;
let mut blocked_sessions = 0usize;
let mut failure: Option<String> = None;
let mut backed_off = false;
for _ in 0..options.max_batches {
// Cheap pre-check, so an exhausted budget does not buy another round of
// selection and stamping. The binding one is taken just before the send.
if TOTAL_BUDGET.checked_sub(started.elapsed()).is_none() {
break;
}
let now = crate::now_ms();
let base = binding.server_url.clone();
let for_cost = binding.clone();
let batch: Batch =
queue.select_batch(MAX_BATCH_ITEMS, MAX_BATCH_BYTES, now, &deferred, |item| {
item_cost(&base, &for_cost, item)
})?;
blocked_sessions = batch.blocked_sessions;
if batch.items.is_empty() {
break;
}
batches += 1;
let keys: Vec<String> = batch.items.iter().map(|i| i.ingest_key.clone()).collect();
// Durable, and before the request: a lost ack must not restart the
// 30-day clock on the next run.
queue.stamp_attempt(&keys, now)?;
let items: Vec<BatchItem> = batch
.items
.iter()
.map(|item| {
Ok(BatchItem {
url: http::item_url(&binding.server_url, &binding, item),
body: serde_json::from_str(&item.body_json)
.context("stored body is no longer valid JSON")?,
})
})
.collect::<Result<Vec<_>>>()?;
// Selection checked the retry window once; a retry inside this batch can
// still cross it. Carry the batch's tightest expiry into the sender so
// every attempt is re-checked against it.
//
// Both clocks are read here, after stamping and JSON building: measuring
// from a `now` captured before that preparation would hand the batch
// however long the preparation took as extra allowance.
let send_now = crate::now_ms();
let ttl_left = batch
.items
.iter()
.map(|item| ttl_remaining(item.first_attempt_ms.unwrap_or(now), send_now))
.min()
.unwrap_or(Duration::ZERO);
let Some(remaining) = TOTAL_BUDGET.checked_sub(started.elapsed()) else {
break;
};
let delivery = match send_with_retries(&sender, &items, remaining, ttl_left) {
Ok(delivery) => delivery,
Err(e) => {
// Transport failure: everything stays pending, by definition.
for item in &batch.items {
queue.record_failure(&item.ingest_key, "transport")?;
}
failure = Some(format!("{e} (all {} item(s) kept)", batch.items.len()));
break;
}
};
// An ack is only meaningful on 200 and 429. Every other status keeps the
// whole batch, whatever its body claims.
if !matches!(delivery.status, 200 | 429) {
let detail = match delivery.status {
401 | 403 => {
if authenticated {
"server rejected the bearer token (AI_MEMORY_AUTH_TOKEN)"
} else {
"server requires authentication; set AI_MEMORY_AUTH_TOKEN"
}
}
413 => {
"the server refused the batch as too large. The relay's own batch bounds \
are fixed; record smaller event bodies at the producer, or raise the \
server's body limit"
}
_ => "unexpected status",
};
for item in &batch.items {
queue.record_failure(&item.ingest_key, &format!("http {}", delivery.status))?;
}
failure = Some(format!(
"HTTP {} from /hook/batch: {detail}; all {} item(s) kept pending",
delivery.status,
batch.items.len()
));
break;
}
let parsed: BatchAck = match serde_json::from_slice(&delivery.body) {
Ok(parsed) => parsed,
Err(_) => {
for item in &batch.items {
queue.record_failure(&item.ingest_key, "malformed ack")?;
}
failure = Some(format!(
"HTTP {} from /hook/batch with an unreadable ack ({} byte(s)); \
all {} item(s) kept pending",
delivery.status,
delivery.body.len(),
batch.items.len()
));
break;
}
};
let accepted = match ack::validate(batch.items.len(), &parsed) {
Ok(accepted) => accepted,
Err(rejected) => {
for item in &batch.items {
queue.record_failure(&item.ingest_key, "inconsistent ack")?;
}
failure = Some(format!(
"HTTP {} from /hook/batch with an inconsistent ack ({rejected}); \
all {} item(s) kept pending",
delivery.status,
batch.items.len()
));
break;
}
};
// Only now, after the whole response was validated, does anything leave
// the queue.
let confirmed: Vec<String> = accepted
.iter()
.filter_map(|idx| batch.items.get(*idx))
.map(|item| item.ingest_key.clone())
.collect();
delivered += queue.confirm(&confirmed, crate::now_ms())?;
if let Some(failed_index) = parsed.failed_index
&& let Some(item) = batch.items.get(failed_index)
{
queue.record_failure(&item.ingest_key, "server reported failed_index")?;
}
// The server stops at failed_index. Keep later, untried sessions
// eligible; otherwise one failed head could block them on every flush.
// Defer the failed head and earlier rate-limited items for this flush.
for (position, item) in batch.items.iter().enumerate() {
if accepted.contains(&position) {
continue;
}
let untried = parsed.failed_index.is_some_and(|failed| position > failed);
if !untried {
deferred.insert(item.session());
}
}
if delivery.status == 429 {
backed_off = true;
break;
}
}
let stats = queue.stats(crate::now_ms())?;
let mut summary = vec![format!(
"acknowledged {delivered} event(s) in {batches} batch(es); {} still pending across {} session(s)",
stats.pending_items, stats.pending_sessions
)];
if backed_off {
summary.push(
"server answered 429: the acknowledged items were applied and the flush stopped. \
Back off before the next run"
.to_owned(),
);
}
if !deferred.is_empty() {
summary.push(format!(
"{} session(s) deferred to a later flush (a failed or rate-limited head); \
their later events were never sent ahead of it",
deferred.len()
));
}
if blocked_sessions > 0 {
summary.push(format!(
"{blocked_sessions} session(s) held: their oldest event passed the 30-day retry \
window. The relay will not resend those automatically; they are retained for you"
));
}
summary.push(
"acknowledged means delivered or dropped by server policy, not that an observation \
was written"
.to_owned(),
);
Ok(Report {
summary,
pending: stats.pending_items > 0,
failure,
})
}
/// How long an event may still be auto-resent, measured from its first durable
/// attempt. Zero means the 30-day window is spent.
///
/// Accuracy depends on a stable host clock: the relay carries no time source of
/// its own, so a clock that jumps stretches or shortens this measurement.
pub fn ttl_remaining(first_attempt_ms: i64, now_ms: i64) -> Duration {
let spent = now_ms.saturating_sub(first_attempt_ms);
let left = identity::RETRY_WINDOW_MS.saturating_sub(spent);
Duration::from_millis(u64::try_from(left).unwrap_or(0))
}
/// Retry a batch a finite number of times under one shared deadline.
///
/// Two clocks bound it. `remaining` is the flush budget left when the batch
/// started; `ttl_left` is how long the batch's oldest item may still be resent.
/// Every request *and* every backoff sleep is drawn from the tighter of the two
/// and recomputed, so three attempts can neither overrun the budget nor put an
/// event on the wire after its retry window closed mid-flush.
fn send_with_retries(
sender: &Sender,
items: &[BatchItem],
remaining: Duration,
ttl_left: Duration,
) -> Result<http::Delivery> {
let start = Instant::now();
let deadline = start + remaining.min(ttl_left);
let ttl_deadline = start + ttl_left;
let mut attempt = 0usize;
loop {
if Instant::now() >= ttl_deadline {
bail!(
"hook batch not sent: the batch's oldest event reached the 30-day retry window \
mid-flush; it is retained, not resent"
);
}
let left = deadline.saturating_duration_since(Instant::now());
if left.is_zero() {
bail!("hook batch transport failure: flush budget exhausted");
}
match sender.post_batch(items, batch_timeout(items.len(), left)) {
Ok(delivery) => return Ok(delivery),
Err(e) => {
if attempt >= TRANSPORT_RETRIES {
return Err(e);
}
let backoff = Duration::from_millis(250)
.saturating_mul(1u32 << attempt.min(8))
.min(MAX_BACKOFF)
.min(deadline.saturating_duration_since(Instant::now()));
if backoff.is_zero() {
return Err(e);
}
attempt += 1;
std::thread::sleep(backoff);
}
}
}
}
/// Wire cost of one item, used for the batch byte budget.
pub fn item_cost(base: &str, binding: &Binding, item: &PendingItem) -> usize {
http::item_url(base, binding, item).len() + item.body_json.len() + ITEM_FRAMING_BYTES
}
/// Per-event timeout scaled by item count, capped twice: by an absolute maximum
/// and by whatever is left of the flush budget. Mirrors core's spool drain.
pub fn batch_timeout(items: usize, remaining: Duration) -> Duration {
let items = u32::try_from(items.max(1)).unwrap_or(u32::MAX);
PER_EVENT_TIMEOUT
.checked_mul(items)
.unwrap_or(Duration::MAX)
.min(MAX_BATCH_TIMEOUT)
.min(remaining)
.max(Duration::from_millis(1))
}
File diff suppressed because it is too large Load Diff
+17
View File
@@ -63,6 +63,23 @@ mutation broker. Browsers should talk to the companion; the companion should tal
to ai-memory with an operator token. That keeps CSRF, confirmation, audit, rate
limits, and UI-specific policy outside the core server.
## `ai-memory-relay`: external lifecycle delivery
[`ai-memory-relay`](../companions/ai-memory-relay) delivers events collected by
an external orchestrator through `/hook/batch`. Its own local queue records events
before sending and retains unacknowledged entries for a later flush. It does not
launch agents, claim handoffs, or open ai-memory's database or wiki.
The orchestrator still maps its events to the native harness payloads and sets
`AI_MEMORY_CAPTURE_OWNER` when launching that harness. The relay uses the native
session identity and derives retry keys from the producer's stable event IDs.
Only the first pending event for each session enters a batch; that session
advances after acknowledgement, even when other sessions are rate-limited.
The package has its own workspace, tests and CLI. Its README defines queue limits,
local data handling and recovery, with an executable test against the real
ai-memory server. Root workspace tests do not run the companion's unit tests.
## `ai-memory-importer`: migration and ingestion companion
This is the companion shape for PR #118. The first implemented companion lives
+3
View File
@@ -149,6 +149,9 @@ Use the tuple recipe when event IDs have narrower scope.
saturation rather than opening an unbounded number of requests.
Maintain a durable producer-side queue if offline/restart recovery matters.
The optional [lifecycle relay](../companions/ai-memory-relay) supplies one through
a separate CLI. It accepts events from the orchestrator and sends only the first
pending event of each session in a batch.
Flush earlier events for a session before sending its terminal event. Serialize
delivery within a session where order matters; different sessions can share the
server concurrently. Late observations and a later terminal event use existing
+863
View File
@@ -0,0 +1,863 @@
#!/usr/bin/env python3
"""Deterministic end-to-end smoke test for the external lifecycle relay.
Run with:
uv run --no-project python tests/e2e/external_relay_smoke.py \
--ai-memory-bin /absolute/ai-memory \
--relay-bin /absolute/ai-memory-relay
The test starts only isolated local processes and gives them a minimal,
synthetic environment. Every invocation preserves its temporary root, including
logs, queue files, payload fixtures, and a redacted diagnostic summary. Pass
``--artifacts-parent DIR`` to choose the parent directory. Exit status zero
means every assertion in the smoke test passed.
"""
from __future__ import annotations
import argparse
import concurrent.futures
import contextlib
import http.client
import http.server
import json
import os
import socket
import sqlite3
import subprocess
import tempfile
import threading
import time
import urllib.error
import urllib.request
import uuid
from pathlib import Path
from typing import Any, Iterable
WORKSPACE = "external-relay-e2e"
PROJECT = "relay-smoke"
ACTOR = "relay-test-actor"
AGENT = "claude-code"
GOOD_TOKEN = "FAKEfake0123456789relaygood"
BAD_TOKEN = "FAKEfake0123456789relaywrong"
MARKER = "external_relay_deterministic_handoff_marker"
EVENTS = ("session-start", "user-prompt-submit", "session-end")
class SmokeFailure(RuntimeError):
"""An assertion or subprocess failed."""
class DropFirstResponseProxy:
"""Forward requests to the real server and drop responses until released."""
def __init__(self, target_port: int, timeout_seconds: float) -> None:
self.target_port = target_port
self.timeout_seconds = timeout_seconds
self.mode_lock = threading.Lock()
self.drop_responses = True
proxy = self
class Handler(http.server.BaseHTTPRequestHandler):
protocol_version = "HTTP/1.1"
def do_POST(self) -> None: # noqa: N802 - BaseHTTPRequestHandler API
length = int(self.headers.get("Content-Length", "0"))
body = self.rfile.read(length)
headers = {
key: value
for key, value in self.headers.items()
if key.lower() not in {"host", "connection", "content-length"}
}
upstream = http.client.HTTPConnection(
"127.0.0.1",
proxy.target_port,
timeout=proxy.timeout_seconds,
)
try:
upstream.request("POST", self.path, body=body, headers=headers)
response = upstream.getresponse()
response_body = response.read()
with proxy.mode_lock:
drop = proxy.drop_responses
if drop:
self.connection.shutdown(socket.SHUT_RDWR)
self.connection.close()
return
self.send_response(response.status)
for key, value in response.getheaders():
if key.lower() not in {
"connection",
"content-length",
"transfer-encoding",
}:
self.send_header(key, value)
self.send_header("Content-Length", str(len(response_body)))
self.end_headers()
self.wfile.write(response_body)
finally:
upstream.close()
def log_message(self, _format: str, *_args: Any) -> None:
return
self.server = http.server.ThreadingHTTPServer(("127.0.0.1", 0), Handler)
self.thread = threading.Thread(
target=self.server.serve_forever,
name="external-relay-drop-response-proxy",
daemon=True,
)
self.thread.start()
@property
def url(self) -> str:
return f"http://127.0.0.1:{self.server.server_port}"
def stop(self) -> None:
self.server.shutdown()
self.server.server_close()
self.thread.join(timeout=5)
if self.thread.is_alive():
raise SmokeFailure("fault proxy thread did not stop")
def forward_responses(self) -> None:
with self.mode_lock:
self.drop_responses = False
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(
description=(
"Run an isolated real-server E2E test for ai-memory-relay. All "
"artifacts are retained; exit status 0 means every assertion passed."
),
epilog=(
"Example: uv run --no-project python tests/e2e/external_relay_smoke.py "
"--ai-memory-bin /absolute/ai-memory "
"--relay-bin /absolute/ai-memory-relay"
),
)
parser.add_argument(
"--ai-memory-bin",
required=True,
type=Path,
help="absolute path to the ai-memory executable",
)
parser.add_argument(
"--relay-bin",
required=True,
type=Path,
help="absolute path to the ai-memory-relay executable",
)
parser.add_argument(
"--timeout-seconds",
type=float,
default=60.0,
help="bounded timeout for readiness and each subprocess (default: 60)",
)
parser.add_argument(
"--artifacts-parent",
type=Path,
help="existing directory under which the retained temp root is created",
)
args = parser.parse_args()
if args.timeout_seconds < 1:
parser.error("--timeout-seconds must be at least 1")
for option in ("ai_memory_bin", "relay_bin"):
value = getattr(args, option)
if not value.is_absolute():
parser.error(f"--{option.replace('_', '-')} must be an absolute path")
if not value.is_file() or not os.access(value, os.X_OK):
parser.error(f"--{option.replace('_', '-')} is not an executable file: {value}")
if args.artifacts_parent is not None:
args.artifacts_parent = args.artifacts_parent.resolve()
if not args.artifacts_parent.is_dir():
parser.error("--artifacts-parent must name an existing directory")
return args
def reserve_port() -> int:
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as listener:
listener.bind(("127.0.0.1", 0))
return int(listener.getsockname()[1])
def write_json(path: Path, value: Any) -> None:
path.write_text(json.dumps(value, indent=2, sort_keys=True) + "\n", encoding="utf-8")
def parse_status(stdout: str, label: str) -> dict[str, int]:
def unique_object(pairs: list[tuple[str, Any]]) -> dict[str, Any]:
result: dict[str, Any] = {}
for key, value in pairs:
if key in result:
raise SmokeFailure(f"{label} status JSON duplicated field {key}")
result[key] = value
return result
try:
parsed = json.loads(stdout, object_pairs_hook=unique_object)
except json.JSONDecodeError as error:
raise SmokeFailure(f"{label} did not emit one JSON object") from error
if not isinstance(parsed, dict):
raise SmokeFailure(f"{label} status JSON is not an object")
pending_items = parsed.get("pending_items")
receipts = parsed.get("receipts")
if type(pending_items) is not int:
raise SmokeFailure(f"{label} status JSON has no integer pending_items")
if type(receipts) is not int:
raise SmokeFailure(f"{label} status JSON has no integer receipts")
return {"pending_items": pending_items, "receipts": receipts}
class Harness:
def __init__(self, args: argparse.Namespace) -> None:
self.args = args
parent = str(args.artifacts_parent) if args.artifacts_parent else None
self.root = Path(tempfile.mkdtemp(prefix="ai-memory-external-relay-", dir=parent))
self.home = self.root / "home"
self.data = self.root / "data"
self.project = self.root / "project"
self.queues = self.root / "queues"
self.payloads = self.root / "payloads"
self.logs = self.root / "logs"
self.tmp = self.root / "tmp"
for directory in (
self.home,
self.data,
self.project,
self.queues,
self.payloads,
self.logs,
self.tmp,
):
directory.mkdir()
(self.project / ".ai-memory.toml").write_text(
f'workspace = "{WORKSPACE}"\nproject = "{PROJECT}"\n', encoding="utf-8"
)
self.empty_gitconfig = self.root / "empty-gitconfig"
self.empty_gitconfig.write_text("", encoding="utf-8")
self.port = reserve_port()
self.base_url = f"http://127.0.0.1:{self.port}"
self.server: subprocess.Popen[bytes] | None = None
self.server_log_handle: Any = None
self.proxies: list[DropFirstResponseProxy] = []
self.command_index = 0
self.command_lock = threading.Lock()
self.checks: dict[str, Any] = {}
def environment(self, token: str | None = GOOD_TOKEN) -> dict[str, str]:
env = {
"HOME": str(self.home),
"USERPROFILE": str(self.home),
"XDG_CONFIG_HOME": str(self.home / ".config"),
"XDG_DATA_HOME": str(self.home / ".local" / "share"),
"TMPDIR": str(self.tmp),
"TEMP": str(self.tmp),
"TMP": str(self.tmp),
"PATH": os.defpath,
"LANG": "C.UTF-8",
"LC_ALL": "C.UTF-8",
"GIT_CONFIG_NOSYSTEM": "1",
"GIT_CONFIG_GLOBAL": str(self.empty_gitconfig),
"GIT_CONFIG_SYSTEM": str(self.empty_gitconfig),
"AI_MEMORY_HOME": str(self.home),
"AI_MEMORY_DATA_DIR": str(self.data),
"AI_MEMORY_EMBEDDING_PROVIDER": "none",
"AI_MEMORY_BACKFILL_ON_START": "false",
"AI_MEMORY_CONSOLIDATE_ON_SESSION_END": "false",
"AI_MEMORY_AUTO_IMPROVE__SCHEDULER__ENABLED": "false",
"RUST_LOG": "off",
}
if os.name == "nt":
for key in ("SystemRoot", "WINDIR", "SystemDrive"):
value = os.environ.get(key)
if value:
env[key] = value
if token is not None:
env["AI_MEMORY_AUTH_TOKEN"] = token
return env
def run(
self,
binary: Path,
arguments: Iterable[str | Path],
*,
token: str | None = GOOD_TOKEN,
input_text: str | None = None,
expect: int | None = 0,
label: str,
) -> subprocess.CompletedProcess[str]:
command = [str(binary), *(str(arg) for arg in arguments)]
started = time.monotonic()
try:
result = subprocess.run(
command,
cwd=self.project,
env=self.environment(token),
input=input_text,
text=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
timeout=self.args.timeout_seconds,
check=False,
)
except subprocess.TimeoutExpired as error:
raise SmokeFailure(f"{label} exceeded {self.args.timeout_seconds:g}s") from error
with self.command_lock:
self.command_index += 1
log_path = self.logs / f"command-{self.command_index:03d}.json"
write_json(
log_path,
{
"label": label,
"argv": command,
"duration_seconds": round(time.monotonic() - started, 3),
"returncode": result.returncode,
"stdout": result.stdout,
"stderr": result.stderr,
},
)
if expect is not None and result.returncode != expect:
raise SmokeFailure(f"{label} exited {result.returncode}, expected {expect}")
return result
def relay(self, arguments: Iterable[str | Path], **kwargs: Any) -> subprocess.CompletedProcess[str]:
return self.run(self.args.relay_bin, arguments, **kwargs)
def queue_init(
self, name: str, producer: str, server_url: str | None = None
) -> Path:
queue = self.queues / name
self.relay(
[
"init",
"--queue-dir",
queue,
"--server-url",
server_url or self.base_url,
"--producer",
producer,
"--actor",
ACTOR,
"--workspace",
WORKSPACE,
"--project",
PROJECT,
],
label=f"init-{name}",
)
return queue
def status(self, queue: Path, label: str) -> tuple[Any, int]:
result = self.relay(["status", "--queue-dir", queue], label=label)
parsed = parse_status(result.stdout, label)
return parsed, parsed["pending_items"]
def enqueue(self, queue: Path, payload_file: Path, label: str) -> None:
self.relay(
["enqueue", "--queue-dir", queue, "--file", payload_file], label=label
)
def flush(
self,
queue: Path,
label: str,
token: str = GOOD_TOKEN,
expect: int = 0,
) -> int:
return self.relay(
["flush", "--queue-dir", queue],
token=token,
expect=expect,
label=label,
).returncode
def start_server(self) -> None:
server_log = self.logs / "server.log"
self.server_log_handle = server_log.open("xb")
self.server = subprocess.Popen(
[
str(self.args.ai_memory_bin),
"--data-dir",
str(self.data),
"serve",
"--transport",
"http",
"--bind",
f"127.0.0.1:{self.port}",
"--workspace",
WORKSPACE,
"--project",
PROJECT,
"--no-watcher",
],
cwd=self.project,
env=self.environment(GOOD_TOKEN),
stdin=subprocess.DEVNULL,
stdout=self.server_log_handle,
stderr=subprocess.STDOUT,
)
deadline = time.monotonic() + self.args.timeout_seconds
while time.monotonic() < deadline:
if self.server.poll() is not None:
raise SmokeFailure(f"ai-memory server exited {self.server.returncode} before readiness")
request = urllib.request.Request(
f"{self.base_url}/mcp",
headers={"Authorization": f"Bearer {GOOD_TOKEN}"},
)
try:
with urllib.request.urlopen(request, timeout=0.5):
return
except urllib.error.HTTPError as error:
if error.code in (400, 404, 405):
return
except (urllib.error.URLError, TimeoutError):
pass
time.sleep(0.05)
raise SmokeFailure("ai-memory server readiness timeout")
def stop_server(self) -> None:
process = self.server
if process is None:
return
self.server = None
if process.poll() is None:
process.terminate()
try:
process.wait(timeout=min(5.0, self.args.timeout_seconds))
except subprocess.TimeoutExpired:
process.kill()
process.wait(timeout=min(5.0, self.args.timeout_seconds))
if self.server_log_handle is not None:
self.server_log_handle.close()
self.server_log_handle = None
def start_drop_response_proxy(self) -> DropFirstResponseProxy:
proxy = DropFirstResponseProxy(self.port, self.args.timeout_seconds)
self.proxies.append(proxy)
return proxy
def stop_proxies(self) -> None:
while self.proxies:
self.proxies.pop().stop()
def db(self) -> sqlite3.Connection:
path = self.data / "db" / "memory.sqlite"
return sqlite3.connect(f"file:{path}?mode=ro", uri=True)
def observation_count(self) -> int:
with contextlib.closing(self.db()) as connection:
return int(connection.execute("SELECT COUNT(*) FROM observations").fetchone()[0])
def hook_spool_json_count(self) -> int:
spool = self.data / "hook-spool"
if not spool.is_dir():
return 0
return sum(1 for path in spool.rglob("*.json") if path.is_file())
def observations(self, session_id: uuid.UUID) -> list[tuple[str, str | None, str | None, str]]:
with contextlib.closing(self.db()) as connection:
return list(
connection.execute(
"SELECT kind, extension, source_event, body FROM observations "
"WHERE session_id = ? ORDER BY created_at, rowid",
(session_id.bytes,),
)
)
def handoff_state(self, from_session: uuid.UUID) -> str | None:
with contextlib.closing(self.db()) as connection:
row = connection.execute(
"SELECT state FROM handoffs WHERE from_session_id = ? ORDER BY created_at DESC LIMIT 1",
(from_session.bytes,),
).fetchone()
return None if row is None else str(row[0])
def native_hook(self, event: str, session_id: uuid.UUID) -> subprocess.CompletedProcess[str]:
payload = json.dumps(
{
"session_id": str(session_id),
"cwd": str(self.project),
"prompt": "native duplicate must remain uncaptured",
}
)
env = self.environment(GOOD_TOKEN)
env["AI_MEMORY_CAPTURE_OWNER"] = "external-relay-e2e"
command = [
str(self.args.ai_memory_bin),
"--data-dir",
str(self.data),
"hook",
"--event",
event,
"--agent",
AGENT,
"--server-url",
self.base_url,
"--auth-token",
GOOD_TOKEN,
]
try:
result = subprocess.run(
command,
cwd=self.project,
env=env,
input=payload,
text=True,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
timeout=self.args.timeout_seconds,
check=False,
)
except subprocess.TimeoutExpired as error:
raise SmokeFailure(f"native-{event} exceeded timeout") from error
self.command_index += 1
write_json(
self.logs / f"command-{self.command_index:03d}.json",
{
"label": f"native-{event}",
"argv": [arg if arg != GOOD_TOKEN else "<redacted>" for arg in command],
"returncode": result.returncode,
"stdout": result.stdout,
"stderr": result.stderr,
},
)
if result.returncode != 0:
raise SmokeFailure(f"native-{event} exited {result.returncode}")
return result
def lifecycle(session_id: uuid.UUID, prefix: str, prompt: str = MARKER) -> list[dict[str, Any]]:
return [
{
"event_id": f"{prefix}-start",
"agent": AGENT,
"event": "session-start",
"body": {"session_id": str(session_id), "cwd": "", "source": prefix},
},
{
"event_id": f"{prefix}-prompt",
"agent": AGENT,
"event": "user-prompt-submit",
"body": {"session_id": str(session_id), "cwd": "", "prompt": prompt},
},
{
"event_id": f"{prefix}-end",
"agent": AGENT,
"event": "session-end",
"body": {"session_id": str(session_id), "cwd": "", "reason": "completed"},
},
]
def materialize(harness: Harness, name: str, events: list[dict[str, Any]]) -> Path:
copied = json.loads(json.dumps(events))
for item in copied:
item["body"]["cwd"] = str(harness.project)
path = harness.payloads / f"{name}.json"
write_json(path, copied)
return path
def require(condition: bool, message: str) -> None:
if not condition:
raise SmokeFailure(message)
def run_smoke(harness: Harness) -> None:
main_session = uuid.uuid4()
main_queue = harness.queue_init("offline-main", "producer-a")
initial_status, initial_pending = harness.status(main_queue, "status-initial")
require(initial_pending == 0, f"new queue has {initial_pending} pending events")
canonical_file = materialize(
harness, "canonical-offline", lifecycle(main_session, "canonical")
)
harness.enqueue(main_queue, canonical_file, "enqueue-offline")
_, offline_pending = harness.status(main_queue, "status-offline-enqueued")
require(offline_pending == 3, f"offline enqueue retained {offline_pending}, expected 3")
harness.start_server()
harness.flush(main_queue, "flush-after-server-start")
_, flushed_pending = harness.status(main_queue, "status-after-flush")
require(flushed_pending == 0, f"successful flush left {flushed_pending} pending")
main_rows = harness.observations(main_session)
require([row[0] for row in main_rows] == ["session-start", "user-prompt", "session-end"], "canonical lifecycle order differs")
require(all(row[1] == "producer-a" for row in main_rows), "canonical producer extension differs")
require([row[2] for row in main_rows] == list(EVENTS), "canonical source events differ")
before_native = harness.observation_count()
spool_before_native = harness.hook_spool_json_count()
receiver = uuid.uuid4()
native_start = harness.native_hook("session-start", receiver)
require(MARKER in native_start.stdout, "capture-owner SessionStart did not receive handoff")
harness.native_hook("user-prompt-submit", receiver)
harness.native_hook("session-end", receiver)
after_native = harness.observation_count()
spool_after_native = harness.hook_spool_json_count()
require(after_native == before_native, "capture-owner native hooks created observations")
require(
spool_after_native == spool_before_native,
"capture-owner native hooks created spool JSON files",
)
require(harness.handoff_state(main_session) == "accepted", "capture-owner hook did not preserve and accept the external handoff")
harness.enqueue(main_queue, canonical_file, "enqueue-identical-replay")
harness.flush(main_queue, "flush-identical-replay")
require(len(harness.observations(main_session)) == 3, "identical replay duplicated observations")
collision_before = harness.observation_count()
collision_queue = harness.queue_init(
"collision-foreign-agent", "producer-collision"
)
collision_status_before, _ = harness.status(
collision_queue, "status-before-collision"
)
collision_receipts_before = collision_status_before["receipts"]
collision_file = materialize(
harness,
"session-collision",
[
{
"event_id": "wrong-agent-collision",
"agent": "codex",
"event": "user-prompt-submit",
"body": {
"session_id": str(main_session),
"cwd": "",
"prompt": "valid wrong-agent event expected to be acknowledged and dropped",
},
}
],
)
harness.enqueue(collision_queue, collision_file, "enqueue-session-collision")
harness.flush(collision_queue, "flush-session-collision")
collision_status_after, collision_pending = harness.status(
collision_queue, "status-after-collision"
)
collision_receipts_after = collision_status_after["receipts"]
require(collision_pending == 0, "acknowledged SessionCollision remained pending")
require(
collision_receipts_after == collision_receipts_before + 1,
"SessionCollision acknowledgement did not create one relay receipt",
)
collision_stored_delta = harness.observation_count() - collision_before
require(collision_stored_delta == 0, "SessionCollision created a stored observation")
same_body = "identical body with distinct external identities"
producer_a_session = uuid.uuid4()
producer_a_file = materialize(
harness,
"distinct-producer-a",
[
{
"event_id": "shared-event-id",
"agent": AGENT,
"event": "user-prompt-submit",
"body": {"session_id": str(producer_a_session), "cwd": "", "prompt": same_body},
},
{
"event_id": "different-event-id",
"agent": AGENT,
"event": "user-prompt-submit",
"body": {"session_id": str(producer_a_session), "cwd": "", "prompt": same_body},
},
],
)
harness.enqueue(main_queue, producer_a_file, "enqueue-distinct-event-ids")
harness.flush(main_queue, "flush-distinct-event-ids")
producer_a_rows = harness.observations(producer_a_session)
require(len(producer_a_rows) == 2, "distinct event IDs collapsed")
producer_b_queue = harness.queue_init("producer-b", "producer-b")
producer_b_file = materialize(
harness,
"distinct-producer-b",
[
{
"event_id": "shared-event-id",
"agent": AGENT,
"event": "user-prompt-submit",
"body": {"session_id": str(producer_a_session), "cwd": "", "prompt": same_body},
}
],
)
harness.enqueue(producer_b_queue, producer_b_file, "enqueue-distinct-producer")
harness.flush(producer_b_queue, "flush-distinct-producer")
combined_producer_rows = harness.observations(producer_a_session)
require(len(combined_producer_rows) == 3, "distinct producer collapsed")
require(
sum(row[1] == "producer-a" for row in combined_producer_rows) == 2
and sum(row[1] == "producer-b" for row in combined_producer_rows) == 1,
"producer namespaces were not preserved",
)
proxy = harness.start_drop_response_proxy()
lost_response_session = uuid.uuid4()
lost_response_queue = harness.queue_init(
"lost-response", "producer-loss", server_url=proxy.url
)
lost_response_file = materialize(
harness,
"lost-response",
[
{
"event_id": "accepted-before-response-loss",
"agent": AGENT,
"event": "user-prompt-submit",
"body": {
"session_id": str(lost_response_session),
"cwd": "",
"prompt": "response loss idempotency marker",
},
}
],
)
harness.enqueue(lost_response_queue, lost_response_file, "enqueue-lost-response")
lost_response_exit = harness.flush(
lost_response_queue, "flush-lost-response", expect=2
)
_, lost_response_pending = harness.status(
lost_response_queue, "status-lost-response"
)
require(lost_response_pending == 1, "response loss did not retain the queue item")
require(
len(harness.observations(lost_response_session)) == 1,
"real server did not accept exactly one observation before response loss",
)
proxy.forward_responses()
harness.flush(lost_response_queue, "flush-lost-response-retry")
_, lost_response_retry_pending = harness.status(
lost_response_queue, "status-lost-response-retry"
)
require(lost_response_retry_pending == 0, "lost-response retry remained pending")
require(
len(harness.observations(lost_response_session)) == 1,
"lost-response retry duplicated the real-server observation",
)
concurrent_sessions: list[tuple[uuid.UUID, Path]] = []
for index in range(15):
session_id = uuid.uuid4()
queue = harness.queue_init(f"concurrent-{index:02d}", f"parallel-{index % 3}")
payload = materialize(
harness,
f"concurrent-{index:02d}",
lifecycle(session_id, f"parallel-{index:02d}", prompt=f"parallel marker {index}"),
)
harness.enqueue(queue, payload, f"enqueue-concurrent-{index:02d}")
concurrent_sessions.append((session_id, queue))
with concurrent.futures.ThreadPoolExecutor(max_workers=15) as executor:
futures = [
executor.submit(harness.flush, queue, f"flush-concurrent-{index:02d}")
for index, (_, queue) in enumerate(concurrent_sessions)
]
for future in futures:
future.result(timeout=harness.args.timeout_seconds + 5)
for session_id, queue in concurrent_sessions:
rows = harness.observations(session_id)
require([row[0] for row in rows] == ["session-start", "user-prompt", "session-end"], f"concurrent session {session_id} lifecycle order differs")
_, pending = harness.status(queue, f"status-concurrent-{session_id.hex[:8]}")
require(pending == 0, f"concurrent queue {session_id} retained {pending}")
auth_session = uuid.uuid4()
auth_queue = harness.queue_init("auth-retry", "producer-auth")
auth_file = materialize(
harness,
"auth-retry",
[
{
"event_id": "auth-event",
"agent": AGENT,
"event": "user-prompt-submit",
"body": {"session_id": str(auth_session), "cwd": "", "prompt": "auth retry marker"},
}
],
)
harness.enqueue(auth_queue, auth_file, "enqueue-auth-retry")
bad_code = harness.flush(
auth_queue, "flush-auth-rejected", token=BAD_TOKEN, expect=2
)
_, rejected_pending = harness.status(auth_queue, "status-auth-rejected")
require(rejected_pending == 1, f"auth rejection retained {rejected_pending}, expected 1")
require(len(harness.observations(auth_session)) == 0, "unauthorized flush wrote an observation")
harness.flush(auth_queue, "flush-auth-retry-correct")
_, retry_pending = harness.status(auth_queue, "status-auth-retry-correct")
require(retry_pending == 0, f"authorized retry left {retry_pending} pending")
require(len(harness.observations(auth_session)) == 1, "authorized retry lost or duplicated event")
final_count = harness.observation_count()
expected_count = 3 + 2 + 1 + 1 + (15 * 3) + 1
require(final_count == expected_count, f"final observation count {final_count}, expected {expected_count}")
harness.checks = {
"initial_queue_status": initial_status,
"initial_pending": initial_pending,
"offline_pending": offline_pending,
"offline_restart_flush": "passed",
"canonical_lifecycle": "passed",
"capture_owner_zero_new_observations": after_native - before_native,
"capture_owner_new_spool_json": spool_after_native - spool_before_native,
"handoff_state": harness.handoff_state(main_session),
"identical_replay_count": len(harness.observations(main_session)),
"session_collision": {
"relay_acknowledged": collision_pending == 0,
"receipt_delta": collision_receipts_after - collision_receipts_before,
"stored_observation_delta": collision_stored_delta,
},
"distinct_event_id_count": len(producer_a_rows),
"distinct_producer_count": sum(
row[1] == "producer-b" for row in combined_producer_rows
),
"lost_response": {
"first_exit": lost_response_exit,
"pending_after_loss": lost_response_pending,
"pending_after_retry": lost_response_retry_pending,
"stored_observations": len(
harness.observations(lost_response_session)
),
},
"concurrent_sessions": len(concurrent_sessions),
"concurrent_observations": sum(
len(harness.observations(sid)) for sid, _ in concurrent_sessions
),
"auth_rejected_exit": bad_code,
"auth_rejected_pending": rejected_pending,
"auth_retry_pending": retry_pending,
"final_observations": final_count,
"expected_observations": expected_count,
}
def main() -> int:
args = parse_args()
harness = Harness(args)
passed = False
error_message: str | None = None
try:
run_smoke(harness)
passed = True
return 0
except Exception as error: # diagnostics must survive every assertion failure
error_message = f"{type(error).__name__}: {error}"
return 1
finally:
try:
harness.stop_proxies()
finally:
harness.stop_server()
summary = {
"status": "passed" if passed else "failed",
"artifacts": str(harness.root),
"checks": harness.checks,
}
if error_message is not None:
summary["error"] = error_message
write_json(harness.root / "summary.json", summary)
print(json.dumps(summary, sort_keys=True))
if __name__ == "__main__":
raise SystemExit(main())