feat(desktop): gate owner-identity egress behind witnessed leases (WSA C1)

Route every owner and managed-agent relay send through the
owner_identity_egress admission layer so no code path can publish
under an owner identity without a rate-limit-witnessed EgressLease.

Privatize AppState.keys behind latch-gated checked accessors
(signing_keys/current_pubkey), thread the EgressLease witness through
all six explicit-key funnels plus the four git-workflow and four
managed-agent sink sites, and rewrite the huddle STT task to
re-resolve keys per send with a mid-huddle recovery-latch break.

The C5 P25/P28 coordinator (journal + three-valued commit outcome +
Indeterminate latch) and its egress-drain wiring stay deferred: each
unreachable item carries a per-item allow(dead_code) naming C5 as its
consumer, with semantics pinned by the owner_identity_egress unit
tests until C5 makes them reachable from production paths.

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-17 22:24:35 -04:00
co-authored by Will Pfleger
parent 159898d8f6
commit 59800a0fb6
37 changed files with 1216 additions and 159 deletions
+47 -1
View File
@@ -17,7 +17,12 @@ use crate::managed_agents::scope::WorkspaceAgentScope;
use crate::managed_agents::{ManagedAgentPairRuntime, ManagedAgentRuntimeKey};
pub struct AppState {
pub keys: Mutex<Keys>,
/// Identity signing keys. PRIVATE (P29-C1): reach through
/// [`AppState::signing_keys`] (refuses in recovery / under the latch),
/// [`AppState::current_pubkey`] (pubkey-only), or
/// [`AppState::identity_lifecycle_keys_guard`] (transition/import commit +
/// lost-state persist only). No signing path may touch the raw field.
keys: Mutex<Keys>,
/// Durable backend holding `keys`. Updated after the key write and before
/// recovery flags are cleared so `get_identity` reports a consistent state.
pub(crate) identity_storage: AtomicU8,
@@ -267,12 +272,53 @@ impl AppState {
until the identity is restored and Buzz is relaunched"
.to_string());
}
// P29-C1: also refuse while owner-identity persistence is latched
// `Indeterminate`. A completed transition that could not prove either
// durable identity canonical must not sign under the unresolved
// identity. The latch is only set by the C5 identity-transition
// coordinator, so this is inert (never `true`) until C5 lands, but the
// gate is in place at every signer now.
if crate::owner_identity_egress::is_identity_indeterminate() {
return Err("owner identity is in an indeterminate recovery state; \
event signing is disabled until the identity is reconciled \
and Buzz is relaunched"
.to_string());
}
self.keys
.lock()
.map_err(|e| e.to_string())
.map(|k| k.clone())
}
/// The current identity's public key, for routing, query filters, and
/// display. Unlike [`signing_keys`](Self::signing_keys) this does NOT
/// refuse in recovery: a public key is not signing capability, and the
/// recovery-reporting surfaces (`get_identity`) need it to describe the
/// recovery state itself. Reads through this accessor rather than the
/// private `keys` field so no caller can reach the secret key for a
/// pubkey-only need.
pub fn current_pubkey(&self) -> Result<nostr::PublicKey, String> {
self.keys
.lock()
.map_err(|e| e.to_string())
.map(|k| k.public_key())
}
/// Raw guard on the identity keys for the identity-lifecycle paths ONLY:
/// the `apply_workspace` and `import_identity` commit stages (which swap
/// keys under their transition guards) and `persist_current_identity`
/// (which clones the ephemeral lost-state key to make it durable). These
/// paths legitimately operate on keys DURING recovery, so they cannot go
/// through [`signing_keys`](Self::signing_keys). Signing and publishing
/// MUST use [`signing_keys`](Self::signing_keys), never this — the field
/// is private so this is the only write door, and it is named to make a
/// signing misuse obvious in review.
pub(crate) fn identity_lifecycle_keys_guard(
&self,
) -> std::sync::LockResult<std::sync::MutexGuard<'_, Keys>> {
self.keys.lock()
}
/// Emit the current huddle state to the frontend via Tauri event.
///
/// Acquires both locks (app_handle + huddle_state), clones a snapshot,
+2 -7
View File
@@ -46,8 +46,7 @@ fn open_db() -> Result<Connection, String> {
}
fn identity_pubkey(state: &AppState) -> Result<String, String> {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
Ok(keys.public_key().to_hex())
Ok(state.current_pubkey()?.to_hex())
}
fn now_secs() -> i64 {
@@ -156,11 +155,7 @@ pub async fn archive_events(
let bucket_results = query_buckets(plan.buckets, state_ref).await;
// ── Phase 3: persist (blocking SQLite) ──────────────────────────────────
let owner_keys = {
let keys_guard = state.keys.lock().map_err(|e| e.to_string())?;
keys_guard.clone()
// guard drops here, before awaiting the blocking commit task.
};
let owner_keys = state.signing_keys()?;
let commit_identity_pk = identity_pk.clone();
let commit_relay_url = relay_url.clone();
run_archive_db_task(move |conn| {
+3 -3
View File
@@ -839,7 +839,7 @@ mod real_relay {
/// is exercised, including NIP-98 signing inside `query_relay`.
fn make_test_app_state(keys: Keys, relay_url: &str) -> AppState {
let state = build_app_state();
*state.keys.lock().unwrap() = keys;
*state.identity_lifecycle_keys_guard().unwrap() = keys;
*state.relay_url_override.lock().unwrap() = Some(relay_url.to_string());
state
}
@@ -920,7 +920,7 @@ mod real_relay {
state: &AppState,
db_path: &Path,
) -> ArchiveBatchResult {
let identity_pk = state.keys.lock().unwrap().public_key().to_hex();
let identity_pk = state.current_pubkey().unwrap().to_hex();
let relay_url = crate::relay::relay_ws_url_with_override(state);
// Phase 1: plan (sync). Connection dropped before any .await.
@@ -936,7 +936,7 @@ mod real_relay {
// Phase 3: persist (sync). Fresh connection, same file.
let conn = store::open_archive_db(db_path).expect("open archive db for commit");
let owner_keys = state.keys.lock().unwrap().clone();
let owner_keys = state.identity_lifecycle_keys_guard().unwrap().clone();
commit_archive(
bucket_results,
plan.ephemeral,
+1 -2
View File
@@ -21,8 +21,7 @@ use crate::{
/// Read the workspace owner pubkey without holding the lock. Used to populate `BUZZ_ACP_AGENT_OWNER`
/// as a fallback for legacy agent records that have no NIP-OA `auth_tag`.
pub(super) fn workspace_owner_hex(state: &AppState) -> Result<String, String> {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
Ok(keys.public_key().to_hex())
Ok(state.current_pubkey()?.to_hex())
}
/// Retain a freshly authored managed-agent event in the local store, flagged
+13 -7
View File
@@ -78,10 +78,7 @@ fn classify_pending_owner(state: &AppState, my_pubkey: &str, d_tag: Option<&str>
#[tauri::command]
pub async fn get_channels(state: State<'_, AppState>) -> Result<Vec<ChannelInfo>, String> {
let _profile_start = std::time::Instant::now();
let my_pubkey = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
keys.public_key().to_hex()
};
let my_pubkey = state.current_pubkey()?.to_hex();
// Step 1: find all kind:39002 (members) events that mention me, then
// pull the channel ids out of their `d` tags.
@@ -511,8 +508,11 @@ async fn ensure_starter_channel_memberships(
}
let channel_uuid = parse_channel_uuid(&channel.id)?;
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let builder = events::build_join(channel_uuid)?;
submit_event_with_keys(builder, state, keys, None).await?;
submit_event_with_keys(builder, state, keys, None, &lease).await?;
channel.is_member = true;
}
@@ -579,7 +579,10 @@ pub async fn create_channel(
// able to retarget the mark onto the new identity.
let creator_keys = state.signing_keys()?;
let creator_pubkey = creator_keys.public_key().to_hex();
submit_event_with_keys(builder, &state, &creator_keys, None).await?;
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
submit_event_with_keys(builder, &state, &creator_keys, None, &lease).await?;
// Mark this channel pending-owner: we just created it, so we know we're
// the owner, but the relay's kind:39002 membership entry (#1761) is
@@ -630,6 +633,9 @@ pub async fn ensure_starter_channels(
let channel_uuid = starter_channel_uuid(&relay_scope, spec.slug);
let channel_uuid_string = channel_uuid.to_string();
starter_ids.push(channel_uuid_string.clone());
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let builder = events::build_create_channel(
channel_uuid,
spec.name,
@@ -639,7 +645,7 @@ pub async fn ensure_starter_channels(
None,
)?;
match submit_event_with_keys(builder, &state, &creator_keys, None).await {
match submit_event_with_keys(builder, &state, &creator_keys, None, &lease).await {
Ok(_) => {
state.mark_pending_owned_channel(&creator_pubkey, &channel_uuid_string);
created_ids.insert(channel_uuid_string.clone());
@@ -230,14 +230,14 @@ fn pending_owner_mark_uses_signer_captured_before_identity_swap() {
// Simulate an in-process identity swap landing during the (here,
// implicit) submit await — e.g. `import_identity` replacing
// `state.keys` while the create request is in flight.
*state.keys.lock().expect("lock keys") = Keys::generate();
*state.identity_lifecycle_keys_guard().expect("lock keys") = Keys::generate();
// The mark must use the captured signer, not whatever `state.keys`
// holds now.
state.mark_pending_owned_channel(&creator_pubkey, "chan-1");
assert!(state.is_pending_owned_channel(&creator_pubkey, "chan-1"));
let post_swap_pubkey = state.keys.lock().expect("lock keys").public_key().to_hex();
let post_swap_pubkey = state.current_pubkey().expect("pubkey").to_hex();
assert!(!state.is_pending_owned_channel(&post_swap_pubkey, "chan-1"));
}
+4 -7
View File
@@ -137,10 +137,7 @@ pub async fn get_agent_memory(
let agent = PublicKey::from_hex(&agent_pubkey)
.map_err(|e| format!("agent pubkey must be 64-hex: {e}"))?;
let viewer_pubkey = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
keys.public_key().to_hex()
};
let viewer_pubkey = state.current_pubkey()?.to_hex();
let managed = load_managed_agents(&app)?;
let is_managed = managed.iter().any(|m| m.pubkey == agent_pubkey);
@@ -160,10 +157,10 @@ pub async fn get_agent_memory(
}
// ── Resolve owner key material ──────────────────────────────────────
// Owner = viewer. Clone the secret key out of the lock immediately so
// we don't hold the mutex across the relay round trip.
// Owner = viewer. `signing_keys()` refuses in recovery / under the latch,
// so a NIP-44 encrypt-to-self here cannot run under an unresolved identity.
let (owner_pubkey, owner_seckey) = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
let keys = state.signing_keys()?;
(keys.public_key(), keys.secret_key().clone())
};
+16 -10
View File
@@ -26,8 +26,7 @@ fn truncated_display_name(pubkey: &PublicKey) -> Result<String, String> {
#[tauri::command]
pub fn get_identity(state: State<'_, AppState>) -> Result<IdentityInfo, String> {
let keys = state.keys.lock().map_err(|error| error.to_string())?;
let pubkey = keys.public_key();
let pubkey = state.current_pubkey()?;
let pubkey_hex = pubkey.to_hex();
let display_name = truncated_display_name(&pubkey)?;
let lost = state
@@ -619,7 +618,10 @@ pub(crate) fn commit_imported_identity(
persist: impl FnOnce(&nostr::Keys) -> Result<crate::app_state::IdentityStorage, String>,
) -> Result<(nostr::PublicKey, crate::app_state::IdentityStorage), String> {
// Capture the previous pubkey up front for post-commit cleanup.
let previous_pubkey = state.keys.lock().map_err(|e| e.to_string())?.public_key();
let previous_pubkey = state
.identity_lifecycle_keys_guard()
.map_err(|e| e.to_string())?
.public_key();
let storage = persist(&keys)?;
@@ -628,7 +630,9 @@ pub(crate) fn commit_imported_identity(
// observing false is guaranteed to see the updated keys.
let pubkey = keys.public_key();
{
let mut active_keys = state.keys.lock().map_err(|e| e.to_string())?;
let mut active_keys = state
.identity_lifecycle_keys_guard()
.map_err(|e| e.to_string())?;
*active_keys = keys;
state.set_identity_storage(storage);
}
@@ -692,7 +696,13 @@ pub async fn persist_current_identity(
}
// Clone current keys without holding the mutex across keyring I/O.
let keys = state.keys.lock().map_err(|e| e.to_string())?.clone();
// Lost-state path: `signing_keys()` would refuse here, but this command
// exists to make the ephemeral lost-state key durable, so it reads the
// raw guard directly.
let keys = state
.identity_lifecycle_keys_guard()
.map_err(|e| e.to_string())?
.clone();
let data_dir = app_handle
.path()
@@ -824,11 +834,7 @@ pub async fn sign_nostr_identity_binding(
&expires_at,
)?;
let keys = state
.keys
.lock()
.map_err(|error| error.to_string())?
.clone();
let keys = state.signing_keys()?;
tauri::async_runtime::spawn_blocking(move || {
let event = build_nostr_identity_binding_event(
@@ -104,10 +104,7 @@ pub async fn resolve_oa_owner(
return Ok(None);
};
let my_pubkey = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
keys.public_key().to_hex()
};
let my_pubkey = state.current_pubkey()?.to_hex();
Ok(Some(OwnerOfAgent {
is_me: my_pubkey.eq_ignore_ascii_case(&owner_hex),
@@ -189,10 +186,7 @@ async fn maybe_owner_auth_tag(
state: &AppState,
target_pubkey: &str,
) -> Result<Option<[String; 4]>, String> {
let my_pubkey = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
keys.public_key().to_hex()
};
let my_pubkey = state.current_pubkey()?.to_hex();
// Self path: never attach auth (spec §Self Requests: if actor==target and
// an `auth` tag is also present, relay MUST treat it as self).
@@ -12,10 +12,7 @@ fn verification_returns_only_public_identity_and_match_status() {
let state = build_app_state();
let backup = create_backup_with_log_n(&state, PASSWORD, FAST_LOG_N).unwrap();
let result = verify_ncryptsec_backup_inner(&state, &backup, PASSWORD).unwrap();
assert_eq!(
result.pubkey,
state.keys.lock().unwrap().public_key().to_hex()
);
assert_eq!(result.pubkey, state.current_pubkey().unwrap().to_hex());
assert!(result.npub.starts_with("npub1"));
assert!(result.matches_current_identity);
}
@@ -112,7 +109,7 @@ fn recovery_mode_blocks_backup_creation() {
#[test]
fn concurrent_identity_swap_vs_backup_is_serialized() {
let state = std::sync::Arc::new(build_app_state());
let key_a = state.keys.lock().unwrap().clone();
let key_a = state.identity_lifecycle_keys_guard().unwrap().clone();
let key_b = Keys::generate();
let swapper = {
@@ -122,7 +119,7 @@ fn concurrent_identity_swap_vs_backup_is_serialized() {
// Mirrors import_identity's locking: mutation guard held
// across the key swap.
let _guard = state.identity_mutation.blocking_lock();
*state.keys.lock().unwrap() = key_b;
*state.identity_lifecycle_keys_guard().unwrap() = key_b;
})
};
+17 -11
View File
@@ -61,10 +61,7 @@ pub async fn get_feed(
.map(|t| t.split(',').any(|s| s.trim() == "needs_action"))
.unwrap_or(true);
let my_pubkey = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
keys.public_key().to_hex()
};
let my_pubkey = state.current_pubkey()?.to_hex();
// Mentions: messages that reference me via #p.
let mut mention_filter = serde_json::json!({
@@ -736,7 +733,7 @@ fn managed_agent_submission_auth_tag(
return Ok(Some(auth_tag));
}
let owner_keys = state.keys.lock().map_err(|error| error.to_string())?;
let owner_keys = state.signing_keys()?;
legacy_managed_agent_auth_tag(&owner_keys, agent_pubkey)
}
@@ -846,6 +843,12 @@ pub async fn send_managed_agent_channel_message(
}
}
let mentions = mention_pubkeys.unwrap_or_default();
// Managed-agent channel send (P29-C1 closed-world sink, widest blast
// radius). Admit BEFORE building the message so the kind:9 `created_at` is
// stamped after the rate-limit wait.
let lease = crate::owner_identity_egress::EgressLease::ManagedAgentKeyed(
crate::owner_identity_egress::admit_managed_agent_egress().await?,
);
let builder = build_managed_agent_channel_message(
channel_uuid,
trimmed,
@@ -853,8 +856,14 @@ pub async fn send_managed_agent_channel_message(
&mentions,
&client_tags,
)?;
let result =
submit_event_with_keys(builder, &state, &keys, submission_auth_tag.as_deref()).await?;
let result = submit_event_with_keys(
builder,
&state,
&keys,
submission_auth_tag.as_deref(),
&lease,
)
.await?;
Ok(SendChannelMessageResponse {
event_id: result.event_id,
@@ -892,10 +901,7 @@ pub async fn remove_reaction(
state: State<'_, AppState>,
) -> Result<(), String> {
// Find our own kind:7 reaction event referencing the target.
let my_pubkey = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
keys.public_key().to_hex()
};
let my_pubkey = state.current_pubkey()?.to_hex();
let target = event_id.trim();
let trimmed_emoji = emoji.trim();
@@ -105,7 +105,7 @@ fn frozen_linkage_corrective_failure_errs_and_boot_reconcile_converges() {
// Mock Tauri app whose active scope points its definitions dir at the
// seeded store; the signing keys are the event's author.
let state = crate::app_state::build_app_state();
*state.keys.lock().unwrap() = owner_keys.clone();
*state.identity_lifecycle_keys_guard().unwrap() = owner_keys.clone();
let app = tauri::test::mock_builder()
.manage(state)
.build(tauri::test::mock_context(tauri::test::noop_assets()))
@@ -95,11 +95,19 @@ async fn publish_prepared_persona(
prepared: PreparedPersonaPublication,
) -> Result<SetPersonaSharedResult, String> {
let api_base_url = crate::relay::relay_http_base_url(&prepared.scope.relay_url);
// The persona event was signed during preparation, so this is admit-before-
// submit (freshness is fixed at preparation time). The lease still gates the
// send under the identity-persistence latch and waits out the rate-limit
// gate before transmit.
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let publish_result = crate::relay::submit_signed_event_at_with_keys(
&prepared.event,
state,
&api_base_url,
&prepared.scope.owner_keys,
&lease,
)
.await;
@@ -870,11 +870,16 @@ pub(crate) async fn submit_engram_event(
crate::egress_guard::assert_no_key_backup_bytes(event_json, "persona snapshot engram submit")?;
// Wait before signing: the relay enforces NIP-98 freshness (±60s) and the
// gate may hold for up to MAX_HINT_SECONDS (300s). Building auth before the
// wait produces a stale `created_at` that the relay will reject.
crate::relay_admission::wait_for_rate_limit().await;
let auth = build_nip98_auth_header_for_keys(agent_keys, &Method::POST, url, event_json)?;
// Managed-agent egress construction site (P29-C1 closed-world sink). Admit
// the interim keyed-egress lease, which waits out the rate-limit gate then
// refuses under the identity-persistence latch/drain. The event is
// pre-signed by the caller (freshness is caller-determined), so this is
// admit-before-submit.
let lease = crate::owner_identity_egress::EgressLease::ManagedAgentKeyed(
crate::owner_identity_egress::admit_managed_agent_egress().await?,
);
let auth =
build_nip98_auth_header_for_keys(agent_keys, &Method::POST, url, event_json, &lease)?;
let mut request = state
.http_client
.post(url)
@@ -252,7 +252,7 @@ fn setup_import_app_with_scope(
let owner_keys = nostr::Keys::generate();
let state = crate::app_state::build_app_state();
{
let mut locked = state.keys.lock().unwrap();
let mut locked = state.identity_lifecycle_keys_guard().unwrap();
*locked = owner_keys.clone();
}
@@ -32,7 +32,7 @@ fn make_scope_with_keys(tmp: &tempfile::TempDir, owner_keys: &nostr::Keys) -> Wo
fn build_import_state(owner_keys: nostr::Keys) -> AppState {
let state = build_app_state();
{
let mut locked = state.keys.lock().unwrap();
let mut locked = state.identity_lifecycle_keys_guard().unwrap();
*locked = owner_keys;
}
state
+120 -6
View File
@@ -114,14 +114,19 @@ pub async fn update_profile_at_relay(
"authors": [expected_pubkey],
"limit": 1
});
let query_lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let prior_events = query_relay_at_with_keys(
&state,
&api_base_url,
std::slice::from_ref(&filter),
&signer,
None,
&query_lease,
)
.await?;
drop(query_lease);
let prior_event = prior_events.first();
let current: Value = prior_event
.and_then(|event| serde_json::from_str::<Value>(&event.content).ok())
@@ -136,10 +141,28 @@ pub async fn update_profile_at_relay(
return Err("profile avatar changed before deferred save".to_string());
}
// Admit BEFORE building the event: `build_deferred_profile_event` stamps
// `created_at` at build time, so admission (which waits out the rate-limit
// gate) must precede it to keep the kind:0 event fresh under a gate hold.
let submit_lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let builder = build_deferred_profile_event(&current, &avatar_url, prior_event)?;
submit_event_at_with_keys(builder, &state, &api_base_url, &signer).await?;
submit_event_at_with_keys(builder, &state, &api_base_url, &signer, &submit_lease).await?;
drop(submit_lease);
let events = query_relay_at_with_keys(&state, &api_base_url, &[filter], &signer, None).await?;
let requery_lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let events = query_relay_at_with_keys(
&state,
&api_base_url,
&[filter],
&signer,
None,
&requery_lease,
)
.await?;
Ok(events
.first()
.map(nostr_convert::profile_info_from_event)
@@ -394,8 +417,7 @@ pub async fn get_presence(
}
fn current_pubkey_hex(state: &AppState) -> Result<String, String> {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
Ok(keys.public_key().to_hex())
Ok(state.current_pubkey()?.to_hex())
}
fn current_pubkey_hex_unwrap(state: &AppState) -> String {
@@ -426,11 +448,11 @@ mod tests {
let captured = capture_expected_signer(&state, &original_pubkey)
.expect("matching identity should be captured");
*state.keys.lock().expect("lock keys") = nostr::Keys::generate();
*state.identity_lifecycle_keys_guard().expect("lock keys") = nostr::Keys::generate();
assert_eq!(captured.public_key().to_hex(), original_pubkey);
assert_ne!(
state.keys.lock().expect("lock keys").public_key().to_hex(),
state.current_pubkey().expect("pubkey").to_hex(),
original_pubkey
);
assert_eq!(
@@ -481,4 +503,96 @@ mod tests {
assert_eq!(filter["limit"], serde_json::json!(25));
assert_eq!(filter["page"], serde_json::json!(1));
}
/// End-to-end drive of `update_profile_at_relay`'s three-lease path
/// (query → submit → re-query) against a per-path counting loopback relay.
///
/// Proves the P29-C1 posture Paul gated C1 close on: the deferred profile
/// save admits exactly one owner-identity egress lease per network op and
/// submits the kind:0 event exactly once. It also exercises the C1c
/// accessor sweep on the live path — `capture_expected_signer` reads the
/// identity through `signing_keys()` (the latch-gated accessor), so a
/// regression that broke the accessor or the funnel witness threading would
/// fail here rather than only in unit isolation.
#[tokio::test]
async fn update_profile_at_relay_submits_exactly_once_across_three_leases() {
use std::io::{Read, Write};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
let _serial = crate::relay_admission::TEST_SERIAL.lock().await;
crate::relay_admission::reset_rate_limit_gate();
let query_count = Arc::new(AtomicUsize::new(0));
let submit_count = Arc::new(AtomicUsize::new(0));
let (qc, sc) = (Arc::clone(&query_count), Arc::clone(&submit_count));
// Loopback relay: the `/query` path answers `[]` (empty profile
// history), any other path (the submit endpoint) answers an accepted
// `SubmitEventResponse`. Both are counted so the test can assert the
// exactly-once submit + two queries of the deferred save flow.
let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let addr = listener.local_addr().unwrap();
let server = std::thread::spawn(move || {
for _ in 0..3 {
let Ok((mut stream, _)) = listener.accept() else {
break;
};
let mut buf = [0u8; 8192];
let n = stream.read(&mut buf).unwrap_or(0);
let request = String::from_utf8_lossy(&buf[..n]);
let request_line = request.lines().next().unwrap_or("");
let body = if request_line.contains("/query") {
qc.fetch_add(1, Ordering::SeqCst);
"[]".to_string()
} else {
sc.fetch_add(1, Ordering::SeqCst);
r#"{"event_id":"deadbeef","accepted":true,"message":"ok"}"#.to_string()
};
let _ = stream.write_all(
format!(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{body}",
body.len()
)
.as_bytes(),
);
let _ = stream.flush();
}
});
let app = tauri::test::mock_builder()
.manage(crate::app_state::build_app_state())
.build(tauri::test::mock_context(tauri::test::noop_assets()))
.expect("mock app");
use tauri::Manager;
let state = app.state::<crate::app_state::AppState>();
let expected_pubkey = state.signing_keys().unwrap().public_key().to_hex();
let result = update_profile_at_relay(
format!("http://{addr}"),
expected_pubkey.clone(),
None,
"https://example.com/avatar.png".to_string(),
state,
)
.await
.expect("deferred profile save must complete");
// Empty re-query returns the empty canonical profile for the identity.
assert_eq!(result.pubkey, expected_pubkey);
// Exactly-once posture: one prior-query, one submit, one re-query.
assert_eq!(
submit_count.load(Ordering::SeqCst),
1,
"kind:0 submitted once"
);
assert_eq!(
query_count.load(Ordering::SeqCst),
2,
"prior + re-query only"
);
server.join().unwrap();
crate::relay_admission::reset_rate_limit_gate();
}
}
@@ -102,6 +102,11 @@ fn normalize_event_id(value: &str) -> Option<String> {
struct ProjectOwnerIdentity {
keys: Keys,
auth_tag: Option<String>,
/// Whether the resolved identity is a managed agent's key (non-viewer
/// branch) rather than the human owner's (viewer branch). Set explicitly at
/// the branch that knows it — the P29-C1 egress lease variant is chosen from
/// this, never inferred from `auth_tag`.
is_managed_agent: bool,
}
fn project_owner_identity(
@@ -114,6 +119,7 @@ fn project_owner_identity(
return Ok(ProjectOwnerIdentity {
keys: viewer_keys,
auth_tag: None,
is_managed_agent: false,
});
}
@@ -140,9 +146,30 @@ fn project_owner_identity(
Ok(ProjectOwnerIdentity {
keys,
auth_tag: record.auth_tag.clone(),
is_managed_agent: true,
})
}
/// Admit the P29-C1 egress lease for a resolved [`ProjectOwnerIdentity`]:
/// the managed-agent keyed lease for the non-viewer branch, the owner-identity
/// lease for the viewer branch. Both wait out the rate-limit gate internally,
/// so callers must admit BEFORE building/signing the event they publish.
async fn admit_project_egress(
identity: &ProjectOwnerIdentity,
) -> Result<crate::owner_identity_egress::EgressLease, String> {
if identity.is_managed_agent {
Ok(
crate::owner_identity_egress::EgressLease::ManagedAgentKeyed(
crate::owner_identity_egress::admit_managed_agent_egress().await?,
),
)
} else {
Ok(crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
))
}
}
fn validate_repo_address(repo_address: &str, owner: &str) -> Result<(), String> {
let prefix = format!("30617:{owner}:");
if repo_address.strip_prefix(&prefix).is_none_or(str::is_empty) {
@@ -438,6 +465,9 @@ pub async fn sign_project_pull_request_status(
return Err("Invalid target repository owner.".to_string());
}
let identity = project_owner_identity(&app, &state, &target_owner)?;
// P29-C1 condition 1: admit BEFORE signing so the NIP-98 freshness window
// opens after the rate-limit wait, not between sign and submit.
let lease = admit_project_egress(&identity).await?;
let event = Event::from_json(build_pull_request_status_event(
&identity.keys,
&input.repo_address,
@@ -447,8 +477,14 @@ pub async fn sign_project_pull_request_status(
input.created_at,
)?)
.map_err(|error| format!("parse signed pull request status: {error}"))?;
submit_signed_event_with_keys(&event, &state, &identity.keys, identity.auth_tag.as_deref())
.await?;
submit_signed_event_with_keys(
&event,
&state,
&identity.keys,
identity.auth_tag.as_deref(),
&lease,
)
.await?;
Ok(())
}
@@ -463,6 +499,8 @@ pub async fn sign_project_pull_request_review_request(
return Err("Invalid target repository owner.".to_string());
}
let identity = project_owner_identity(&app, &state, &target_owner)?;
// P29-C1 condition 1: admit before signing.
let lease = admit_project_egress(&identity).await?;
let event = Event::from_json(build_review_request_event(
&identity.keys,
&input.repo_address,
@@ -471,8 +509,14 @@ pub async fn sign_project_pull_request_review_request(
&input.reviewer_label,
)?)
.map_err(|error| format!("parse signed review request: {error}"))?;
submit_signed_event_with_keys(&event, &state, &identity.keys, identity.auth_tag.as_deref())
.await?;
submit_signed_event_with_keys(
&event,
&state,
&identity.keys,
identity.auth_tag.as_deref(),
&lease,
)
.await?;
Ok(())
}
@@ -495,8 +539,22 @@ pub async fn publish_project_pull_request_merged_status(
return Err("Invalid merged pull request status event.".to_string());
}
let identity = project_owner_identity(&app, &state, &target_owner)?;
submit_signed_event_with_keys(&event, &state, &identity.keys, identity.auth_tag.as_deref())
.await?;
// P29-C1 condition 2 (documented pre-signed exception): this command
// receives an ALREADY-SIGNED kind-1631 event from input (verified above:
// kind, pubkey == target_owner, signature). Admission-before-signing is
// impossible here — the signature is fixed and its freshness is
// input-determined by `status_created_at` — so admit-before-submit is all
// this site can do. Chosen, not missed; noted in the C8 closed-world
// evidence so the exception is auditable.
let lease = admit_project_egress(&identity).await?;
submit_signed_event_with_keys(
&event,
&state,
&identity.keys,
identity.auth_tag.as_deref(),
&lease,
)
.await?;
Ok(())
}
@@ -663,11 +721,17 @@ pub async fn merge_project_pull_request(
)?;
let signed_status = Event::from_json(&status_event)
.map_err(|error| format!("parse signed merged status: {error}"))?;
// P29-C1 condition 1: the git clone→merge→push above ran in spawn_blocking
// for seconds to minutes; admit only now, after it completes and before the
// status event is built/signed, so the lease spans sign→auth→transmit and
// never the blocking git op.
let lease = admit_project_egress(&owner_identity).await?;
let status_publication_error = submit_signed_event_with_keys(
&signed_status,
&state,
&owner_identity.keys,
owner_identity.auth_tag.as_deref(),
&lease,
)
.await
.err();
@@ -64,10 +64,7 @@ pub async fn list_relay_members(state: State<'_, AppState>) -> Result<serde_json
pub async fn get_my_relay_membership(
state: State<'_, AppState>,
) -> Result<serde_json::Value, String> {
let my_pubkey = {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
keys.public_key().to_hex()
};
let my_pubkey = state.current_pubkey()?.to_hex();
let events = query_relay(
&state,
@@ -55,7 +55,7 @@ fn setup_team_import_app_with_scope(
let owner_keys = nostr::Keys::generate();
let state = crate::app_state::build_app_state();
{
let mut locked = state.keys.lock().unwrap();
let mut locked = state.identity_lifecycle_keys_guard().unwrap();
*locked = owner_keys.clone();
}
@@ -778,7 +778,7 @@ mod egress_guard_boundary {
) -> crate::app_state::AppState {
let state = crate::app_state::build_app_state();
{
let mut locked = state.keys.lock().unwrap();
let mut locked = state.identity_lifecycle_keys_guard().unwrap();
*locked = owner_keys;
}
if let Some(s) = scope {
+1 -2
View File
@@ -290,8 +290,7 @@ pub async fn deny_approval(
// ── Helpers (pure, unit-tested in workflows_tests.rs) ─────────────────────────
fn current_pubkey_hex(state: &AppState) -> Result<String, String> {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
Ok(keys.public_key().to_hex())
Ok(state.current_pubkey()?.to_hex())
}
fn now_secs() -> i64 {
+4 -9
View File
@@ -56,11 +56,11 @@ pub struct ActiveWorkspaceInfo {
/// Returns the current active workspace info (relay URL + pubkey).
#[tauri::command]
pub fn get_active_workspace(state: State<'_, AppState>) -> Result<ActiveWorkspaceInfo, String> {
let keys = state.keys.lock().map_err(|e| e.to_string())?;
let pubkey = state.current_pubkey()?;
let relay_url = relay::relay_ws_url_with_override(&state);
Ok(ActiveWorkspaceInfo {
relay_url,
pubkey: keys.public_key().to_hex(),
pubkey: pubkey.to_hex(),
})
}
@@ -193,12 +193,7 @@ async fn apply_workspace_body(
let base_dir = crate::managed_agents::managed_agents_base_dir(&app).unwrap_or_default();
let effective_owner_pubkey = match &parsed_keys {
Some(keys) => keys.public_key().to_hex(),
None => state
.keys
.lock()
.map_err(|e| e.to_string())?
.public_key()
.to_hex(),
None => state.current_pubkey()?.to_hex(),
};
let target_scope_id =
crate::managed_agents::scope::derive_scope_id(&relay_url, &effective_owner_pubkey);
@@ -297,7 +292,7 @@ async fn apply_workspace_body(
return Ok(WorkspaceApplyResult::drain_failed(msg));
}
};
let mut keys_guard = match state.keys.lock() {
let mut keys_guard = match state.identity_lifecycle_keys_guard() {
Ok(g) => g,
Err(e) => {
drop(override_guard);
+11 -3
View File
@@ -91,6 +91,7 @@ async fn boundary_submit_event_at_with_keys_blocks_ncryptsec() {
&state,
"http://127.0.0.1:9", // discard port — must never be reached
&keys,
&crate::owner_identity_egress::test_owner_egress_lease(),
)
.await
.unwrap_err();
@@ -130,6 +131,7 @@ async fn boundary_submit_signed_event_at_with_keys_blocks_ncryptsec() {
&state,
"http://127.0.0.1:9", // discard port — must never be reached
&keys,
&crate::owner_identity_egress::test_owner_egress_lease(),
)
.await
.unwrap_err();
@@ -145,9 +147,15 @@ async fn boundary_submit_signed_event_with_keys_blocks_ncryptsec() {
let event = nostr::EventBuilder::new(nostr::Kind::Custom(9), NCRYPTSEC)
.sign_with_keys(&keys)
.unwrap();
let err = crate::relay::submit_signed_event_with_keys(&event, &state, &keys, None)
.await
.unwrap_err();
let err = crate::relay::submit_signed_event_with_keys(
&event,
&state,
&keys,
None,
&crate::owner_identity_egress::test_owner_egress_lease(),
)
.await
.unwrap_err();
assert_guard_error(&err);
}
+4 -6
View File
@@ -301,9 +301,8 @@ pub async fn start_huddle(
successful_agents.clone();
hs.maybe_auto_enable_transcription_for_agents();
let own_pubkey = state
.keys
.lock()
.map(|k| k.public_key().to_hex())
.current_pubkey()
.map(|pk| pk.to_hex())
.unwrap_or_default();
let mut participants = successful_agents.clone();
if !own_pubkey.is_empty() && !participants.contains(&own_pubkey) {
@@ -421,9 +420,8 @@ pub async fn join_huddle(
// Seed participant list with own pubkey as a fallback until relay responds.
let own_pubkey = state
.keys
.lock()
.map(|k| k.public_key().to_hex())
.current_pubkey()
.map(|pk| pk.to_hex())
.unwrap_or_default();
let committed = {
+35 -8
View File
@@ -12,7 +12,7 @@ use std::{
};
use nostr::JsonUtil;
use tauri::State;
use tauri::{Manager, State};
use uuid::Uuid;
use crate::app_state::AppState;
@@ -632,8 +632,16 @@ pub(crate) fn spawn_transcription_task(
let spawned_gen = session_generation.load(Ordering::Acquire);
let http_client = state.http_client.clone();
let keys = match state.keys.lock() {
Ok(k) => k.clone(),
// Capture the AppHandle (stable for the process lifetime) rather than a
// long-lived owner Keys clone: P29-C1 forbids the STT task from pinning
// owner identity across the huddle. Keys are re-resolved per send through
// the recovery-gated accessor below, so an identity that enters recovery
// mid-huddle stops publishing at once.
let app = match state.app_handle.lock() {
Ok(guard) => match guard.clone() {
Some(app) => app,
None => return,
},
Err(_) => return,
};
let relay_base_url = crate::relay::relay_api_base_url_with_override(state);
@@ -667,11 +675,29 @@ pub(crate) fn spawn_transcription_task(
continue;
}
};
// Wait before signing: the relay enforces NIP-98 freshness (±60s)
// and the gate may hold for up to MAX_HINT_SECONDS (300s). Sign
// the kind event and build NIP-98 auth after the wait so both
// timestamps are fresh — single clean order: wait → sign → auth → send.
crate::relay_admission::wait_for_rate_limit().await;
// Re-resolve keys per send through the recovery-gated accessor and
// admit an owner-identity egress lease. try_admit_owner_identity_egress
// waits out the rate-limit gate internally (NIP-98 freshness ±60s
// under a ≤300s hold), so the lease is born after the wait and spans
// only sign → auth → transmit. Acquiring per send means no owner Keys
// outlive a single publish, and a mid-huddle recovery latch refuses
// the lease so this task stops publishing.
let app_state = app.state::<AppState>();
let keys = match app_state.signing_keys() {
Ok(k) => k,
Err(e) => {
eprintln!("buzz-desktop: STT signing key unavailable: {e}");
break;
}
};
let lease = match crate::owner_identity_egress::try_admit_owner_identity_egress().await
{
Ok(lease) => crate::owner_identity_egress::EgressLease::OwnerIdentity(lease),
Err(e) => {
eprintln!("buzz-desktop: STT egress refused: {e}");
break;
}
};
let body_bytes = match sign_and_guard_stt_body(builder, &keys) {
Ok(b) => b,
Err(e) => {
@@ -685,6 +711,7 @@ pub(crate) fn spawn_transcription_task(
&reqwest::Method::POST,
&url,
&body_bytes,
&lease,
) {
Ok(h) => h,
Err(e) => {
+1 -1
View File
@@ -55,7 +55,7 @@ pub(crate) async fn connect_audio_relay(
let relay_url = crate::relay::relay_ws_url_with_override(state);
let ws_url = format!("{relay_url}/huddle/{channel_id}/audio");
let keys = state.keys.lock().map_err(|e| e.to_string())?.clone();
let keys = state.signing_keys()?;
// TTS interrupt flags — recv task cancels TTS when remote humans speak.
let (tts_cancel, tts_active) = {
+1
View File
@@ -28,6 +28,7 @@ mod models;
mod native_websocket;
mod nostr_bind;
pub mod nostr_convert;
mod owner_identity_egress;
mod prevent_sleep;
mod ptt_shortcut;
mod relay;
@@ -298,6 +298,13 @@ async fn flush_pending_events_at(
// timestamp at publish time; kind, tags, and content are preserved,
// and `mark_synced` below still compares against the retained row's
// original `created_at`/`content`, which are untouched.
// Admit BEFORE re-signing: archive requests get a fresh `created_at`
// at publish time (relay ±120s freshness), so admission — which waits
// out the rate-limit gate — must precede the re-sign to keep the
// request fresh under a gate hold.
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let is_archive_request =
buzz_core_pkg::kind::is_identity_archive_request_kind(current.kind);
let event = if is_archive_request {
@@ -311,6 +318,7 @@ async fn flush_pending_events_at(
state,
&relay_api_base,
owner_keys,
&lease,
)
.await
.is_err()
@@ -855,7 +855,7 @@ mod flush_barrier {
.sign_with_keys(&keys)
.unwrap();
let state = build_app_state();
*state.keys.lock().unwrap() = keys;
*state.identity_lifecycle_keys_guard().unwrap() = keys;
let fresh = resign_with_fresh_timestamp(&stale, &state).unwrap();
@@ -914,7 +914,7 @@ mod flush_barrier {
}
let state = build_app_state();
*state.keys.lock().unwrap() = keys;
*state.identity_lifecycle_keys_guard().unwrap() = keys;
*state.relay_url_override.lock().unwrap() = Some(spawn_stub_relay().await);
let flushed = flush_pending_events(&db_path, &state).await.expect("flush");
@@ -236,13 +236,8 @@ pub async fn restore_managed_agents_on_launch(
// Snapshot the workspace owner pubkey once for the legacy auth_tag fallback.
// Read outside the per-agent spawn loop so all parallel spawns see the same
// value and we don't lock `state.keys` repeatedly.
let owner_hex: Option<String> = state
.keys
.lock()
.map_err(|e| e.to_string())
.ok()
.map(|k| k.public_key().to_hex());
// value and we don't re-read the identity repeatedly.
let owner_hex: Option<String> = state.current_pubkey().ok().map(|pk| pk.to_hex());
#[cfg(feature = "mesh-llm")]
let agents_to_start = {
@@ -290,11 +290,7 @@ fn start_pair_under_held_locks<R: tauri::Runtime>(
runtimes.remove(&key);
terminate_untracked_pair_runtime(app, &key)?;
let owner = state
.keys
.lock()
.ok()
.map(|keys| keys.public_key().to_hex());
let owner = state.current_pubkey().ok().map(|pk| pk.to_hex());
let scope_id = state
.capture_active_scope()
.map(|scope| scope.scope_id.clone());
@@ -421,6 +417,11 @@ async fn probe_agent_relay_access(
let keys = nostr::Keys::parse(record.private_key_nsec.trim())
.map_err(|error| format!("invalid managed-agent key: {error}"))?;
let api_base = crate::relay::relay_http_base_url(&key.relay_url);
// Managed-agent egress construction site (P29-C1 closed-world sink). Admit
// the interim keyed-egress lease before the probe query.
let lease = crate::owner_identity_egress::EgressLease::ManagedAgentKeyed(
crate::owner_identity_egress::admit_managed_agent_egress().await?,
);
tokio::time::timeout(
std::time::Duration::from_secs(10),
crate::relay::query_relay_at_with_keys(
@@ -429,6 +430,7 @@ async fn probe_agent_relay_access(
&[serde_json::json!({"kinds": [39002], "#p": [record.pubkey]})],
&keys,
record.auth_tag.as_deref(),
&lease,
),
)
.await
@@ -480,12 +480,16 @@ async fn publish_status_report_at(
payload: serde_json::Value,
) -> Result<(), String> {
let api_base_url = crate::relay::relay_http_base_url(relay_url);
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let keys = state.signing_keys()?;
crate::relay::submit_event_at_with_keys(
build_status_report_event(payload)?,
state,
&api_base_url,
&keys,
&lease,
)
.await
.map(|_| ())
@@ -0,0 +1,738 @@
//! Owner-identity egress registry — revocation of pre-latch capability (P29-C1).
//!
//! # Why this exists
//!
//! Making `AppState.keys` private (P28-C1) closes capability ACQUISITION but
//! cannot revoke owner `Keys` already CLONED before an identity transition. At
//! base, long-lived tasks (huddle STT) capture the owner keys once and then
//! sign, build NIP-98 authorization, and transmit indefinitely without
//! consulting `AppState` again, and the explicit-key egress funnels
//! (`submit_signed_event_at_with_keys`, `submit_event_at_with_keys`,
//! `submit_signed_event_with_keys`, `submit_event_with_keys`,
//! `query_relay_at_with_keys`, `build_nip98_auth_header{,_for_keys}`) accept an
//! already-issued `&Keys` with no latch or generation check. A task that
//! captured keys under identity A and resumes after the coordinator has latched
//! an ambiguous transition would sign and transmit as the unresolved identity.
//!
//! This module makes owner-identity egress itself latch-aware, reusing the
//! admission + RAII lease + drain-barrier shape rather than inventing a
//! parallel mechanism:
//!
//! - A **process-global registry** carries an *identity-persistence
//! generation* and an admission [`IdentityPersistenceState`].
//! - Every owner-identity sign / NIP-98-build / network-submit funnel acquires
//! an [`OwnerIdentityEgressLease`] at the last irreversible boundary via
//! [`try_admit_owner_identity_egress`]; admission validates — immediately
//! before the operation — that the state is [`Live`](IdentityPersistenceState::Live).
//! The lease is held only across sign → auth → transmit, never across
//! rate-limit waits, which bounds the drain below.
//! - The identity-transition coordinator, after the runtime drain and BEFORE
//! the journal write + durable B dispatch, calls [`begin_egress_drain`] to
//! bump the generation and refuse new admission, then [`await_egress_drain`]
//! to await every in-flight lease — so no send begun under identity A can
//! cross the durable-dispatch boundary. The TOCTOU window is closed by
//! linearization (admission and state transitions share one lock), not by
//! prose. New admission resumes only at a proven exit via
//! [`resume_egress_live`] (`Committed` / `DefinitelyUnchanged` /
//! verified-rollback) or is latched fail-closed via
//! [`latch_identity_indeterminate`] on the ambiguous branch (P28-C1).
//!
//! # Scope of this landing (C1a)
//!
//! This is the P29 substrate only: the registry, the per-send RAII lease, and
//! the coordinator drain/latch API. No production path drives a transition yet
//! — the generation never bumps and the state stays
//! [`Live`](IdentityPersistenceState::Live) until the identity-transition
//! coordinator wires the drain — so this landing is behavior-preserving. The
//! generic `OwnerIdentityCapability<T>` with its `Session`/`Bearer`/`Artifact`
//! policies is a later extension of this same registry; the per-send lease is
//! its `BoundedLease` instantiation, and its public name is fixed here so the
//! funnel witness signatures stay stable across that extension.
//!
//! # Closed-world egress construction sites (C1c condition 3)
//!
//! The `EgressLease` witness is only ever minted at these eight sites — the
//! complete owner/managed-agent egress sink set. The `relay_admission`
//! module doc and the C8 closed-world sink test enumerate exactly this set; a
//! new send that skips admission fails to compile (no witness) and fails the
//! sink-test enumeration.
//!
//! Owner-identity (`try_admit_owner_identity_egress`):
//! - `submit_event` / `query_relay_at` / `build_nip98_auth_header` — the
//! owner-default wrappers self-admit, so their transitive callers need no
//! witness.
//! - `mesh_llm::coordinator` status-report publication.
//! - `commands::profile::update_profile_at_relay` — query → submit → re-query,
//! three admissions (the exactly-once-per-op live case).
//! - `managed_agents::persona_events`, `commands::personas::sharing`,
//! `commands::channels` (×3), and the huddle STT task (per-send inline
//! owner lease).
//!
//! Managed-agent keyed (`admit_managed_agent_egress`) — the four sink sites
//! plus the git-workflow branch, variant selected at runtime by
//! `ProjectOwnerIdentity::is_managed_agent` (never derived from `auth_tag`):
//! 1. `submit_engram_event` (snapshot import).
//! 2. `sync_managed_agent_profile`.
//! 3. the agent query probe (`runtime_commands`).
//! 4. `commands::messages` reaction publish.
//! 5. the four `project_git_workflow` PR-status/merge sends (items 5–8).
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{LazyLock, Mutex};
use tokio::sync::Notify;
/// Admission state of owner-identity egress — the third recovery state
/// (`Indeterminate`) extends the existing `AppState::identity_lost` /
/// `keyring_locked` model (P28-C1).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum IdentityPersistenceState {
/// Persistence is settled: egress admission succeeds and the checked key
/// accessors serve keys normally.
Live,
/// A transition is draining in-flight egress before its durable dispatch.
/// New lease admission is refused; the checked key accessors still serve
/// keys, because already-admitted in-flight leases hold their own captures
/// and complete under the outgoing identity before the durable dispatch.
/// Transient and in-memory only — cleared at a proven exit.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
Draining,
/// A completed transition could not prove EITHER durable identity
/// canonical. Latched durable fail-closed: owner-identity egress admission
/// AND the checked key accessors refuse until reconciliation proves one
/// durable identity canonical. This is the P28-C1 third recovery state,
/// gating all checked identity access.
Indeterminate,
}
/// Mutable registry state, guarded by a single mutex so admission and every
/// state transition linearize (closing the check-then-send TOCTOU window).
struct RegistryInner {
state: IdentityPersistenceState,
/// Number of admitted leases not yet dropped. The drain awaits this
/// reaching zero.
in_flight: u64,
}
/// Process-global owner-identity egress registry. One per process, mirroring
/// the scope-generation authority in [`crate::managed_agents::scope`].
struct Registry {
inner: Mutex<RegistryInner>,
/// Monotonic identity-persistence generation. Read lock-free on the hot
/// admission path; only ever advanced by [`begin_egress_drain`] under
/// `inner`.
generation: AtomicU64,
/// Notified when `in_flight` reaches zero so the drain can wake.
drained: Notify,
}
static REGISTRY: LazyLock<Registry> = LazyLock::new(|| Registry {
inner: Mutex::new(RegistryInner {
state: IdentityPersistenceState::Live,
in_flight: 0,
}),
generation: AtomicU64::new(0),
drained: Notify::new(),
});
fn lock_inner() -> std::sync::MutexGuard<'static, RegistryInner> {
// A poisoned egress lock means a lease-holding thread panicked mid-send.
// The counters remain consistent (the panicking thread's lease Drop still
// runs), so recover the guard rather than propagating the poison and
// wedging every subsequent admission and drain.
REGISTRY
.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
/// The current identity-persistence generation. A lease is valid only while
/// this equals the generation captured at its admission.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub fn current_identity_persistence_generation() -> u64 {
REGISTRY.generation.load(Ordering::Acquire)
}
/// The current admission state.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub fn identity_persistence_state() -> IdentityPersistenceState {
lock_inner().state
}
/// Whether owner-identity persistence is latched `Indeterminate`. The checked
/// key accessors (P28-C1) refuse when this is `true`.
pub fn is_identity_indeterminate() -> bool {
lock_inner().state == IdentityPersistenceState::Indeterminate
}
/// A per-send owner-identity egress lease — the `BoundedLease` capability.
///
/// Acquired at the last irreversible egress boundary and held only across
/// sign → auth → transmit. Dropping it decrements the in-flight count so the
/// coordinator's drain can complete. Not `Clone` and carries no keys: the
/// caller pairs it with a per-send key acquisition through the checked
/// accessor, so a lease can never outlive the send it authorizes.
#[derive(Debug)]
#[must_use = "an egress lease authorizes exactly one sign/auth/transmit window; \
hold it across that window and drop it immediately after"]
pub struct OwnerIdentityEgressLease {
/// The identity-persistence generation current at admission.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
generation: u64,
}
impl OwnerIdentityEgressLease {
/// The identity-persistence generation this lease was admitted under.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub fn generation(&self) -> u64 {
self.generation
}
}
impl Drop for OwnerIdentityEgressLease {
fn drop(&mut self) {
let mut inner = lock_inner();
inner.in_flight = inner.in_flight.saturating_sub(1);
let drained = inner.in_flight == 0;
drop(inner);
if drained {
// Wake a drain awaiting the last in-flight lease. Harmless when no
// drain is in progress (no waiter registered).
REGISTRY.drained.notify_waiters();
}
}
}
/// Admit an owner-identity egress send, returning a lease bound to the current
/// identity-persistence generation.
///
/// Admission performs the relay rate-limit wait FIRST, then admits: the lease
/// is born *after* the wait completes, so it never spans the wait by
/// construction (spec L4505–4508) and the coordinator drain it participates in
/// stays bounded. Because the explicit-key funnels require an
/// [`EgressLease`] witness and the only constructors of one are these `admit_*`
/// entry points, a send can never reach the funnel without having waited — no
/// call site can forget the rate-limit wait (the funnels no longer wait
/// themselves).
///
/// Succeeds only in [`IdentityPersistenceState::Live`]. Returns `Err` when a
/// transition is [`Draining`](IdentityPersistenceState::Draining) (new
/// admission refused while in-flight sends complete) or the state is latched
/// [`Indeterminate`](IdentityPersistenceState::Indeterminate) (fail-closed
/// after an unprovable transition). The admission check and the generation read
/// happen under one lock AFTER the wait, so the lease reflects the persistence
/// state immediately before the operation — a wait that overlaps a drain
/// refuses on wake rather than admitting a stale generation.
pub async fn try_admit_owner_identity_egress() -> Result<OwnerIdentityEgressLease, String> {
crate::relay_admission::wait_for_rate_limit().await;
admit_owner_identity_after_wait()
}
/// The post-wait owner-identity admission check + in-flight bump, under one
/// lock. Split from the wait so the state-machine invariants are unit-testable
/// without driving the async rate-limit gate.
fn admit_owner_identity_after_wait() -> Result<OwnerIdentityEgressLease, String> {
let mut inner = lock_inner();
match inner.state {
IdentityPersistenceState::Live => {
inner.in_flight += 1;
Ok(OwnerIdentityEgressLease {
generation: REGISTRY.generation.load(Ordering::Acquire),
})
}
IdentityPersistenceState::Draining => Err(
"owner-identity egress is draining for an identity transition; \
signing and publishing are paused until it resolves"
.to_string(),
),
IdentityPersistenceState::Indeterminate => Err(
"owner identity is in an indeterminate recovery state; event \
signing is disabled until the identity is reconciled and Buzz is \
relaunched"
.to_string(),
),
}
}
/// Begin the coordinator's egress drain: bump the identity-persistence
/// generation and close new admission (state → `Draining`). Returns the new
/// generation. Call this after the runtime drain and BEFORE the journal write
/// + durable B dispatch, then [`await_egress_drain`] to await in-flight leases.
///
/// Idempotent-guard: only transitions from `Live`. Returns `Err` if a
/// transition is already in flight or latched — the coordinator serializes
/// transitions under its own guards, so this is a defensive invariant check.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub fn begin_egress_drain() -> Result<u64, String> {
let mut inner = lock_inner();
if inner.state != IdentityPersistenceState::Live {
return Err(format!(
"cannot begin egress drain from state {:?}; a transition is already \
in flight or latched",
inner.state
));
}
inner.state = IdentityPersistenceState::Draining;
// Bump under the lock so no `Live` admission can capture the pre-bump
// generation after the state has flipped to `Draining`.
Ok(REGISTRY.generation.fetch_add(1, Ordering::AcqRel) + 1)
}
/// Await every in-flight lease admitted before the drain began. Returns once
/// no lease is outstanding. Holds no lock across its await, and an in-flight
/// lease never re-acquires a state guard mid-lease, so the drain cannot
/// deadlock against a send it is waiting on.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub async fn await_egress_drain() {
loop {
let notified = REGISTRY.drained.notified();
tokio::pin!(notified);
// Register the waiter BEFORE reading the count so a lease Drop that
// reaches zero between the read and the await cannot be lost.
notified.as_mut().enable();
if lock_inner().in_flight == 0 {
return;
}
notified.await;
}
}
/// Resume admission at a proven transition exit (`Committed` finished, or
/// `DefinitelyUnchanged` / verified-rollback): state → `Live`. The generation
/// is already current for the winning identity, so leases admitted from here
/// carry it.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub fn resume_egress_live() {
lock_inner().state = IdentityPersistenceState::Live;
}
/// Latch the durable fail-closed state on the ambiguous branch (P28-C1): state
/// → `Indeterminate`. Owner-identity egress admission and the checked key
/// accessors refuse until [`resume_egress_live`] is called after reconciliation
/// proves one durable identity canonical.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub fn latch_identity_indeterminate() {
lock_inner().state = IdentityPersistenceState::Indeterminate;
}
/// A witness that an owner-identity egress funnel caller holds a valid egress
/// capability for the send it is about to perform.
///
/// The explicit-key funnels (`submit_signed_event_at_with_keys`,
/// `submit_event_at_with_keys`, `submit_signed_event_with_keys`,
/// `submit_event_with_keys`, `query_relay_at_with_keys`,
/// `build_nip98_auth_header{,_for_keys}`) take an
/// `&EgressLease` so NO unleased caller shape exists — a bypass is a compile
/// error, not a convention (P29-C1). Two caller classes share those funnels,
/// so the witness accommodates both:
///
/// - [`OwnerIdentity`](EgressLease::OwnerIdentity) — the fully-real owner
/// capability: a generation-qualified [`OwnerIdentityEgressLease`] admitted
/// by [`try_admit_owner_identity_egress`], which refuses under the P28-C1
/// latch and while a transition drains.
/// - [`ManagedAgentKeyed`](EgressLease::ManagedAgentKeyed) — an interim token
/// for managed-agent-key callers (the genuine agent callers of
/// `build_nip98_auth_header_for_keys`: `sync_managed_agent_profile`,
/// `submit_engram_event`, and the agent query-probe). See
/// [`ManagedAgentEgressLease`] for the Phase-4 contract this interim shape
/// defers.
#[derive(Debug)]
pub enum EgressLease {
/// The fully-real owner-identity capability (P29-C1).
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
OwnerIdentity(OwnerIdentityEgressLease),
/// The interim managed-agent-key capability (Phase 4 completes it).
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
ManagedAgentKeyed(ManagedAgentEgressLease),
}
/// Interim managed-agent keyed-egress capability — the seam Phase 4 completes.
///
/// # Phase-4 contract (the deferred half)
///
/// The §3.3a P29 contract requires the shared explicit-key funnels to admit a
/// "§3.3 keyed-egress lease for managed-agent callers … which already refuse
/// under the latch." The full §3.3 keyed-egress substrate — split
/// admission/drain keys, per-scope generation gating, and the
/// `admit_keyed_egress()` constructor — is Phase 4 work (§8 item 4), sequenced
/// AFTER Phase 3b. This interim newtype exists so the funnels take a uniform
/// [`EgressLease`] witness from Phase 3b onward and their signatures never
/// churn across the phase boundary.
///
/// **`admit_managed_agent_egress` is the ONLY constructor and will be replaced
/// by Phase 4's `admit_keyed_egress`, which adds the deferred scope-generation
/// gate.** This interim form implements only the latch/drain half: it refuses
/// admission whenever owner-identity persistence is not [`Live`] (the P28-C1
/// latch, or a transition drain). Phase 4 ADDS scope-generation gating to an
/// already-fail-closed token — it never loosens this behavior. Managed-agent
/// scope is keyed by `(relay, owner)`; when owner identity persistence is
/// `Indeterminate`, which scope is legitimately active is itself indeterminate,
/// so agent egress must fail closed there — the exact cross-scope leak this arc
/// exists to prevent.
#[derive(Debug)]
#[must_use = "a managed-agent egress capability authorizes exactly one \
agent-key sign/auth/transmit window"]
pub struct ManagedAgentEgressLease {
/// The owner identity-persistence generation current at admission. Held so
/// the coordinator's drain awaits agent sends in flight at the barrier,
/// exactly as it awaits owner leases.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
generation: u64,
}
impl ManagedAgentEgressLease {
/// The owner identity-persistence generation this capability was admitted
/// under.
// Consumed by C5 (P25/P28 coordinator); remove allow when C5 lands.
#[allow(dead_code)]
pub fn generation(&self) -> u64 {
self.generation
}
}
impl Drop for ManagedAgentEgressLease {
fn drop(&mut self) {
let mut inner = lock_inner();
inner.in_flight = inner.in_flight.saturating_sub(1);
let drained = inner.in_flight == 0;
drop(inner);
if drained {
REGISTRY.drained.notify_waiters();
}
}
}
/// Admit a managed-agent keyed-egress send — the interim P29 seam (see
/// [`ManagedAgentEgressLease`] for the Phase-4 contract).
///
/// Waits out the relay rate-limit gate FIRST, then admits, so the interim
/// capability shares the wait-in-admission shape of
/// [`try_admit_owner_identity_egress`]: the lease never spans the wait and no
/// agent call site can forget it.
///
/// Refuses whenever owner-identity persistence is not [`Live`] (latched
/// `Indeterminate`, or a transition draining), matching the "already refuse
/// under the latch" requirement. Phase 4's `admit_keyed_egress` replaces this
/// and adds the deferred scope-generation gate.
pub async fn admit_managed_agent_egress() -> Result<ManagedAgentEgressLease, String> {
crate::relay_admission::wait_for_rate_limit().await;
admit_managed_agent_after_wait()
}
/// The post-wait managed-agent admission check + in-flight bump, under one
/// lock. Split from the wait so the state-machine invariants are unit-testable
/// without driving the async rate-limit gate.
fn admit_managed_agent_after_wait() -> Result<ManagedAgentEgressLease, String> {
let mut inner = lock_inner();
match inner.state {
IdentityPersistenceState::Live => {
inner.in_flight += 1;
Ok(ManagedAgentEgressLease {
generation: REGISTRY.generation.load(Ordering::Acquire),
})
}
IdentityPersistenceState::Draining => Err(
"owner-identity egress is draining for an identity transition; \
managed-agent publishing is paused until it resolves"
.to_string(),
),
IdentityPersistenceState::Indeterminate => Err(
"owner identity is in an indeterminate recovery state; the active \
agent scope is unresolved, so managed-agent egress is disabled \
until the identity is reconciled and Buzz is relaunched"
.to_string(),
),
}
}
/// Process-global mutex serializing tests that mutate the registry's
/// process-global state. Any test that admits a lease, drives a drain, or
/// latches must hold this guard for its whole duration and reset via
/// [`reset_registry_for_test`] on entry, because all tests share one registry.
#[cfg(test)]
pub(crate) static EGRESS_REGISTRY_TEST_LOCK: Mutex<()> = Mutex::new(());
/// Reset the registry to its baseline (`Live`, zero in-flight) for a test.
/// The generation is monotonic and never resets, so tests assert relative
/// generation movement, never absolute values.
#[cfg(test)]
pub(crate) fn reset_registry_for_test() {
let mut inner = lock_inner();
inner.state = IdentityPersistenceState::Live;
inner.in_flight = 0;
}
/// Mint an owner-identity [`EgressLease`] for tests that exercise a funnel's
/// non-admission behavior (e.g. the NIP-49 egress-guard boundary or NIP-98
/// freshness) and only need a witness value. Admits directly through the
/// post-wait core so no rate-limit gate is driven.
#[cfg(test)]
pub(crate) fn test_owner_egress_lease() -> EgressLease {
EgressLease::OwnerIdentity(
admit_owner_identity_after_wait().expect("test lease admits when live"),
)
}
#[cfg(test)]
mod tests {
use super::*;
fn guard() -> std::sync::MutexGuard<'static, ()> {
let g = EGRESS_REGISTRY_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
reset_registry_for_test();
g
}
#[test]
fn admits_when_live() {
let _g = guard();
let lease = admit_owner_identity_after_wait().expect("live admits");
assert_eq!(
lease.generation(),
current_identity_persistence_generation()
);
assert_eq!(identity_persistence_state(), IdentityPersistenceState::Live);
}
#[test]
fn lease_drop_decrements_in_flight() {
let _g = guard();
{
let _lease = admit_owner_identity_after_wait().unwrap();
assert_eq!(lock_inner().in_flight, 1);
}
assert_eq!(lock_inner().in_flight, 0, "drop must decrement");
}
#[test]
fn begin_drain_bumps_generation_and_refuses_new_admission() {
let _g = guard();
let before = current_identity_persistence_generation();
let new_gen = begin_egress_drain().expect("live drains");
assert_eq!(new_gen, before + 1, "drain bumps the generation by one");
assert_eq!(
identity_persistence_state(),
IdentityPersistenceState::Draining
);
assert!(
admit_owner_identity_after_wait().is_err(),
"no new admission while draining"
);
}
#[test]
fn begin_drain_refuses_when_not_live() {
let _g = guard();
begin_egress_drain().unwrap();
assert!(
begin_egress_drain().is_err(),
"cannot begin a second drain while one is in flight"
);
}
#[test]
fn resume_live_reopens_admission() {
let _g = guard();
begin_egress_drain().unwrap();
assert!(admit_owner_identity_after_wait().is_err());
resume_egress_live();
assert!(
admit_owner_identity_after_wait().is_ok(),
"a proven exit reopens admission"
);
}
#[test]
fn latch_indeterminate_refuses_egress_and_reports_state() {
let _g = guard();
begin_egress_drain().unwrap();
latch_identity_indeterminate();
assert!(is_identity_indeterminate());
assert_eq!(
identity_persistence_state(),
IdentityPersistenceState::Indeterminate
);
assert!(
admit_owner_identity_after_wait().is_err(),
"the indeterminate latch refuses all owner-identity egress"
);
}
#[test]
fn draining_does_not_report_indeterminate() {
// The checked key accessors refuse on `Indeterminate` only; a
// transient `Draining` must NOT trip the accessor refusal, so
// already-admitted in-flight leases can still obtain keys per send.
let _g = guard();
begin_egress_drain().unwrap();
assert!(
!is_identity_indeterminate(),
"draining is not the indeterminate latch"
);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] // EGRESS_REGISTRY_TEST_LOCK serialises parallel tests
async fn drain_returns_immediately_with_no_in_flight_leases() {
let _g = guard();
begin_egress_drain().unwrap();
// No leases outstanding — the drain completes without blocking.
await_egress_drain().await;
assert_eq!(lock_inner().in_flight, 0);
}
#[tokio::test]
#[allow(clippy::await_holding_lock)] // EGRESS_REGISTRY_TEST_LOCK serialises parallel tests
async fn drain_awaits_an_in_flight_lease_then_completes() {
let _g = guard();
let lease = admit_owner_identity_after_wait().unwrap();
begin_egress_drain().unwrap();
assert_eq!(lock_inner().in_flight, 1);
// The drain must not complete while the lease is held; drop it from a
// spawned task after the drain has begun awaiting.
let handle = tokio::spawn(async move {
tokio::task::yield_now().await;
drop(lease);
});
await_egress_drain().await;
assert_eq!(lock_inner().in_flight, 0, "drain awaited the lease drop");
handle.await.unwrap();
}
#[test]
fn lease_admitted_before_drain_carries_pre_bump_generation() {
let _g = guard();
let lease = admit_owner_identity_after_wait().unwrap();
let admitted_gen = lease.generation();
let new_gen = begin_egress_drain().unwrap();
assert!(
admitted_gen < new_gen,
"a lease admitted under A is stale once the drain bumps to B"
);
}
#[test]
fn managed_agent_admits_when_live() {
let _g = guard();
let lease = admit_managed_agent_after_wait().expect("live admits agent egress");
assert_eq!(
lease.generation(),
current_identity_persistence_generation()
);
}
#[test]
fn managed_agent_refuses_under_the_latch() {
// Condition (2): the interim managed-agent variant refuses under the
// P28-C1 latch, matching the spec's "already refuse under the latch".
let _g = guard();
begin_egress_drain().unwrap();
latch_identity_indeterminate();
assert!(
admit_managed_agent_after_wait().is_err(),
"the indeterminate latch refuses managed-agent egress"
);
}
#[test]
fn managed_agent_refuses_while_draining() {
let _g = guard();
begin_egress_drain().unwrap();
assert!(
admit_managed_agent_after_wait().is_err(),
"no new agent admission while a transition drains"
);
}
#[test]
fn managed_agent_lease_participates_in_the_drain() {
// The coordinator's drain must await agent sends in flight at the
// barrier exactly as it awaits owner leases.
let _g = guard();
{
let _lease = admit_managed_agent_after_wait().unwrap();
assert_eq!(lock_inner().in_flight, 1);
}
assert_eq!(lock_inner().in_flight, 0, "agent lease drop decrements");
}
#[test]
fn egress_lease_enum_wraps_both_variants() {
let _g = guard();
let owner = EgressLease::OwnerIdentity(admit_owner_identity_after_wait().unwrap());
let agent = EgressLease::ManagedAgentKeyed(admit_managed_agent_after_wait().unwrap());
assert!(matches!(owner, EgressLease::OwnerIdentity(_)));
assert!(matches!(agent, EgressLease::ManagedAgentKeyed(_)));
assert_eq!(lock_inner().in_flight, 2, "both variants hold a lease");
}
/// Wait-in-admission: an armed rate-limit gate holds admission until the
/// window clears, and the lease is only born afterward. This is the
/// structural guarantee that a leased send can never span the wait and no
/// call site can forget it — the wait lives inside the sole lease
/// constructor.
#[tokio::test(start_paused = true)]
#[allow(clippy::await_holding_lock)] // EGRESS_REGISTRY_TEST_LOCK serialises parallel tests
async fn admission_waits_out_the_rate_limit_gate_before_leasing() {
// Hold BOTH the registry guard and the gate serial: this test drives
// the shared process-wide rate-limit static.
let _serial = crate::relay_admission::TEST_SERIAL.lock().await;
let _g = guard();
crate::relay_admission::reset_rate_limit_gate();
crate::relay_admission::activate_rate_limit(Some(30));
let start = tokio::time::Instant::now();
let lease = try_admit_owner_identity_egress()
.await
.expect("admits once the window clears");
assert_eq!(
tokio::time::Instant::now() - start,
std::time::Duration::from_secs(30),
"admission must wait out the full armed window before leasing"
);
// The lease exists only post-wait, so it is in-flight now, not during.
assert_eq!(lock_inner().in_flight, 1);
drop(lease);
crate::relay_admission::reset_rate_limit_gate();
}
/// A gate armed while the state is latched still refuses AFTER the wait —
/// the wait does not admit a send the coordinator has fenced off. Proves
/// the check happens post-wait, closing the drain-during-wait window.
#[tokio::test(start_paused = true)]
#[allow(clippy::await_holding_lock)] // EGRESS_REGISTRY_TEST_LOCK serialises parallel tests
async fn admission_refuses_after_wait_when_draining() {
let _serial = crate::relay_admission::TEST_SERIAL.lock().await;
let _g = guard();
crate::relay_admission::reset_rate_limit_gate();
crate::relay_admission::activate_rate_limit(Some(5));
begin_egress_drain().unwrap();
assert!(
try_admit_owner_identity_egress().await.is_err(),
"a drain begun before/during the wait refuses admission on wake"
);
assert_eq!(
lock_inner().in_flight,
0,
"a refused admission leases nothing"
);
crate::relay_admission::reset_rate_limit_gate();
}
}
+36 -11
View File
@@ -103,16 +103,26 @@ pub fn build_nip98_auth_header(
url: &str,
body: &[u8],
state: &AppState,
lease: &crate::owner_identity_egress::EgressLease,
) -> Result<String, String> {
let keys = state.keys.lock().map_err(|error| error.to_string())?;
build_nip98_auth_header_for_keys(&keys, method, url, body)
let keys = state.signing_keys()?;
build_nip98_auth_header_for_keys(&keys, method, url, body, lease)
}
/// Build a NIP-98 HTTP-auth header signed with an explicit identity.
///
/// Requires an [`EgressLease`](crate::owner_identity_egress::EgressLease)
/// witness (P29-C1): this is one of the explicit-key egress funnels the
/// spec names, so no caller can sign NIP-98 auth with a raw `&Keys` without
/// first proving a lease. The lease is admitted post-rate-limit-wait, so
/// signing here always follows the wait — freshness (NIP-98 ±60s) holds even
/// under a ≤300s gate hold.
pub fn build_nip98_auth_header_for_keys(
keys: &Keys,
method: &Method,
url: &str,
body: &[u8],
_lease: &crate::owner_identity_egress::EgressLease,
) -> Result<String, String> {
let payload_hash = hex::encode(Sha256::digest(body));
@@ -317,11 +327,17 @@ pub async fn query_relay_at(
api_base_url: &str,
filters: &[serde_json::Value],
) -> Result<Vec<nostr::Event>, String> {
crate::relay_admission::wait_for_rate_limit().await;
// Owner-default query: admit the owner-identity egress lease (which waits
// out the rate-limit gate internally, then validates the latch) and hold
// it across sign → auth → transmit. The wrapper self-admits so its many
// transitive callers need no witness.
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let url = format!("{}/query", api_base_url);
let body_bytes =
serde_json::to_vec(filters).map_err(|e| format!("filter serialization failed: {e}"))?;
let auth = build_nip98_auth_header(&Method::POST, &url, &body_bytes, state)?;
let auth = build_nip98_auth_header(&Method::POST, &url, &body_bytes, state, &lease)?;
let response = state
.http_client
@@ -346,12 +362,12 @@ pub async fn query_relay_at_with_keys(
filters: &[serde_json::Value],
keys: &Keys,
auth_tag: Option<&str>,
lease: &crate::owner_identity_egress::EgressLease,
) -> Result<Vec<nostr::Event>, String> {
crate::relay_admission::wait_for_rate_limit().await;
let url = format!("{}/query", api_base_url);
let body_bytes =
serde_json::to_vec(filters).map_err(|e| format!("filter serialization failed: {e}"))?;
let auth = build_nip98_auth_header_for_keys(keys, &Method::POST, &url, &body_bytes)?;
let auth = build_nip98_auth_header_for_keys(keys, &Method::POST, &url, &body_bytes, lease)?;
let mut request = state
.http_client
.post(&url)
@@ -445,7 +461,13 @@ pub async fn sync_managed_agent_profile(
avatar_url: Option<&str>,
auth_tag: Option<&str>, // NIP-OA auth tag JSON
) -> Result<(), String> {
crate::relay_admission::wait_for_rate_limit().await;
// Managed-agent egress construction site (P29-C1 closed-world sink). Admit
// the interim keyed-egress lease, which waits out the rate-limit gate then
// refuses under the identity-persistence latch/drain, and hold it across
// sign → auth → transmit.
let lease = crate::owner_identity_egress::EgressLease::ManagedAgentKeyed(
crate::owner_identity_egress::admit_managed_agent_egress().await?,
);
// Build a signed kind:0 profile event (with optional NIP-OA auth tag).
let event = build_profile_event(agent_keys, display_name, avatar_url, auth_tag)?;
let event_json = event.as_json();
@@ -453,7 +475,8 @@ pub async fn sync_managed_agent_profile(
crate::egress_guard::assert_no_key_backup_bytes(&body_bytes, "agent profile sync")?;
let url = format!("{}/events", relay_http_base_url(relay_url));
let auth = build_nip98_auth_header_for_keys(agent_keys, &Method::POST, &url, &body_bytes)?;
let auth =
build_nip98_auth_header_for_keys(agent_keys, &Method::POST, &url, &body_bytes, &lease)?;
let mut request = state
.http_client
@@ -547,11 +570,12 @@ pub async fn submit_event_with_keys(
state: &AppState,
keys: &Keys,
auth_tag: Option<&str>,
lease: &crate::owner_identity_egress::EgressLease,
) -> Result<SubmitEventResponse, String> {
let event = builder
.sign_with_keys(keys)
.map_err(|e| format!("failed to sign event: {e}"))?;
submit_signed_event_with_keys(&event, state, keys, auth_tag).await
submit_signed_event_with_keys(&event, state, keys, auth_tag, lease).await
}
/// POST an already-signed event using the same explicit identity for NIP-98.
@@ -560,15 +584,16 @@ pub async fn submit_signed_event_with_keys(
state: &AppState,
keys: &Keys,
auth_tag: Option<&str>,
lease: &crate::owner_identity_egress::EgressLease,
) -> Result<SubmitEventResponse, String> {
if event.pubkey != keys.public_key() {
return Err("signed event does not match the publishing identity".to_string());
}
crate::relay_admission::wait_for_rate_limit().await;
let url = format!("{}/events", relay_api_base_url_with_override(state));
let body_bytes = event.as_json().into_bytes();
crate::egress_guard::assert_no_key_backup_bytes(&body_bytes, "signed event submit (keys)")?;
let auth_header = build_nip98_auth_header_for_keys(keys, &Method::POST, &url, &body_bytes)?;
let auth_header =
build_nip98_auth_header_for_keys(keys, &Method::POST, &url, &body_bytes, lease)?;
let mut request = state
.http_client
+13 -4
View File
@@ -18,15 +18,16 @@ pub async fn submit_signed_event_at_with_keys(
state: &AppState,
api_base_url: &str,
keys: &nostr::Keys,
lease: &crate::owner_identity_egress::EgressLease,
) -> Result<SubmitEventResponse, String> {
if event.pubkey != keys.public_key() {
return Err("signed event does not match the publishing identity".to_string());
}
crate::relay_admission::wait_for_rate_limit().await;
let url = format!("{}/events", api_base_url.trim_end_matches('/'));
let body_bytes = event.as_json().into_bytes();
crate::egress_guard::assert_no_key_backup_bytes(&body_bytes, "relay event submit")?;
let auth_header = build_nip98_auth_header_for_keys(keys, &Method::POST, &url, &body_bytes)?;
let auth_header =
build_nip98_auth_header_for_keys(keys, &Method::POST, &url, &body_bytes, lease)?;
let response = state
.http_client
@@ -60,11 +61,12 @@ pub async fn submit_event_at_with_keys(
state: &AppState,
api_base_url: &str,
keys: &nostr::Keys,
lease: &crate::owner_identity_egress::EgressLease,
) -> Result<SubmitEventResponse, String> {
let event = builder
.sign_with_keys(keys)
.map_err(|e| format!("failed to sign event: {e}"))?;
submit_signed_event_at_with_keys(&event, state, api_base_url, keys).await
submit_signed_event_at_with_keys(&event, state, api_base_url, keys, lease).await
}
/// Build and submit an event to the currently active workspace relay.
@@ -72,7 +74,14 @@ pub async fn submit_event(
builder: nostr::EventBuilder,
state: &AppState,
) -> Result<SubmitEventResponse, String> {
// Owner-default submit: admit the owner-identity egress lease (waits out the
// rate-limit gate internally, then validates the latch) and hold it across
// sign → auth → transmit. The wrapper self-admits so its transitive callers
// need no witness.
let lease = crate::owner_identity_egress::EgressLease::OwnerIdentity(
crate::owner_identity_egress::try_admit_owner_identity_egress().await?,
);
let api_base_url = relay_api_base_url_with_override(state);
let keys = state.signing_keys()?;
submit_event_at_with_keys(builder, state, &api_base_url, &keys).await
submit_event_at_with_keys(builder, state, &api_base_url, &keys, &lease).await
}
+25 -11
View File
@@ -4,18 +4,25 @@
//! sends until the quota window clears — matching the TS-side gate in
//! `relayRateLimitGate.ts` that already governs WebSocket operations.
//!
//! **Coverage:** all entry points in `relay.rs` (`query_relay_at`,
//! `submit_event`, `submit_signed_event`, `submit_signed_event_with_keys`,
//! `sync_managed_agent_profile`) and the three previously-direct senders
//! (`submit_engram_event` in snapshot import + team_snapshot, huddle STT)
//! all call `wait_for_rate_limit()` before `.send()`.
//! **Coverage:** the wait is centralized in the two owner-identity egress
//! admission constructors — `try_admit_owner_identity_egress` and
//! `admit_managed_agent_egress` (`owner_identity_egress`). Each waits out the
//! gate FIRST, then admits its lease. Because the explicit-key relay funnels
//! (`submit_signed_event_at_with_keys`, `submit_event_at_with_keys`,
//! `submit_signed_event_with_keys`, `submit_event_with_keys`,
//! `query_relay_at_with_keys`, and the NIP-98 builders) require an
//! `EgressLease` witness whose ONLY constructors are those admission entry
//! points, no send can reach `.send()` without having waited — the wait can no
//! longer be forgotten at a call site. The eight owner/managed-agent egress
//! construction sites this closes over are enumerated in the
//! `owner_identity_egress` module doc.
//!
//! **Media upload/download and `/info`** call `relay_error_message()` on
//! non-200 responses, so their 429s arm the shared gate as conservative
//! back-off (any relay overload signal is worth honouring across domains).
//! They do not call `wait_for_rate_limit()` themselves — their operations
//! are driven by user-initiated file transfers rather than bridge event flow,
//! and they have independent retry logic.
//! They do not wait on the gate themselves — their operations are driven by
//! user-initiated file transfers rather than bridge event flow, and they have
//! independent retry logic.
//!
//! **Community scope:** the gate is reset on every `apply_workspace` call,
//! mirroring the TS gate's `resetRateLimitGate()` on community switch in
@@ -100,14 +107,20 @@ pub fn reset_rate_limit_gate() {
*GATE_EXPIRY.lock().unwrap_or_else(|e| e.into_inner()) = None;
}
/// Serializes every test that arms the process-wide gate static — including
/// `owner_identity_egress`'s wait-in-admission test — so armed expiries never
/// bleed between parallel test threads.
#[cfg(test)]
pub(crate) static TEST_SERIAL: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
#[cfg(test)]
mod tests {
use super::*;
// The gate is a process-wide static shared by every test in this binary,
// so all gate tests serialize on one async lock to keep armed expiries
// from bleeding between parallel test threads.
pub(crate) static TEST_SERIAL: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
// so all gate tests serialize on the module-level `TEST_SERIAL` (shared
// with cross-module gate consumers) to keep armed expiries from bleeding
// between parallel test threads.
#[tokio::test(start_paused = true)]
async fn wait_returns_immediately_when_gate_is_inactive() {
@@ -398,6 +411,7 @@ mod tests {
&reqwest::Method::POST,
"https://relay.example.com/events",
b"{}",
&crate::owner_identity_egress::test_owner_egress_lease(),
)
.expect("header build must succeed");