feat(desktop): add cross-workspace agent-library data model (Phase 1)

The Phase 0 seam made every persona writer merge-preserving; Phase 1 lands
the data model that seam was built for: `library.json` and the scope-local
crash journals, plus their quarantine-preserving IO. No operation mutates
them yet — share, edit, materialize, delete, and deploy land in Phases 2-5.

`library.rs` models the versioned document envelope, `SharedDefinition`
(the allowlisted content a share may carry — credentials, identity,
`env_vars`, activation, and projection metadata are structurally absent),
`LibraryEntry` with its per-scope `ProjectionState` machine, verified
`IdentityBinding`s, and the permanent `DeferredArchive`/`RemovalManifest`
retirement markers. `load_library` classifies per §2.1: an absent file is a
valid empty library, whole-document corruption is preserved as `.invalid`
and blocks all mutation, an unknown version is read-only fail, and a single
malformed / semantically invalid / identity-colliding entry is quarantined
raw while healthy siblings stay usable. Identity collisions on `library_id`,
a live origin, or a non-terminal `(scope, slug)` claim group-quarantine every
collider — never first-wins.

`library/journals.rs` adds the scope-local `pending-agent-keys.json` and
`deploy-intents.json`, which share the owning workspace's failure domain
rather than the library's so a broken `library.json` never blocks a
workspace-local create or deploy. Their failure domains are asymmetric: a
pending-keys failure degrades only that scope's create/import path, while a
deploy-intents failure — unknown version, syntax error, or a duplicate-pubkey
mutex violation validated on read — fails the scope's destructive and deploy
paths closed.

`apply_shared_definition` is the sole writer of shared content onto a scoped
keyless record, assigning exactly the shared slots plus revision/timestamp and
mirroring `into_agent_record` so a populated record is byte-identical to a
freshly projected one. `ManagedAgentRecord` gains
`last_completed_deploy_attempt_id`, the deploy-provenance stamp that forms an
inseparable pair with `backend_agent_id` in `copy_runtime_state`.

Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
Duncan
2026-08-11 06:06:20 -04:00
co-authored by Will Pfleger
parent 1ac3f847c7
commit 887793ce04
36 changed files with 1512 additions and 7 deletions
@@ -114,6 +114,7 @@ fn agent_record() -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -26,6 +26,13 @@ fn copy_runtime_state(from: &ManagedAgentRecord, to: &mut ManagedAgentRecord) {
to.runtime_pid = from.runtime_pid;
to.backend = from.backend.clone();
to.backend_agent_id.clone_from(&from.backend_agent_id);
// `backend_agent_id` and `last_completed_deploy_attempt_id` are one
// inseparable deploy-provenance pair (§2.6, P14-I2): a stamp missing here
// would let a rename rollback manufacture a record whose backend id and
// attempt stamp disagree, or make `same_configuration` reject the rollback
// as an unrelated change. Copy both, never one without the other.
to.last_completed_deploy_attempt_id
.clone_from(&from.last_completed_deploy_attempt_id);
to.provider_binary_path
.clone_from(&from.provider_binary_path);
to.last_started_at.clone_from(&from.last_started_at);
+3 -3
View File
@@ -771,9 +771,8 @@ pub async fn create_managed_agent(
agent_command_override,
agent_args,
mcp_command,
// BUZZ_ACP_TURN_TIMEOUT is deprecated and ignored by the harness;
// store the schema default only. Use idle_timeout_seconds or
// max_turn_duration_seconds for actual turn-length control.
// BUZZ_ACP_TURN_TIMEOUT is deprecated and ignored by the harness; store
// the schema default. Use idle/max_turn_duration_seconds for turn length.
turn_timeout_seconds: DEFAULT_AGENT_TURN_TIMEOUT_SECONDS,
// 0 or None → harness uses its own default (320s idle, 3600s max), and the CLI also clamps 0 → minimum.
idle_timeout_seconds: input.idle_timeout_seconds.filter(|s| *s > 0),
@@ -825,6 +824,7 @@ pub async fn create_managed_agent(
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -61,6 +61,7 @@ fn bare_agent_record(
auto_restart_on_config_change: false,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -104,6 +104,7 @@ async fn test_full_tail_stop_spawn_receipt_register_save() {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Default::default(),
definition_parallelism: None,
@@ -544,6 +545,7 @@ async fn test_relay_mesh_preflight_precedes_stop() {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Default::default(),
definition_parallelism: None,
@@ -631,6 +631,7 @@ async fn test_record_mesh_change_after_preflight_aborts_before_stop() {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Default::default(),
definition_parallelism: None,
@@ -69,6 +69,7 @@ fn make_agent(
auto_restart_on_config_change: false,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -213,6 +213,7 @@ fn local_agent() -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -62,6 +62,7 @@ fn make_definition(slug: &str) -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -647,6 +647,7 @@ where
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: respond_to_wire.clone(),
definition_respond_to_allowlist: minted.respond_to_allowlist.clone(),
definition_parallelism: minted_parallelism,
@@ -71,6 +71,7 @@ fn make_definition(slug: &str) -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -56,6 +56,7 @@ fn agent(persona_id: &str, name: &str, display_name: Option<&str>) -> ManagedAge
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -639,6 +639,7 @@ where
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: respond_to_wire.clone(),
definition_respond_to_allowlist: definition.respond_to_allowlist.clone(),
definition_parallelism: minted_parallelism,
@@ -227,6 +227,7 @@ fn team_export_with_instance_and_memory_level_uses_supplied_entries() {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -214,6 +214,7 @@ mod tests {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -414,6 +414,7 @@ mod tests {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -69,6 +69,7 @@ fn minimal_record() -> ManagedAgentRecord {
source_team_persona_slug: Some("lep".to_string()), // MUST NOT appear
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: Some("allowlist".to_string()),
catalog_source: None,
definition_respond_to_allowlist: vec!["abc123def".to_string()],
@@ -113,6 +113,7 @@ fn test_record() -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -280,6 +280,7 @@ fn record_with(
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -303,8 +304,7 @@ fn record_agent_command_override_beats_runtime() {
#[test]
fn record_agent_command_legacy_persona_fallback() {
// Pre-migration record: persona_id set, no runtime — resolves through
// the legacy persona path unchanged.
// Pre-migration record: persona_id set, no runtime — legacy path unchanged.
let personas = vec![persona_with_runtime("p1", Some("goose"))];
let record = record_with(None, Some("p1"), None);
assert_eq!(record_agent_command(&record, &personas), "goose");
@@ -91,6 +91,7 @@ fn record(
auto_restart_on_config_change: false,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -351,6 +351,7 @@ fn bare_record() -> ManagedAgentRecord {
auto_restart_on_config_change: false,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -0,0 +1,500 @@
//! Cross-workspace agent library — device-local, unscoped metadata (§2).
//!
//! `library.json` (`<agents-base>/library.json`, `0o600`, never published) is
//! the single source of truth for shared agent definitions, their per-scope
//! projections, and the crash-safe identity bindings that carry one npub across
//! a same-owner's workspaces. This module is the Phase-1 data model and its
//! quarantine-preserving IO; the operations that mutate it (share, edit,
//! materialize, delete, deploy) land in Phases 2-5. The scope-local journals
//! that must share a *workspace's* failure domain rather than the library's
//! (P9-I1) live in the [`journals`] submodule.
//!
//! Fault model (§2.1):
//! - Absent file = a valid empty library (version 1, no entries).
//! - Whole-document syntax failure → preserved as `library.json.invalid`
//! (copy, not rename — the `storage.rs` loud-failure discipline) and NO
//! library or scoped mutation proceeds; workspaces keep running from their
//! scoped caches.
//! - Unknown/forward `version` → read-only fail: no library or scoped mutation.
//! - A single malformed, semantically invalid, or identity-colliding entry is
//! quarantined (kept raw, surfaced as a degradation, never merged, never a
//! winner) while its healthy siblings stay fully usable, and every healthy
//! rewrite preserves the quarantined entry value-equivalently (P2-I2).
//!
//! Items are `pub(crate)` for the Phase 2-5 callers; suppress the dead-code
//! lint until those land (matches `team_snapshot.rs`).
#![allow(dead_code)]
use std::collections::{BTreeMap, HashMap};
use std::path::{Path, PathBuf};
use nostr::PublicKey;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use super::storage::{atomic_write_json_restricted, backup_invalid_store};
use super::types::ManagedAgentRecord;
pub(crate) mod journals;
/// The only `library.json` schema version v1 understands. A document whose
/// `version` differs is a forward format: read-only fail, never migrated.
pub(crate) const SUPPORTED_LIBRARY_VERSION: u32 = 1;
/// Resolve the unscoped `library.json` path under the agents base directory.
/// The library is a peer of `scopes/`, never inside a scope — it is device
/// metadata shared across every workspace (§2, layout).
pub(crate) fn library_path(base_dir: &Path) -> PathBuf {
base_dir.join("library.json")
}
// ── Document envelope ─────────────────────────────────────────────────────────
/// The versioned `library.json` envelope (§2.1). Entries are retained RAW and
/// decoded individually so one malformed entry can never reject the whole file,
/// and re-serializing only the valid ones can never silently erase a
/// quarantined sibling.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub(crate) struct LibraryDocument {
/// `1` in v1; an unknown value is a forward format (read-only fail).
pub version: u32,
/// Entries retained RAW and decoded one at a time (§2.1 quarantine).
pub entries: Vec<Value>,
/// Orphan-key journal (§2.5): pubkeys minted for a binding whose commit may
/// not have completed. Non-secret (pubkeys only). The ONLY journal in
/// `library.json` — binding state is library-linked by definition; plain
/// crash journals live scope-local (P9-I1, [`journals`]).
#[serde(default)]
pub orphan_keys: Vec<String>,
}
impl LibraryDocument {
/// A valid empty library (the absent-file and freshly-initialized shape).
pub fn empty() -> Self {
Self {
version: SUPPORTED_LIBRARY_VERSION,
entries: Vec::new(),
orphan_keys: Vec::new(),
}
}
}
// ── Shared content ────────────────────────────────────────────────────────────
/// Allowlisted shared agent content (§2.2). Deliberately NOT `AgentDefinition`:
/// `env_vars`, credentials, identity, `is_active`, catalog `shared`,
/// catalog/team provenance, runtime state, relay, memory, timestamps, and
/// projection metadata are structurally absent — they cannot ride a share.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct SharedDefinition {
pub display_name: String,
pub avatar_url: Option<String>,
pub system_prompt: String,
pub runtime: Option<String>,
pub model: Option<String>,
pub provider: Option<String>,
pub name_pool: Vec<String>,
/// NIP-AP behavioral defaults, WIRE shape (kebab-case string / optional u32)
/// — parsed only at the mint boundary, so an unknown future mode string
/// round-trips byte-identically.
pub respond_to: Option<String>,
pub respond_to_allowlist: Vec<String>,
pub parallelism: Option<u32>,
}
/// The ONLY writer of shared content onto a scoped keyless definition record
/// (§2.2). Assigns exactly the [`SharedDefinition`] slots, plus
/// `library_applied_revision` and `updated_at`; every other field of `local` —
/// identity, `env_vars`, `is_active`, runtime state, `library_ref` — is
/// untouched by construction. The value mapping mirrors
/// `AgentDefinition::into_agent_record` so a record populated this way is
/// byte-identical to one freshly projected.
pub(crate) fn apply_shared_definition(
local: &mut ManagedAgentRecord,
shared: &SharedDefinition,
revision: u64,
) {
local.name = shared.display_name.clone();
local.display_name = Some(shared.display_name.clone());
local.avatar_url = shared.avatar_url.clone();
local.system_prompt = (!shared.system_prompt.is_empty()).then(|| shared.system_prompt.clone());
local.runtime = shared.runtime.clone();
local.model = shared.model.clone();
local.provider = shared.provider.clone();
local.name_pool = shared.name_pool.clone();
local.definition_respond_to = shared.respond_to.clone();
local.definition_respond_to_allowlist = shared.respond_to_allowlist.clone();
local.definition_parallelism = shared.parallelism;
local.library_applied_revision = Some(revision);
local.updated_at = crate::util::now_iso();
}
// ── Library entry ─────────────────────────────────────────────────────────────
/// A shared library entry (§2.3). `library_id` routes scoped records
/// (`library_ref`/`lib-<id>`, §2.6); `origin` is the share idempotency key.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct LibraryEntry {
/// UUID v4, minted once at share commit.
pub library_id: String,
/// `(origin scope_id, origin slug)` — share idempotency key.
pub origin: OriginKey,
/// Monotonic integer; wall-clock is display-only.
pub revision: u64,
/// Tombstone; kept forever in v1.
pub deleted: bool,
/// Provenance guard only — never an identity lookup key.
pub owner_pubkey_at_share: String,
pub shared: SharedDefinition,
/// Owner-keyed PUBLIC identity bindings — pubkey + auth tag, NEVER an nsec.
/// A binding exists only after its nsec is write-and-read-back verified in
/// the device keyring (§2.5); there is no "pending binding" state.
pub identity_bindings: BTreeMap<String, IdentityBinding>,
/// Archive obligations deferred by a protected removal (§3.6 step 3).
/// PERMANENT journaled retirement markers in v1 (P17-C1): no v1 path
/// discharges them; their presence keeps `key_archive_protected(pubkey)`
/// true. SET semantics on `(library_id, scope_id, agent_pubkey)` — every
/// append is an upsert (P15-MINOR); reads treat legacy duplicate rows as
/// one obligation.
#[serde(default)]
pub deferred_archives: Vec<DeferredArchive>,
/// Authoritative per-scope membership — sole source for delete-confirm
/// enumeration, tombstone completion, and key lifetime.
pub projections: BTreeMap<String, ProjectionEntry>,
}
/// `(origin scope_id, origin slug)` — the share idempotency key (§2.3).
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct OriginKey {
pub scope_id: String,
pub slug: String,
}
/// A verified owner→agent identity binding (§2.5). `auth_tag` is the NIP-OA
/// `auth` tag JSON; the map key is the owner pubkey it embeds.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct IdentityBinding {
pub agent_pubkey: String,
pub auth_tag: String,
}
/// One skipped archive (§2.3): the identity `agent_pubkey`, live on `scope_id`'s
/// relay, whose archive request was skipped by a protected removal in THAT
/// scope. A permanent v1 retirement marker — never discharged by any v1 path.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct DeferredArchive {
pub scope_id: String,
pub agent_pubkey: String,
}
/// This scope's projection of a library entry (§2.3).
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct ProjectionEntry {
pub state: ProjectionState,
/// This scope's actual definition slug, fixed at `Pending` commit: the
/// deterministic `lib-<library_id>` for a lazily materialized scope, or the
/// origin record's original human slug for the sharing scope. Always equals
/// the scoped record's `slug` and its instances' `persona_id`.
pub local_slug: String,
/// Non-secret metadata for the ruling-4 confirm dialog, captured at this
/// scope's own activation (or at share, for the sharing scope).
pub relay_url: String,
pub workspace_label: Option<String>,
}
/// The uniform journal-first projection state machine (§2.3, P2-C1/C2). Every
/// bracketed scoped-file step is flanked by library writes: intent precedes the
/// scoped mutation, confirmation follows it.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) enum ProjectionState {
/// Journal intent: materialization not yet confirmed.
Pending,
Materialized {
revision: u64,
},
/// Journal intent: user removed in this workspace; scoped cascade not yet
/// confirmed (P2-C2). `manifest` journaled BEFORE any deletion (P4-C3).
ExcludePending {
manifest: Option<RemovalManifest>,
},
/// Terminal: cascade confirmed; NEVER rematerialize.
Excluded,
/// Library deleted; this scope's cascade not yet confirmed.
DeletePending {
manifest: Option<RemovalManifest>,
},
/// Terminal.
Deleted,
}
impl ProjectionState {
/// `Excluded`/`Deleted` are terminal; the identity index only conflicts on
/// non-terminal `(scope_id, local_slug)` ownership (§2.1 index rule (c)).
pub fn is_terminal(&self) -> bool {
matches!(self, ProjectionState::Excluded | ProjectionState::Deleted)
}
}
/// Everything a removal retry needs after the local records have left disk
/// (§2.3, P4-C3). Captured by the record-owning scope from its own readable
/// cache, in the same or an earlier library write than the FIRST destructive
/// scoped step. Non-secret (pubkeys/coordinates only).
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct RemovalManifest {
/// Persona tombstone coordinate: `(owner pubkey, definition d-tag)`.
pub definition_coordinate: (String, String),
/// Every linked instance pubkey at capture time; each requires the full
/// per-pubkey cleanup set (30177 tombstone + archive, key deletion only per
/// the §2.5 entry-wide lifetime, re-evaluated against the index at execution
/// time, never trusted stale).
pub instances: Vec<String>,
}
// ── Quarantine-preserving load (§2.1) ─────────────────────────────────────────
/// The outcome of reading `library.json` (§2.1). The four states the
/// failure-domain rules key on.
#[derive(Debug)]
pub(crate) enum LibraryLoad {
/// Absent file — a valid empty library (version 1, no entries).
Empty,
/// A healthy version-1 document, classified into healthy + quarantined.
Loaded(LoadedLibrary),
/// Whole-document syntax/shape failure. Preserved as `library.json.invalid`;
/// NO library or scoped mutation may proceed.
Corrupt,
/// Parsed, but `version` is unknown/forward: read-only fail.
UnknownVersion(u32),
}
/// A successfully parsed version-1 library, partitioned by the §2.1 rules. The
/// raw `quarantined` values are retained verbatim so a healthy rewrite
/// preserves them value-equivalently (P2-I2).
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct LoadedLibrary {
/// Individually valid, semantically sound, non-colliding entries.
pub healthy: Vec<LibraryEntry>,
/// Decode / semantic / identity-collision failures, kept RAW.
pub quarantined: Vec<Value>,
pub orphan_keys: Vec<String>,
/// Human-readable reasons, one per quarantined entry, for the UI.
pub degradations: Vec<String>,
}
impl LoadedLibrary {
/// Reconstruct the on-disk document, re-serializing healthy entries and
/// re-attaching every quarantined raw value unchanged (P2-I2 preservation).
pub fn rebuild_document(&self) -> Result<LibraryDocument, String> {
let mut entries = Vec::with_capacity(self.healthy.len() + self.quarantined.len());
for entry in &self.healthy {
entries.push(
serde_json::to_value(entry)
.map_err(|e| format!("serialize library entry {}: {e}", entry.library_id))?,
);
}
entries.extend(self.quarantined.iter().cloned());
Ok(LibraryDocument {
version: SUPPORTED_LIBRARY_VERSION,
entries,
orphan_keys: self.orphan_keys.clone(),
})
}
}
/// Read and classify `library.json` under `base_dir` (§2.1). Never mutates the
/// file; whole-document corruption is preserved as `.invalid` as a read side
/// effect so the loud-failure evidence survives.
pub(crate) fn load_library(base_dir: &Path) -> LibraryLoad {
let path = library_path(base_dir);
let bytes = match std::fs::read(&path) {
Ok(bytes) => bytes,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return LibraryLoad::Empty,
Err(_) => {
// An unreadable-but-present file is preserved and treated as
// corrupt: no mutation may proceed against an unknown on-disk state.
backup_invalid_store(&path);
return LibraryLoad::Corrupt;
}
};
let document: LibraryDocument = match serde_json::from_slice(&bytes) {
Ok(document) => document,
Err(_) => {
backup_invalid_store(&path);
return LibraryLoad::Corrupt;
}
};
if document.version != SUPPORTED_LIBRARY_VERSION {
return LibraryLoad::UnknownVersion(document.version);
}
LibraryLoad::Loaded(classify_entries(document))
}
/// Persist a document (§2.1): atomic temp-file + rename + `0o600`, same as
/// `managed-agents.json` (`storage.rs` `write_agent_store_to_path`). Creates the
/// parent agents dir if needed.
pub(crate) fn save_library_document(
base_dir: &Path,
document: &LibraryDocument,
) -> Result<(), String> {
let path = library_path(base_dir);
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)
.map_err(|e| format!("create agents dir {}: {e}", parent.display()))?;
}
let payload =
serde_json::to_vec_pretty(document).map_err(|e| format!("serialize library.json: {e}"))?;
atomic_write_json_restricted(&path, &payload)
}
/// Per-entry decode + semantic validation + document-wide identity index
/// (§2.1). Every failure quarantines its entry (raw-preserved, degradation);
/// identity collisions group-quarantine ALL colliders (never first-wins).
fn classify_entries(document: LibraryDocument) -> LoadedLibrary {
let LibraryDocument {
entries,
orphan_keys,
..
} = document;
// Pass 1: individual decode + semantic binding validation. Survivors keep
// their ORIGINAL raw value alongside the decoded form, so a later
// group-quarantine preserves the exact on-disk bytes (P2-I2) rather than a
// re-serialization that could drift or fail.
let mut decoded: Vec<(LibraryEntry, Value)> = Vec::new();
let mut quarantined: Vec<Value> = Vec::new();
let mut degradations: Vec<String> = Vec::new();
for raw in entries {
match serde_json::from_value::<LibraryEntry>(raw.clone()) {
Err(e) => {
degradations.push(format!("library entry failed to decode: {e}"));
quarantined.push(raw);
}
Ok(entry) => match validate_entry_bindings(&entry) {
Err(reason) => {
degradations.push(format!(
"library entry {} quarantined: {reason}",
entry.library_id
));
quarantined.push(raw);
}
Ok(()) => decoded.push((entry, raw)),
},
}
}
// Pass 2: document-wide identity index — collisions group-quarantine every
// collider (§2.1 rule; P6-I3), preserved via the pass-1 original raw.
let colliders = identity_collisions(decoded.iter().map(|(entry, _)| entry));
let mut healthy = Vec::new();
for (index, (entry, raw)) in decoded.into_iter().enumerate() {
if colliders.contains(&index) {
degradations.push(format!(
"library entry {} quarantined: identity-index collision",
entry.library_id
));
quarantined.push(raw);
} else {
healthy.push(entry);
}
}
LoadedLibrary {
healthy,
quarantined,
orphan_keys,
degradations,
}
}
/// Semantic binding validation (§2.5 read validation): every binding's auth tag
/// must verify against its `agent_pubkey` AND embed exactly the owner pubkey it
/// is keyed under. Also rejects an entry that binds one `agent_pubkey` under two
/// different owners. Malformed pubkeys/tags fail closed.
fn validate_entry_bindings(entry: &LibraryEntry) -> Result<(), String> {
let mut seen_agents: HashMap<String, String> = HashMap::new();
for (owner_hex, binding) in &entry.identity_bindings {
let owner = PublicKey::from_hex(owner_hex)
.map_err(|_| format!("binding owner {owner_hex} is not a valid pubkey"))?;
let agent = PublicKey::from_hex(&binding.agent_pubkey).map_err(|_| {
format!(
"binding agent {} is not a valid pubkey",
binding.agent_pubkey
)
})?;
let embedded =
buzz_sdk_pkg::nip_oa::verify_auth_tag(&binding.auth_tag, &agent).map_err(|e| {
format!(
"binding auth tag for {} failed verification: {e}",
binding.agent_pubkey
)
})?;
if embedded != owner {
return Err(format!(
"binding auth tag for {} embeds {} but is keyed under {owner_hex}",
binding.agent_pubkey,
embedded.to_hex()
));
}
if let Some(prev) = seen_agents.insert(binding.agent_pubkey.clone(), owner_hex.clone()) {
return Err(format!(
"agent {} is bound under two owners ({prev} and {owner_hex})",
binding.agent_pubkey
));
}
}
Ok(())
}
/// Build the document-wide identity index over the individually valid entries
/// and return the indices of every collider (§2.1 rule; P6-I3). Collisions on
/// (a) `library_id`, (b) live (non-`deleted`) [`OriginKey`], or (c) a
/// non-terminal `(scope_id, local_slug)` projection claim group-quarantine ALL
/// participants — picking a winner would silently rewrite a scope.
fn identity_collisions<'a>(
entries: impl Iterator<Item = &'a LibraryEntry>,
) -> std::collections::HashSet<usize> {
let mut by_library_id: HashMap<&str, Vec<usize>> = HashMap::new();
let mut by_origin: HashMap<&OriginKey, Vec<usize>> = HashMap::new();
let mut by_scope_slug: HashMap<(&str, &str), Vec<usize>> = HashMap::new();
for (index, entry) in entries.enumerate() {
by_library_id
.entry(entry.library_id.as_str())
.or_default()
.push(index);
if !entry.deleted {
by_origin.entry(&entry.origin).or_default().push(index);
}
for (scope_id, projection) in &entry.projections {
if !projection.state.is_terminal() {
by_scope_slug
.entry((scope_id.as_str(), projection.local_slug.as_str()))
.or_default()
.push(index);
}
}
}
let mut colliders = std::collections::HashSet::new();
for group in by_library_id
.values()
.chain(by_origin.values())
.chain(by_scope_slug.values())
{
if group.len() > 1 {
colliders.extend(group.iter().copied());
}
}
colliders
}
#[cfg(test)]
mod tests;
@@ -0,0 +1,285 @@
//! Scope-local crash journals (§2.1, P9-I1). These share the OWNING
//! WORKSPACE's failure domain, not the library's: an unknown-version or corrupt
//! `library.json` must never block a purely workspace-local plain create,
//! import, or deploy, so `pending-agent-keys.json` and `deploy-intents.json`
//! live in `<agents-base>/scopes/<scope_id>/` (the scope's `definitions_dir`),
//! readable and writable whenever that scope is.
//!
//! `scope_id` is NOT stored in any row — it is the directory identity, exactly
//! as with `managed-agents.json`; rows can never migrate across scopes.
//!
//! Failure domains are ASYMMETRIC by operation class (P10-C3):
//! - `pending-agent-keys.json` unreadable → that scope's create/import paths
//! degrade loudly (they cannot journal a fresh key); no other scope and no
//! library operation is affected.
//! - `deploy-intents.json` unreadable OR semantically invalid (unknown version,
//! syntax failure, or a duplicate `agent_pubkey` row — the at-most-one mutex
//! is VALIDATED on read and group-rejected, never first-wins) → the scope's
//! destructive removal and new-deploy paths FAIL CLOSED, because an unreadable
//! journal may hide an intent equivalent to a live remote deployment.
//!
//! Unlike `library.json`, a corrupt journal is NOT copied to `.invalid`: both
//! failure modes above refuse to write, so the corrupt file is never
//! overwritten and needs no separate forensic snapshot. Absent = valid empty.
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use crate::managed_agents::storage::atomic_write_json_restricted;
/// The only journal schema version v1 understands (both files).
pub(crate) const SUPPORTED_JOURNAL_VERSION: u32 = 1;
/// `<scope>/pending-agent-keys.json` — pubkeys minted for a scoped record whose
/// keyring write may have outrun its JSON commit (§2.5 step (2), P8-I1).
pub(crate) fn pending_keys_path(definitions_dir: &Path) -> PathBuf {
definitions_dir.join("pending-agent-keys.json")
}
/// `<scope>/deploy-intents.json` — the per-`agent_pubkey` deployment mutex and
/// durable phase machine (§3.3, P8-C1).
pub(crate) fn deploy_intents_path(definitions_dir: &Path) -> PathBuf {
definitions_dir.join("deploy-intents.json")
}
// ── pending-agent-keys.json ───────────────────────────────────────────────────
/// Non-secret journal of agent pubkeys whose keyring entry may have no
/// committed record yet (§2.1). Reaped at the owning scope's next activation.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct PendingKeysJournal {
pub version: u32,
pub pending: Vec<String>,
}
impl PendingKeysJournal {
pub fn empty() -> Self {
Self {
version: SUPPORTED_JOURNAL_VERSION,
pending: Vec::new(),
}
}
}
/// Read outcome for `pending-agent-keys.json` (§2.1). `Unreadable` degrades the
/// owning scope's create/import paths loudly; no other scope is affected.
#[derive(Debug)]
pub(crate) enum PendingKeysLoad {
/// Absent file — a valid empty journal.
Empty,
Loaded(PendingKeysJournal),
/// Unknown version or syntax failure — degrade loudly.
Unreadable(String),
}
/// Read and classify `pending-agent-keys.json` under `definitions_dir` (§2.1).
pub(crate) fn load_pending_keys(definitions_dir: &Path) -> PendingKeysLoad {
let path = pending_keys_path(definitions_dir);
let bytes = match std::fs::read(&path) {
Ok(bytes) => bytes,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return PendingKeysLoad::Empty,
Err(e) => return PendingKeysLoad::Unreadable(format!("read pending-agent-keys.json: {e}")),
};
let journal: PendingKeysJournal = match serde_json::from_slice(&bytes) {
Ok(journal) => journal,
Err(e) => {
return PendingKeysLoad::Unreadable(format!("parse pending-agent-keys.json: {e}"))
}
};
if journal.version != SUPPORTED_JOURNAL_VERSION {
return PendingKeysLoad::Unreadable(format!(
"pending-agent-keys.json version {} is not supported (expected {SUPPORTED_JOURNAL_VERSION})",
journal.version
));
}
PendingKeysLoad::Loaded(journal)
}
/// Persist `pending-agent-keys.json`: atomic temp-file + rename + `0o600`.
pub(crate) fn save_pending_keys(
definitions_dir: &Path,
journal: &PendingKeysJournal,
) -> Result<(), String> {
save_journal(&pending_keys_path(definitions_dir), journal)
}
// ── deploy-intents.json ───────────────────────────────────────────────────────
/// The per-`agent_pubkey` deployment-intent journal (§2.1). AT MOST ONE row per
/// pubkey — the invariant is VALIDATED on read (see [`load_deploy_intents`]).
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct DeployIntentsJournal {
pub version: u32,
pub intents: Vec<DeployIntent>,
}
impl DeployIntentsJournal {
pub fn empty() -> Self {
Self {
version: SUPPORTED_JOURNAL_VERSION,
intents: Vec::new(),
}
}
}
/// One outstanding deploy attempt (§2.1). `provider_config` OWNS the routing the
/// possible remote lives under until the row is resolved (P13-I2); the canonical
/// record supplies only the fresh payload, never a reconstruction of the lost
/// attempt.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct DeployIntent {
/// Mutex key (`scope_id` is the file's directory).
pub agent_pubkey: String,
pub phase: DeployIntentPhase,
/// The provider this attempt targets — the row owns routing until resolved.
pub provider_id: String,
/// Provider context (never secret; validated by `validate_provider_config`).
pub provider_config: serde_json::Value,
/// ISO timestamp, for degradation display.
pub created_at: String,
}
/// Durable deploy phase machine (§2.1, P10-C1).
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) enum DeployIntentPhase {
/// A deploy attempt owns this row; only the matching `attempt_id` may
/// clear/resolve it (P9-C3). Recovery of an orphaned `Running` row is a
/// residue clear or replay-as-recovery (§3.3 step 3).
Running { attempt_id: String },
/// The user explicitly accepted orphaning the possible remote (P10-C1).
/// Journaled BEFORE any record destruction; fenced by `consent_id`.
/// Recovery finishes the deletion — NEVER replays — from the row's own
/// durable data plus the pre-journaled library obligation `cleanup` points
/// to (P11-I1).
OrphanConsentPendingRemoval {
consent_id: String,
cleanup: ConsentCleanup,
},
}
/// The full per-pubkey cleanup policy a consented removal must durably produce,
/// captured inline from still-intact canonical state at consent time (§2.1,
/// P11-I1) — never re-derived from post-crash disk.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct ConsentCleanup {
/// A 30177 tombstone must be durably retained for this pubkey.
pub tombstone_required: bool,
/// NIP-IA archive retained now (plain / unprotected pubkey).
pub archive_required: bool,
/// `library_id` of EVERY entry protecting this pubkey (§2.5 predicate) whose
/// immediate archive was skipped — each row `DeferredArchive` journaled on
/// its entry BEFORE this consent transition (obligation-first, universal —
/// plain carriers included; P14-I1). Empty = no protection applied.
pub deferred_archive_entries: Vec<String>,
/// Exact cleanup coordinate for a projected instance; `None` for plain.
pub projected: Option<ProjectedCleanupRef>,
}
/// The pre-journaled library obligation a projected consent points to (§2.1).
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) struct ProjectedCleanupRef {
/// Owning entry.
pub library_id: String,
/// This scope's projection slug (the consenting scope and the projection's
/// scope are the same scope by construction — the journal's directory).
pub local_slug: String,
/// Which durable library obligation this consent points to, written BEFORE
/// the consent phase transition (§3.3 step 4 ordering) so recovery never
/// guesses.
pub obligation: ProjectedObligation,
}
/// The kind of projection obligation a consent owns (§2.1, P14-I1). Deferred
/// archives are NEVER named here — they live in
/// [`ConsentCleanup::deferred_archive_entries`] for plain and projected alike.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub(crate) enum ProjectedObligation {
/// Projection is `ExcludePending`/`DeletePending`: the pubkey was merged
/// into that `ProjectionEntry`'s `RemovalManifest` before consent.
ManifestMembership,
/// Projection is `Materialized` (direct forced delete): no removal manifest
/// exists and none is owed.
None,
}
/// Read outcome for `deploy-intents.json` (§2.1, P10-C3). `Unreadable` — unknown
/// version, syntax failure, OR a duplicate-pubkey integrity violation — fails
/// the owning scope's destructive and new-deploy paths CLOSED.
#[derive(Debug)]
pub(crate) enum DeployIntentsLoad {
/// Absent file — a valid empty journal.
Empty,
/// A healthy version-1 journal with at most one row per pubkey.
Loaded(DeployIntentsJournal),
Unreadable(String),
}
/// Read and validate `deploy-intents.json` under `definitions_dir` (§2.1). The
/// at-most-one-per-pubkey mutex is validated on read: any duplicate pubkey
/// group-rejects the WHOLE journal (never first-wins), because an ambiguous
/// journal may hide an intent equivalent to a live remote deployment (P10-C3).
pub(crate) fn load_deploy_intents(definitions_dir: &Path) -> DeployIntentsLoad {
let path = deploy_intents_path(definitions_dir);
let bytes = match std::fs::read(&path) {
Ok(bytes) => bytes,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return DeployIntentsLoad::Empty,
Err(e) => return DeployIntentsLoad::Unreadable(format!("read deploy-intents.json: {e}")),
};
let journal: DeployIntentsJournal = match serde_json::from_slice(&bytes) {
Ok(journal) => journal,
Err(e) => return DeployIntentsLoad::Unreadable(format!("parse deploy-intents.json: {e}")),
};
if journal.version != SUPPORTED_JOURNAL_VERSION {
return DeployIntentsLoad::Unreadable(format!(
"deploy-intents.json version {} is not supported (expected {SUPPORTED_JOURNAL_VERSION})",
journal.version
));
}
if let Some(pubkey) = duplicate_pubkey(&journal.intents) {
return DeployIntentsLoad::Unreadable(format!(
"deploy-intents.json has duplicate rows for {pubkey}: the at-most-one mutex is violated"
));
}
DeployIntentsLoad::Loaded(journal)
}
/// Persist `deploy-intents.json`: atomic temp-file + rename + `0o600`.
pub(crate) fn save_deploy_intents(
definitions_dir: &Path,
journal: &DeployIntentsJournal,
) -> Result<(), String> {
save_journal(&deploy_intents_path(definitions_dir), journal)
}
/// The first `agent_pubkey` that appears in more than one row, if any.
fn duplicate_pubkey(intents: &[DeployIntent]) -> Option<&str> {
let mut counts: HashMap<&str, usize> = HashMap::new();
for intent in intents {
*counts.entry(intent.agent_pubkey.as_str()).or_default() += 1;
}
intents
.iter()
.map(|intent| intent.agent_pubkey.as_str())
.find(|pubkey| counts[pubkey] > 1)
}
/// Shared atomic `0o600` writer for both journals: create the parent scope dir
/// if needed, then temp-file + rename via `storage.rs`.
fn save_journal<T: Serialize>(path: &Path, journal: &T) -> Result<(), String> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)
.map_err(|e| format!("create scope dir {}: {e}", parent.display()))?;
}
let payload = serde_json::to_vec_pretty(journal)
.map_err(|e| format!("serialize {}: {e}", path.display()))?;
atomic_write_json_restricted(path, &payload)
}
@@ -0,0 +1,669 @@
//! Phase-1 data-model tests (§7): the document envelope's quarantine-preserving
//! load/rebuild (P2-I2), per-entry decode + semantic + identity-index
//! quarantine (P4-I3/P6-I3), `apply_shared_definition` field isolation (§2.2),
//! and the scope-local journals' asymmetric failure domains (P9-I1/P10-C3).
use std::collections::BTreeMap;
use nostr::Keys;
use serde_json::{json, Value};
use tempfile::tempdir;
use super::journals::*;
use super::*;
use crate::managed_agents::types::AgentDefinition;
// ── fixtures ──────────────────────────────────────────────────────────────────
/// A blank [`AgentDefinition`] — the seed for `apply_shared_definition` field
/// tests. `AgentDefinition` derives no `Default`, and a full literal is robust
/// to future field additions in a way a `..Default::default()` shorthand can't
/// be (there is nothing to spread).
fn minimal_definition() -> AgentDefinition {
AgentDefinition {
id: "seed".to_string(),
display_name: "Seed".to_string(),
avatar_url: None,
system_prompt: "seed prompt".to_string(),
runtime: None,
model: None,
provider: None,
name_pool: Vec::new(),
is_builtin: false,
is_active: true,
shared: false,
source_team: None,
source_team_persona_slug: None,
catalog_source: None,
env_vars: BTreeMap::new(),
respond_to: None,
respond_to_allowlist: Vec::new(),
parallelism: None,
created_at: "2026-08-11T00:00:00Z".to_string(),
updated_at: "2026-08-11T00:00:00Z".to_string(),
}
}
/// A verified binding under `owner`, plus the (owner, agent) hex pair.
fn binding(owner: &Keys, agent: &Keys) -> (String, IdentityBinding, String) {
let owner_hex = owner.public_key().to_hex();
let agent_hex = agent.public_key().to_hex();
let auth_tag = buzz_sdk_pkg::nip_oa::compute_auth_tag(owner, &agent.public_key(), "")
.expect("compute auth tag");
(
owner_hex,
IdentityBinding {
agent_pubkey: agent_hex.clone(),
auth_tag,
},
agent_hex,
)
}
/// A minimal healthy entry: `library_id`, origin `(scope, slug)`, one
/// `Materialized` projection in `scope`, and (optionally) one verified binding.
fn entry(library_id: &str, scope: &str, slug: &str, bind: Option<(&Keys, &Keys)>) -> LibraryEntry {
let mut identity_bindings = BTreeMap::new();
let mut owner_at_share = "00".repeat(32);
if let Some((owner, agent)) = bind {
let (owner_hex, b, _) = binding(owner, agent);
owner_at_share = owner_hex.clone();
identity_bindings.insert(owner_hex, b);
}
let mut projections = BTreeMap::new();
projections.insert(
scope.to_string(),
ProjectionEntry {
state: ProjectionState::Materialized { revision: 1 },
local_slug: slug.to_string(),
relay_url: "wss://relay.example".to_string(),
workspace_label: Some("Home".to_string()),
},
);
LibraryEntry {
library_id: library_id.to_string(),
origin: OriginKey {
scope_id: scope.to_string(),
slug: slug.to_string(),
},
revision: 1,
deleted: false,
owner_pubkey_at_share: owner_at_share,
shared: SharedDefinition {
display_name: "Aria".to_string(),
avatar_url: None,
system_prompt: "be helpful".to_string(),
runtime: Some("acp".to_string()),
model: None,
provider: None,
name_pool: vec!["Aria".to_string()],
respond_to: Some("mentions".to_string()),
respond_to_allowlist: vec![],
parallelism: Some(1),
},
deferred_archives: vec![],
identity_bindings,
projections,
}
}
fn value_of(entry: &LibraryEntry) -> Value {
serde_json::to_value(entry).expect("serialize entry")
}
/// A projection claiming `local_slug` in some scope, in the given state.
fn proj(state: ProjectionState, local_slug: &str) -> ProjectionEntry {
ProjectionEntry {
state,
local_slug: local_slug.to_string(),
relay_url: "wss://relay.example".to_string(),
workspace_label: None,
}
}
fn doc(entries: Vec<Value>) -> LibraryDocument {
LibraryDocument {
version: SUPPORTED_LIBRARY_VERSION,
entries,
orphan_keys: vec![],
}
}
fn write_doc(base: &std::path::Path, document: &LibraryDocument) {
save_library_document(base, document).expect("save library.json");
}
// ── envelope: absent / version / corruption ────────────────────────────────────
#[test]
fn test_absent_library_loads_as_empty() {
let dir = tempdir().unwrap();
assert!(matches!(load_library(dir.path()), LibraryLoad::Empty));
}
#[test]
fn test_unknown_version_is_read_only_fail() {
let dir = tempdir().unwrap();
let mut document = doc(vec![]);
document.version = 999;
write_doc(dir.path(), &document);
match load_library(dir.path()) {
LibraryLoad::UnknownVersion(v) => assert_eq!(v, 999),
other => panic!("expected UnknownVersion, got {other:?}"),
}
// Read-only fail preserves the file untouched — no `.invalid` snapshot.
assert!(!library_path(dir.path())
.with_extension("json.invalid")
.exists());
}
#[test]
fn test_whole_document_corruption_is_preserved_and_blocks() {
let dir = tempdir().unwrap();
let path = library_path(dir.path());
std::fs::write(&path, b"{ this is not json").unwrap();
assert!(matches!(load_library(dir.path()), LibraryLoad::Corrupt));
// Loud-failure discipline: the malformed bytes survive as `.invalid`.
let invalid = path.with_extension("json.invalid");
assert_eq!(std::fs::read(&invalid).unwrap(), b"{ this is not json");
// And the original is left in place (copy, not rename).
assert!(path.exists());
}
#[test]
fn test_saved_document_round_trips_at_0600() {
let dir = tempdir().unwrap();
let owner = Keys::generate();
let agent = Keys::generate();
let document = doc(vec![value_of(&entry(
"lib-a",
"scope-a",
"aria",
Some((&owner, &agent)),
))]);
write_doc(dir.path(), &document);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
let mode = std::fs::metadata(library_path(dir.path()))
.unwrap()
.permissions()
.mode();
assert_eq!(mode & 0o777, 0o600);
}
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert_eq!(loaded.healthy.len(), 1);
assert!(loaded.quarantined.is_empty());
assert_eq!(loaded.healthy[0].library_id, "lib-a");
}
other => panic!("expected Loaded, got {other:?}"),
}
}
// ── quarantine: decode / semantic / identity index ─────────────────────────────
#[test]
fn test_unknown_field_entry_quarantines_while_sibling_stays_healthy() {
let dir = tempdir().unwrap();
let owner = Keys::generate();
let agent = Keys::generate();
let mut bad = value_of(&entry("lib-bad", "scope-b", "bob", None));
bad.as_object_mut()
.unwrap()
.insert("surprise".to_string(), json!("unexpected"));
let good = value_of(&entry(
"lib-good",
"scope-g",
"gwen",
Some((&owner, &agent)),
));
write_doc(dir.path(), &doc(vec![bad.clone(), good]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert_eq!(loaded.healthy.len(), 1);
assert_eq!(loaded.healthy[0].library_id, "lib-good");
assert_eq!(loaded.quarantined, vec![bad]);
assert_eq!(loaded.degradations.len(), 1);
}
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_binding_with_wrong_owner_key_quarantines() {
let dir = tempdir().unwrap();
let owner = Keys::generate();
let other_owner = Keys::generate();
let agent = Keys::generate();
// Tag is valid under `owner`, but keyed under `other_owner` — owner mismatch.
let mut e = entry("lib-mismatch", "scope-m", "mia", None);
let (_, b, _) = binding(&owner, &agent);
e.identity_bindings
.insert(other_owner.public_key().to_hex(), b);
let raw = value_of(&e);
write_doc(dir.path(), &doc(vec![raw.clone()]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert!(loaded.healthy.is_empty());
assert_eq!(loaded.quarantined, vec![raw]);
}
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_binding_with_wrong_agent_pubkey_quarantines() {
let dir = tempdir().unwrap();
let owner = Keys::generate();
let agent = Keys::generate();
let wrong_agent = Keys::generate();
// Tag authorizes `agent`, but the binding claims `wrong_agent`.
let owner_hex = owner.public_key().to_hex();
let auth_tag = buzz_sdk_pkg::nip_oa::compute_auth_tag(&owner, &agent.public_key(), "").unwrap();
let mut e = entry("lib-wrongagent", "scope-w", "wes", None);
e.identity_bindings.insert(
owner_hex,
IdentityBinding {
agent_pubkey: wrong_agent.public_key().to_hex(),
auth_tag,
},
);
let raw = value_of(&e);
write_doc(dir.path(), &doc(vec![raw.clone()]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert!(loaded.healthy.is_empty());
assert_eq!(loaded.quarantined, vec![raw]);
}
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_one_agent_bound_under_two_owners_quarantines() {
let dir = tempdir().unwrap();
let owner_a = Keys::generate();
let owner_b = Keys::generate();
let agent = Keys::generate();
let mut e = entry("lib-two-owners", "scope-t", "tom", None);
let (owner_a_hex, ba, _) = binding(&owner_a, &agent);
let (owner_b_hex, bb, _) = binding(&owner_b, &agent);
e.identity_bindings.insert(owner_a_hex, ba);
e.identity_bindings.insert(owner_b_hex, bb);
let raw = value_of(&e);
write_doc(dir.path(), &doc(vec![raw.clone()]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert!(loaded.healthy.is_empty());
assert_eq!(loaded.quarantined, vec![raw]);
}
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_duplicate_library_id_group_quarantines_all_colliders() {
let dir = tempdir().unwrap();
let a = value_of(&entry("lib-dup", "scope-1", "one", None));
let b = value_of(&entry("lib-dup", "scope-2", "two", None));
let unrelated = value_of(&entry("lib-solo", "scope-3", "three", None));
write_doc(dir.path(), &doc(vec![a.clone(), b.clone(), unrelated]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert_eq!(loaded.healthy.len(), 1);
assert_eq!(loaded.healthy[0].library_id, "lib-solo");
assert_eq!(loaded.quarantined.len(), 2);
assert!(loaded.quarantined.contains(&a));
assert!(loaded.quarantined.contains(&b));
}
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_duplicate_live_origin_group_quarantines() {
let dir = tempdir().unwrap();
// Distinct library_ids, same LIVE origin — share resumption is ambiguous.
let mut a = entry("lib-x", "scope-shared", "sameslug", None);
let mut b = entry("lib-y", "scope-shared", "sameslug", None);
// Give them distinct non-terminal projections so ONLY the origin collides.
a.projections.clear();
b.projections.clear();
let a = value_of(&a);
let b = value_of(&b);
let unrelated = value_of(&entry("lib-z", "scope-else", "elsewhere", None));
write_doc(dir.path(), &doc(vec![a.clone(), b.clone(), unrelated]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert_eq!(loaded.healthy.len(), 1);
assert_eq!(loaded.healthy[0].library_id, "lib-z");
assert_eq!(loaded.quarantined.len(), 2);
}
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_deleted_origin_does_not_collide() {
let dir = tempdir().unwrap();
// Same origin, but one entry is a tombstone — a `deleted` origin is not live,
// so there is no live-origin collision.
let mut live = entry("lib-live", "scope-o", "shared", None);
live.projections.clear();
let mut tomb = entry("lib-tomb", "scope-o", "shared", None);
tomb.deleted = true;
tomb.projections.clear();
write_doc(dir.path(), &doc(vec![value_of(&live), value_of(&tomb)]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => assert_eq!(loaded.healthy.len(), 2),
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_conflicting_non_terminal_scope_slug_group_quarantines() {
let dir = tempdir().unwrap();
// Distinct library_ids AND distinct origins, but both non-terminal
// projections claim `(scope-c, dup-slug)` — a scope would be rewritten from
// an arbitrary entry.
let mut a = entry("lib-p", "origin-a", "slug-a", None);
let mut b = entry("lib-q", "origin-b", "slug-b", None);
a.projections = BTreeMap::from([(
"scope-c".to_string(),
proj(ProjectionState::Materialized { revision: 1 }, "dup-slug"),
)]);
b.projections = BTreeMap::from([(
"scope-c".to_string(),
proj(ProjectionState::Pending, "dup-slug"),
)]);
write_doc(dir.path(), &doc(vec![value_of(&a), value_of(&b)]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => {
assert!(loaded.healthy.is_empty());
assert_eq!(loaded.quarantined.len(), 2);
}
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_terminal_projection_claims_do_not_collide() {
let dir = tempdir().unwrap();
// Same `(scope, slug)` claim but both terminal — terminal claims never
// conflict (the index only guards non-terminal ownership); distinct origins
// keep origin out of it.
let mut a = entry("lib-ta", "origin-ta", "slug-ta", None);
let mut b = entry("lib-tb", "origin-tb", "slug-tb", None);
a.projections = BTreeMap::from([(
"scope-c".to_string(),
proj(ProjectionState::Excluded, "dup-slug"),
)]);
b.projections = BTreeMap::from([(
"scope-c".to_string(),
proj(ProjectionState::Deleted, "dup-slug"),
)]);
write_doc(dir.path(), &doc(vec![value_of(&a), value_of(&b)]));
match load_library(dir.path()) {
LibraryLoad::Loaded(loaded) => assert_eq!(loaded.healthy.len(), 2),
other => panic!("expected Loaded, got {other:?}"),
}
}
// ── P2-I2 acceptance: preservation across a healthy rewrite ─────────────────────
#[test]
fn test_quarantined_entry_preserved_value_equivalently_across_edit_and_restart() {
let dir = tempdir().unwrap();
let owner = Keys::generate();
let agent = Keys::generate();
let mut bad = value_of(&entry("lib-quar", "scope-q", "quinn", None));
bad.as_object_mut()
.unwrap()
.insert("mystery".to_string(), json!({"nested": [1, 2, 3]}));
let good = entry("lib-edit", "scope-e", "evan", Some((&owner, &agent)));
write_doc(dir.path(), &doc(vec![bad.clone(), value_of(&good)]));
// Load, edit the healthy entry, rebuild, and persist — the quarantined raw
// must ride through unchanged.
let mut loaded = match load_library(dir.path()) {
LibraryLoad::Loaded(l) => l,
other => panic!("expected Loaded, got {other:?}"),
};
loaded.healthy[0].revision = 2;
let rebuilt = loaded.rebuild_document().expect("rebuild");
write_doc(dir.path(), &rebuilt);
match load_library(dir.path()) {
LibraryLoad::Loaded(l) => {
assert_eq!(l.healthy.len(), 1);
assert_eq!(l.healthy[0].revision, 2);
assert_eq!(l.quarantined, vec![bad]);
}
other => panic!("expected Loaded, got {other:?}"),
}
}
// ── §2.2: apply_shared_definition writes only shared slots ──────────────────────
#[test]
fn test_apply_shared_definition_writes_only_shared_slots() {
let mut record = minimal_definition().into_agent_record();
record.pubkey = "agent-pubkey".to_string();
record.private_key_nsec = "nsec-secret".to_string();
record.is_active = true;
record.env_vars = BTreeMap::from([("SECRET".to_string(), "value".to_string())]);
record.library_ref = Some("lib-existing".to_string());
record.backend_agent_id = Some("backend-123".to_string());
record.last_completed_deploy_attempt_id = Some("attempt-9".to_string());
let shared = SharedDefinition {
display_name: "Renamed".to_string(),
avatar_url: Some("https://a.example/x.png".to_string()),
system_prompt: "new prompt".to_string(),
runtime: Some("acp".to_string()),
model: Some("gpt".to_string()),
provider: Some("openai".to_string()),
name_pool: vec!["Renamed".to_string()],
respond_to: Some("all".to_string()),
respond_to_allowlist: vec!["npub1".to_string()],
parallelism: Some(3),
};
apply_shared_definition(&mut record, &shared, 7);
// Shared slots overwritten.
assert_eq!(record.name, "Renamed");
assert_eq!(record.display_name.as_deref(), Some("Renamed"));
assert_eq!(record.system_prompt.as_deref(), Some("new prompt"));
assert_eq!(record.model.as_deref(), Some("gpt"));
assert_eq!(record.definition_respond_to.as_deref(), Some("all"));
assert_eq!(record.definition_parallelism, Some(3));
assert_eq!(record.library_applied_revision, Some(7));
// Identity, secrets, activation, linkage, deploy provenance untouched.
assert_eq!(record.pubkey, "agent-pubkey");
assert_eq!(record.private_key_nsec, "nsec-secret");
assert!(record.is_active);
assert_eq!(
record.env_vars,
BTreeMap::from([("SECRET".to_string(), "value".to_string())])
);
assert_eq!(record.library_ref.as_deref(), Some("lib-existing"));
assert_eq!(record.backend_agent_id.as_deref(), Some("backend-123"));
assert_eq!(
record.last_completed_deploy_attempt_id.as_deref(),
Some("attempt-9")
);
}
#[test]
fn test_apply_shared_definition_empty_prompt_maps_to_none() {
let mut record = minimal_definition().into_agent_record();
let shared = SharedDefinition {
display_name: "N".to_string(),
avatar_url: None,
system_prompt: String::new(),
runtime: None,
model: None,
provider: None,
name_pool: vec![],
respond_to: None,
respond_to_allowlist: vec![],
parallelism: None,
};
apply_shared_definition(&mut record, &shared, 1);
// Mirrors `into_agent_record`: an empty prompt becomes `None`, not `Some("")`.
assert_eq!(record.system_prompt, None);
}
// ── scope-local journals (P9-I1 / P10-C3) ───────────────────────────────────────
#[test]
fn test_absent_journals_load_as_empty() {
let dir = tempdir().unwrap();
assert!(matches!(
load_pending_keys(dir.path()),
PendingKeysLoad::Empty
));
assert!(matches!(
load_deploy_intents(dir.path()),
DeployIntentsLoad::Empty
));
}
#[test]
fn test_pending_keys_round_trip() {
let dir = tempdir().unwrap();
let journal = PendingKeysJournal {
version: SUPPORTED_JOURNAL_VERSION,
pending: vec!["pk1".to_string(), "pk2".to_string()],
};
save_pending_keys(dir.path(), &journal).unwrap();
match load_pending_keys(dir.path()) {
PendingKeysLoad::Loaded(j) => assert_eq!(j, journal),
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_pending_keys_unknown_version_is_unreadable() {
let dir = tempdir().unwrap();
std::fs::write(
pending_keys_path(dir.path()),
br#"{"version":2,"pending":[]}"#,
)
.unwrap();
assert!(matches!(
load_pending_keys(dir.path()),
PendingKeysLoad::Unreadable(_)
));
}
#[test]
fn test_deploy_intents_round_trip_preserves_phase() {
let dir = tempdir().unwrap();
let journal = DeployIntentsJournal {
version: SUPPORTED_JOURNAL_VERSION,
intents: vec![
DeployIntent {
agent_pubkey: "pk-run".to_string(),
phase: DeployIntentPhase::Running {
attempt_id: "attempt-1".to_string(),
},
provider_id: "fly".to_string(),
provider_config: json!({"region": "iad"}),
created_at: "2026-08-11T00:00:00Z".to_string(),
},
DeployIntent {
agent_pubkey: "pk-consent".to_string(),
phase: DeployIntentPhase::OrphanConsentPendingRemoval {
consent_id: "consent-1".to_string(),
cleanup: ConsentCleanup {
tombstone_required: true,
archive_required: false,
deferred_archive_entries: vec!["lib-a".to_string()],
projected: Some(ProjectedCleanupRef {
library_id: "lib-a".to_string(),
local_slug: "lib-a-slug".to_string(),
obligation: ProjectedObligation::ManifestMembership,
}),
},
},
provider_id: "fly".to_string(),
provider_config: json!(null),
created_at: "2026-08-11T00:01:00Z".to_string(),
},
],
};
save_deploy_intents(dir.path(), &journal).unwrap();
match load_deploy_intents(dir.path()) {
DeployIntentsLoad::Loaded(j) => assert_eq!(j, journal),
other => panic!("expected Loaded, got {other:?}"),
}
}
#[test]
fn test_deploy_intents_duplicate_pubkey_group_rejects_whole_journal() {
let dir = tempdir().unwrap();
let row = |pubkey: &str| DeployIntent {
agent_pubkey: pubkey.to_string(),
phase: DeployIntentPhase::Running {
attempt_id: "a".to_string(),
},
provider_id: "fly".to_string(),
provider_config: json!({}),
created_at: "t".to_string(),
};
let journal = DeployIntentsJournal {
version: SUPPORTED_JOURNAL_VERSION,
intents: vec![row("dup"), row("dup"), row("other")],
};
save_deploy_intents(dir.path(), &journal).unwrap();
// The at-most-one mutex is validated on read: any duplicate rejects the
// WHOLE journal, never first-wins (P10-C3).
match load_deploy_intents(dir.path()) {
DeployIntentsLoad::Unreadable(reason) => assert!(reason.contains("dup")),
other => panic!("expected Unreadable, got {other:?}"),
}
}
#[test]
fn test_deploy_intents_unknown_field_is_unreadable() {
let dir = tempdir().unwrap();
std::fs::write(
deploy_intents_path(dir.path()),
br#"{"version":1,"intents":[],"surprise":true}"#,
)
.unwrap();
assert!(matches!(
load_deploy_intents(dir.path()),
DeployIntentsLoad::Unreadable(_)
));
}
#[test]
fn test_journals_written_at_0600() {
let dir = tempdir().unwrap();
save_pending_keys(dir.path(), &PendingKeysJournal::empty()).unwrap();
save_deploy_intents(dir.path(), &DeployIntentsJournal::empty()).unwrap();
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
for path in [
pending_keys_path(dir.path()),
deploy_intents_path(dir.path()),
] {
let mode = std::fs::metadata(&path).unwrap().permissions().mode();
assert_eq!(mode & 0o777, 0o600, "{}", path.display());
}
}
}
@@ -16,6 +16,7 @@ pub(crate) mod effective_config;
mod env_vars;
pub(crate) mod git_bash;
pub(crate) mod global_config;
pub(crate) mod library;
mod managed_node_paths;
mod nest;
pub(crate) mod parallelism;
@@ -500,6 +500,7 @@ fn make_agent(name: &str, persona_id: Option<&str>) -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -115,6 +115,7 @@ mod tests {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -56,6 +56,7 @@ pub(super) fn sample_record() -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -1526,6 +1526,7 @@ mod tests {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -1550,8 +1551,7 @@ mod tests {
#[test]
fn buzz_agent_databricks_v2_with_databricks_model_but_no_buzz_agent_model_is_ready() {
// The baked buzz-releases env sets DATABRICKS_MODEL but not BUZZ_AGENT_MODEL.
// An agent with only DATABRICKS_MODEL must pass the readiness gate.
// Baked buzz-releases env has DATABRICKS_MODEL but no BUZZ_AGENT_MODEL; an agent with only DATABRICKS_MODEL must pass the gate.
let env = make_env(
"buzz-agent",
env_with(&[
@@ -87,6 +87,7 @@ pub(super) fn fixture(
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -89,6 +89,7 @@ fn test_compensate_drain_writer_vs_compensation_deterministic() {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Default::default(),
definition_parallelism: None,
@@ -339,6 +340,7 @@ fn test_compensate_drain_concurrent_start_is_blocked() {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Default::default(),
definition_parallelism: None,
@@ -68,6 +68,7 @@ fn record() -> ManagedAgentRecord {
catalog_source: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: Vec::new(),
definition_parallelism: None,
@@ -306,6 +306,7 @@ mod tests {
source_team_persona_slug: Some("SENTINEL_SLUG".to_string()), // MUST NOT appear
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
catalog_source: None,
definition_respond_to_allowlist: vec![],
@@ -215,6 +215,7 @@ fn managed_agent(name: &str) -> ManagedAgentRecord {
relay_mesh: None,
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: None,
definition_respond_to_allowlist: vec![],
definition_parallelism: None,
@@ -324,6 +324,20 @@ pub struct ManagedAgentRecord {
/// through a persona save.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub library_applied_revision: Option<u64>,
/// Deploy-attempt provenance, not projection metadata (§2.6). Stamps the
/// `attempt_id` of the last provider deploy that durably landed a
/// `backend_agent_id` — written ONLY alongside `backend_agent_id` in the
/// same success write (§3.3 step 2, P12-I1). A failed or ambiguous attempt
/// never touches it, so an older stamp survives across failures and a
/// `Running` deploy-intent row whose `attempt_id` mismatches this stamp
/// still routes to replay. A legacy record reads as unstamped (`None`),
/// which recovery already treats as "no residue proof → replay". Forms one
/// inseparable deploy-provenance pair with `backend_agent_id` in every
/// copy/rollback/normalization helper (`copy_runtime_state`, P14-I2).
/// `#[serde(default, skip_serializing_if)]` keeps records without it
/// byte-identical to head (invariant 4).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub last_completed_deploy_attempt_id: Option<String>,
/// NIP-AP definition-level behavioral defaults, absorbed from
/// `AgentDefinition` in WIRE shape (kebab-case string / optional u32),
/// distinct from the instance-side `respond_to`/`respond_to_allowlist`/
@@ -72,6 +72,7 @@ impl AgentDefinition {
// a freshly projected definition carries none.
library_ref: None,
library_applied_revision: None,
last_completed_deploy_attempt_id: None,
definition_respond_to: self.respond_to,
definition_respond_to_allowlist: self.respond_to_allowlist,
definition_parallelism: self.parallelism,