mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): area-3+4 snapshot import seams, transition helpers
Area 3 — snapshot import phase seams: - Extract confirm_agent_snapshot_import_core and confirm_team_snapshot_import_core, generic over tauri::Runtime. - Add before_store (post-entry-capture, pre-Phase-3a-lock) and after_store (post-Phase-3a-lock, pre-Phase-3b) hooks; no-ops in production. - ProfilePublish<'a>/MemoryPublish<'a> borrowed arg structs — no secret-key cloning across closures. - Thin Tauri commands call cores with no-op hooks and real relay adapters; app.state::<AppState>() inside async move blocks avoids non-'static borrows. - Delete duplicate submit_engram_event from team_snapshot.rs; reuse import.rs version (pub(crate)) for both agent and team engram boundaries. - Add retain_team_pending_in_scope(scope, team) sibling in teams.rs; team snapshot Phase 3 calls it instead of live retain_team_pending to avoid re-resolving active scope after a possible switch. - Move egress-guard + avatar tests from import.rs into import_tests.rs (included via #[path]); update egress inventory and allowlists. - Add SCOPE_GENERATION_TEST_LOCK in scope.rs for cross-module serialization of generation-sensitive tests; agent + team seam tests share this lock. - 4 named seam tests all pass concurrently (verified with cargo test --lib). Area 4 — workspace transition helpers: - Extract with_workspace_transition_preflight: acquires workspace_transition, runs fail_if_client_mesh_active, invokes transition_body under guard. - Extract install_client_under_workspace_transition: acquires workspace_transition, validates full captured scope identity (scope_id, relay, owner_pubkey, generation) under guard, calls install. - Both helpers live in mesh_llm_scope.rs alongside fail_if_client_mesh_active. - Area 4 tests (3 named) implemented in mesh_llm_tests.rs (separate commit). Also: remove unused AppHandle import from nest.rs (zero warnings). Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
co-authored by
Will Pfleger
parent
1ec287f935
commit
1dc537dad5
@@ -4,6 +4,8 @@
|
||||
//! Most items here are `pub(super)` so they remain private to the module;
|
||||
//! `mesh_stop_client` is `pub(crate)` and re-exported as `pub` from `mesh_llm`.
|
||||
|
||||
use std::future::Future;
|
||||
|
||||
use tauri::{AppHandle, Manager, State};
|
||||
|
||||
use crate::app_state::AppState;
|
||||
@@ -95,6 +97,91 @@ pub(crate) async fn fail_if_client_mesh_active<R: tauri::Runtime>(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Acquire the `workspace_transition` lock, run the production
|
||||
/// `fail_if_client_mesh_active` preflight, then invoke `transition_body` while
|
||||
/// the guard remains held.
|
||||
///
|
||||
/// Used by `apply_workspace` and live identity import so neither duplicates the
|
||||
/// acquire + preflight sequence inline. Tests call this helper directly to prove
|
||||
/// the serialization contract against `install_client_under_workspace_transition`.
|
||||
///
|
||||
/// The transition body runs synchronously (or dispatches to `spawn_blocking`)
|
||||
/// while the async guard is alive on the current task. If the preflight fails,
|
||||
/// `transition_body` is never called.
|
||||
pub(crate) async fn with_workspace_transition_preflight<R, F, T>(
|
||||
app: &AppHandle<R>,
|
||||
transition_body: F,
|
||||
) -> Result<T, String>
|
||||
where
|
||||
R: tauri::Runtime,
|
||||
F: FnOnce() -> Result<T, String>,
|
||||
{
|
||||
let state = app.state::<AppState>();
|
||||
let _transition_guard = state.workspace_transition.lock().await;
|
||||
|
||||
// Fail closed if a client-mode Mesh runtime is active.
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
fail_if_client_mesh_active(app).await?;
|
||||
|
||||
transition_body()
|
||||
}
|
||||
|
||||
/// Acquire the `workspace_transition` lock, validate the full captured scope
|
||||
/// identity under the guard — `(scope_id, normalized relay, owner_pubkey,
|
||||
/// generation)` — then call the injected `install` closure.
|
||||
///
|
||||
/// Fails if:
|
||||
/// - no active scope exists at the time of validation;
|
||||
/// - any identity field of the captured scope differs from the current active scope;
|
||||
/// - the generation counter has advanced (a workspace switch occurred).
|
||||
///
|
||||
/// Used by `ensure_relay_mesh_for_record` after capturing scope + discovering
|
||||
/// the bootstrap target. Tests inject `install` directly so no port is touched.
|
||||
pub(crate) async fn install_client_under_workspace_transition<R, I, Fut>(
|
||||
app: &AppHandle<R>,
|
||||
captured_scope: &crate::managed_agents::scope::WorkspaceAgentScope,
|
||||
install: I,
|
||||
) -> Result<(), String>
|
||||
where
|
||||
R: tauri::Runtime,
|
||||
I: FnOnce() -> Fut,
|
||||
Fut: Future<Output = Result<(), String>>,
|
||||
{
|
||||
let state = app.state::<AppState>();
|
||||
let _transition_guard = state.workspace_transition.lock().await;
|
||||
|
||||
// Validate the captured scope's generation and full identity under the guard.
|
||||
// If the workspace switched after discovery but before we acquired the lock,
|
||||
// abort without invoking install.
|
||||
crate::managed_agents::scope::validate_scope_generation(captured_scope)
|
||||
.map_err(|e| format!("mesh client install: captured scope stale: {e}"))?;
|
||||
|
||||
let active = state.capture_active_scope().ok_or(
|
||||
"mesh client install: no active workspace scope after lock acquisition".to_string(),
|
||||
)?;
|
||||
|
||||
// Validate the full scope identity — not just the generation counter.
|
||||
let normalize = crate::managed_agents::scope::normalize_relay_for_scope;
|
||||
if active.scope_id != captured_scope.scope_id
|
||||
|| normalize(&active.relay_url) != normalize(&captured_scope.relay_url)
|
||||
|| active.owner_pubkey != captured_scope.owner_pubkey
|
||||
{
|
||||
return Err(format!(
|
||||
"mesh client install: captured scope identity mismatch \
|
||||
(captured scope_id={}, relay={}, owner={}; \
|
||||
active scope_id={}, relay={}, owner={})",
|
||||
captured_scope.scope_id,
|
||||
captured_scope.relay_url,
|
||||
captured_scope.owner_pubkey,
|
||||
active.scope_id,
|
||||
active.relay_url,
|
||||
active.owner_pubkey,
|
||||
));
|
||||
}
|
||||
|
||||
install().await
|
||||
}
|
||||
|
||||
/// Stop the local Mesh **client** (consuming) runtime.
|
||||
///
|
||||
/// Only tears down a client-mode runtime. Serve-mode and absent runtimes are
|
||||
|
||||
@@ -308,7 +308,7 @@ pub async fn set_persona_active(
|
||||
|
||||
pub(crate) const PNG_MAGIC: [u8; 4] = [0x89, 0x50, 0x4E, 0x47];
|
||||
mod card;
|
||||
mod snapshot;
|
||||
pub(crate) mod snapshot;
|
||||
pub use card::*;
|
||||
#[cfg(test)]
|
||||
pub(crate) use snapshot::import::decode_snapshot_from_bytes;
|
||||
|
||||
@@ -6,9 +6,10 @@
|
||||
//! registered in `lib.rs` through the same `personas::` path as the export
|
||||
//! commands.
|
||||
|
||||
use futures_util::future::BoxFuture;
|
||||
use nostr::ToBech32;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tauri::{AppHandle, Emitter, State};
|
||||
use tauri::{AppHandle, Emitter, Manager, State};
|
||||
|
||||
use crate::{
|
||||
app_state::AppState,
|
||||
@@ -24,6 +25,29 @@ use crate::{
|
||||
util::now_iso,
|
||||
};
|
||||
|
||||
// ── Outbound adapter arg structs ──────────────────────────────────────────────
|
||||
|
||||
/// Arguments passed to an injected profile-publish callback.
|
||||
///
|
||||
/// Borrows all fields to avoid cloning `nostr::Keys` across closures.
|
||||
pub(crate) struct ProfilePublish<'a> {
|
||||
pub relay_url: &'a str,
|
||||
pub agent_keys: &'a nostr::Keys,
|
||||
pub display_name: &'a str,
|
||||
pub avatar_url: Option<&'a str>,
|
||||
pub auth_tag: Option<&'a str>,
|
||||
}
|
||||
|
||||
/// Arguments passed to an injected engram-submit callback.
|
||||
///
|
||||
/// Borrows all fields to avoid cloning `nostr::Keys` across closures.
|
||||
pub(crate) struct MemoryPublish<'a> {
|
||||
pub relay_url: &'a str,
|
||||
pub event_json: &'a [u8],
|
||||
pub agent_keys: &'a nostr::Keys,
|
||||
pub auth_tag: Option<&'a str>,
|
||||
}
|
||||
|
||||
/// Maximum snapshot file size accepted before decode (5 MiB for JSON,
|
||||
/// 10 MiB for PNG). Mirrors the established persona-import limits.
|
||||
pub(crate) const MAX_SNAPSHOT_JSON_BYTES: usize = 5 * 1024 * 1024;
|
||||
@@ -407,40 +431,40 @@ pub(crate) use import_entry::capture_agent_snapshot_import_entry;
|
||||
|
||||
// ── `confirm_agent_snapshot_import` ──────────────────────────────────────────
|
||||
|
||||
/// Import a `buzz-agent-snapshot v1` file as a brand-new agent.
|
||||
/// Testable core of [`confirm_agent_snapshot_import`].
|
||||
///
|
||||
/// Phase sequence:
|
||||
/// 1. Validate — decode the manifest and reject early on any error.
|
||||
/// 2. Mint — generate a new keypair + NIP-OA auth tag; create a
|
||||
/// `AgentDefinition` + `ManagedAgentRecord` through the same primitives
|
||||
/// used by the normal create flow.
|
||||
/// 3. Publish — kind:30175 definition via retention path; kind:0 profile
|
||||
/// via `sync_managed_agent_profile`.
|
||||
/// 4. Memory — for each opted-in entry, build a fresh `kind:30174` event
|
||||
/// with `engram::build_event` under the new agent↔owner conversation
|
||||
/// key and POST it to the relay. Failures are collected and returned as
|
||||
/// `memory_errors`; the agent itself is already created.
|
||||
/// `before_store` — called after entry capture, immediately before Phase 3a
|
||||
/// acquires `managed_agents_store_lock`. Used in tests to inject a concurrent
|
||||
/// workspace switch; no-op in production.
|
||||
///
|
||||
/// Importing the same file twice yields two distinct agents with different
|
||||
/// keypairs. No source identity material (pubkey, nsec, auth_tag, relay_url,
|
||||
/// env_vars, backend, lineage) is consumed.
|
||||
#[tauri::command]
|
||||
pub async fn confirm_agent_snapshot_import(
|
||||
/// `after_store` — called after Phase 3a releases `managed_agents_store_lock`,
|
||||
/// immediately before Phase 3b's first outbound call. Used in tests to prove
|
||||
/// Phase 3b reads captured variables, not live state; no-op in production.
|
||||
///
|
||||
/// `profile_sync` and `submit_memory` are the outbound adapters; production
|
||||
/// passes real relay calls while tests inject assertions over captured fields.
|
||||
pub(crate) async fn confirm_agent_snapshot_import_core<R, Before, After, Profile, Memory>(
|
||||
input: AgentSnapshotImportConfirm,
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<AgentSnapshotImportResult, String> {
|
||||
// Capture the active scope and verify owner-key agreement at entry.
|
||||
// `capture_agent_snapshot_import_entry` is the production boundary guard;
|
||||
// it is also called directly by unit tests.
|
||||
let entry = capture_agent_snapshot_import_entry(&state)?;
|
||||
app: &tauri::AppHandle<R>,
|
||||
state: &AppState,
|
||||
before_store: Before,
|
||||
after_store: After,
|
||||
profile_sync: Profile,
|
||||
submit_memory: Memory,
|
||||
) -> Result<AgentSnapshotImportResult, String>
|
||||
where
|
||||
R: tauri::Runtime,
|
||||
Before: Fn() + Send + Sync,
|
||||
After: Fn() + Send + Sync,
|
||||
Profile: for<'a> Fn(ProfilePublish<'a>) -> BoxFuture<'a, Result<(), String>>,
|
||||
Memory: for<'a> Fn(MemoryPublish<'a>) -> BoxFuture<'a, Result<(), String>>,
|
||||
{
|
||||
let entry = capture_agent_snapshot_import_entry(state)?;
|
||||
let captured_scope = entry.captured_scope;
|
||||
let captured_owner_keys = entry.captured_owner_keys;
|
||||
let definitions_dir = captured_scope.definitions_dir.clone();
|
||||
|
||||
// ── Phase 1: validate (no writes) ────────────────────────────────────────
|
||||
// Locked cards unlock only via this machine's exact key endpoints;
|
||||
// anything else fails closed here, before key generation.
|
||||
let snapshot = {
|
||||
let records = {
|
||||
let _store_guard = state
|
||||
@@ -457,7 +481,6 @@ pub async fn confirm_agent_snapshot_import(
|
||||
return Err("Snapshot display name is empty.".to_string());
|
||||
}
|
||||
|
||||
// ── Resolve behavioral defaults ──────────────────────────────────────────
|
||||
let minted = resolve_snapshot_import_behavior(
|
||||
snapshot.definition.respond_to.as_deref(),
|
||||
&snapshot.definition.respond_to_allowlist,
|
||||
@@ -466,15 +489,11 @@ pub async fn confirm_agent_snapshot_import(
|
||||
)?;
|
||||
let minted_parallelism = minted.parallelism;
|
||||
|
||||
// Profile metadata must contain a hosted URL. Inline avatar data can be far
|
||||
// larger than the relay's kind:0 content limit, so upload imported pixels
|
||||
// before minting or persisting the new agent. Failing here keeps import
|
||||
// atomic instead of creating an agent whose profile can never publish.
|
||||
let effective_avatar = materialize_import_avatar(
|
||||
snapshot.profile.avatar_data_url.as_deref(),
|
||||
snapshot.profile.avatar_url.as_deref(),
|
||||
|avatar_bytes| async {
|
||||
crate::commands::media::upload_image_bytes(avatar_bytes, &state)
|
||||
crate::commands::media::upload_image_bytes(avatar_bytes, state)
|
||||
.await
|
||||
.map(|descriptor| descriptor.url)
|
||||
.map_err(|error| format!("Could not upload the imported avatar: {error}"))
|
||||
@@ -482,15 +501,13 @@ pub async fn confirm_agent_snapshot_import(
|
||||
)
|
||||
.await?;
|
||||
|
||||
// Wire-format string for the persona definition's respond_to field.
|
||||
// Omit when it is the default (owner-only) to keep definitions clean.
|
||||
let respond_to_wire: Option<String> = if minted.respond_to != RespondTo::default() {
|
||||
Some(minted.respond_to.as_str().to_string())
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
// ── Phase 2: mint keys + auth tag (sync, outside lock) ───────────────────
|
||||
// ── Phase 2: mint keys + auth tag ────────────────────────────────────────
|
||||
let (agent_keys, private_key_nsec, pubkey, auth_tag, owner_pubkey_hex) = {
|
||||
let agent_keys = nostr::Keys::generate();
|
||||
let pubkey = agent_keys.public_key().to_hex();
|
||||
@@ -498,8 +515,6 @@ pub async fn confirm_agent_snapshot_import(
|
||||
.secret_key()
|
||||
.to_bech32()
|
||||
.map_err(|e| format!("failed to encode agent private key: {e}"))?;
|
||||
|
||||
// NIP-OA auth tag: bridge nostr 0.37 → 0.36 (buzz-sdk) via hex round-trip.
|
||||
let compat_owner = nostr::Keys::parse(&captured_owner_keys.secret_key().to_secret_hex())
|
||||
.map_err(|e| format!("failed to bridge owner keys: {e}"))?;
|
||||
let compat_agent = nostr::PublicKey::from_hex(&pubkey)
|
||||
@@ -518,25 +533,23 @@ pub async fn confirm_agent_snapshot_import(
|
||||
)
|
||||
};
|
||||
|
||||
// ── Phase 3a: create AgentDefinition + ManagedAgentRecord (sync lock) ──────
|
||||
// ── Phase 3a: create AgentDefinition + ManagedAgentRecord (sync lock) ────
|
||||
// `before_store` fires after entry capture and before lock acquisition so
|
||||
// a test-injected workspace switch arrives here — not via stale entry setup.
|
||||
before_store();
|
||||
let (persona, record) = {
|
||||
let _store_guard = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
// Re-validate the captured scope's generation before any write; abort
|
||||
// if the workspace switched since Phase 1.
|
||||
crate::managed_agents::scope::validate_scope_generation(&captured_scope)
|
||||
.map_err(|e| format!("confirm_agent_snapshot_import: {e}"))?;
|
||||
|
||||
// Re-verify owner key under lock — concurrent identity swap must not split writes.
|
||||
if captured_owner_keys.public_key().to_hex() != captured_scope.owner_pubkey {
|
||||
return Err("confirm_agent_snapshot_import: owner key mismatch under lock".to_string());
|
||||
}
|
||||
|
||||
// Resolve retention scope from the captured scope so retain calls
|
||||
// write to the correct scope's DB even if live state diverged.
|
||||
let retention_scope = crate::managed_agents::retention::retention_scope_from_captured(
|
||||
&captured_scope,
|
||||
captured_owner_keys.clone(),
|
||||
@@ -545,7 +558,6 @@ pub async fn confirm_agent_snapshot_import(
|
||||
let mut personas = crate::managed_agents::load_personas_at(&definitions_dir)?;
|
||||
let mut records = crate::managed_agents::storage::load_managed_agents_at(&definitions_dir)?;
|
||||
|
||||
// Guard against duplicate pubkey (astronomically unlikely but safe).
|
||||
if records.iter().any(|r| r.pubkey == pubkey) {
|
||||
return Err(format!("generated pubkey {pubkey} already exists — retry"));
|
||||
}
|
||||
@@ -553,7 +565,6 @@ pub async fn confirm_agent_snapshot_import(
|
||||
let now = now_iso();
|
||||
let persona_id = uuid::Uuid::new_v4().to_string();
|
||||
|
||||
// Build persona from snapshot definition.
|
||||
let persona = AgentDefinition {
|
||||
id: persona_id.clone(),
|
||||
display_name: display_name.clone(),
|
||||
@@ -583,12 +594,8 @@ pub async fn confirm_agent_snapshot_import(
|
||||
|
||||
personas.push(persona.clone());
|
||||
crate::managed_agents::save_personas_at(&definitions_dir, &personas)?;
|
||||
|
||||
// Enqueue the kind:30175 persona event via the retention path.
|
||||
super::super::pending::retain_persona_pending_in_scope(&retention_scope, &persona);
|
||||
|
||||
// Build the managed agent record — no machine-local commands, no
|
||||
// secrets, no lineage from the snapshot.
|
||||
let record = ManagedAgentRecord {
|
||||
pubkey: pubkey.clone(),
|
||||
name: display_name.clone(),
|
||||
@@ -597,10 +604,8 @@ pub async fn confirm_agent_snapshot_import(
|
||||
persona_id: Some(persona_id.clone()),
|
||||
private_key_nsec: private_key_nsec.clone(),
|
||||
auth_tag: auth_tag.clone(),
|
||||
relay_url: String::new(), // resolves to workspace relay at runtime
|
||||
relay_url: String::new(),
|
||||
avatar_url: effective_avatar.clone(),
|
||||
// Machine-local commands: derive from the runtime catalog at
|
||||
// spawn time — never manufacture from snapshot data.
|
||||
acp_command: crate::managed_agents::DEFAULT_ACP_COMMAND.to_string(),
|
||||
agent_command: String::new(),
|
||||
agent_command_override: None,
|
||||
@@ -632,9 +637,6 @@ pub async fn confirm_agent_snapshot_import(
|
||||
last_exit_code: None,
|
||||
last_error: None,
|
||||
last_error_code: None,
|
||||
// Instance-level behavioral defaults agree with the resolved
|
||||
// definition: both come from the single minted struct so they
|
||||
// are always consistent at mint time.
|
||||
respond_to: minted.respond_to,
|
||||
respond_to_allowlist: minted.respond_to_allowlist.clone(),
|
||||
is_builtin: false,
|
||||
@@ -653,33 +655,25 @@ pub async fn confirm_agent_snapshot_import(
|
||||
|
||||
records.push(record.clone());
|
||||
crate::managed_agents::storage::save_managed_agents_at(&definitions_dir, &records)?;
|
||||
|
||||
// Enqueue the kind:30177 managed-agent event via retention.
|
||||
// (Uses the same pattern as agents.rs::retain_managed_agent_pending
|
||||
// inlined here to avoid cross-module private-fn access.)
|
||||
retain_agent_pending(&retention_scope, &record);
|
||||
|
||||
crate::managed_agents::try_regenerate_nest(&app).ok();
|
||||
|
||||
// Notify other mounted clients of local persona+managed-agent writes,
|
||||
// matching the contract used by other local managed-agent mutations.
|
||||
crate::managed_agents::try_regenerate_nest(app).ok();
|
||||
let _ = app.emit("agents-data-changed", ());
|
||||
|
||||
(persona, record)
|
||||
};
|
||||
// Phase 3a lock released. `after_store` fires before Phase 3b so a test
|
||||
// can advance scope generation and verify Phase 3b still reads captured vars.
|
||||
after_store();
|
||||
|
||||
// ── Phase 3b: publish kind:0 profile (async, outside lock) ───────────────
|
||||
// Use the captured scope's relay URL so profile publication targets the
|
||||
// same workspace where the definition was written in Phase 3a.
|
||||
let relay_url = effective_agent_relay_url(&record.relay_url, &captured_scope.relay_url);
|
||||
let profile_sync_error = sync_managed_agent_profile(
|
||||
&state,
|
||||
&relay_url,
|
||||
&agent_keys,
|
||||
&display_name,
|
||||
effective_avatar.as_deref(),
|
||||
auth_tag.as_deref(),
|
||||
)
|
||||
let profile_sync_error = profile_sync(ProfilePublish {
|
||||
relay_url: &relay_url,
|
||||
agent_keys: &agent_keys,
|
||||
display_name: &display_name,
|
||||
avatar_url: effective_avatar.as_deref(),
|
||||
auth_tag: auth_tag.as_deref(),
|
||||
})
|
||||
.await
|
||||
.err();
|
||||
|
||||
@@ -691,9 +685,6 @@ pub async fn confirm_agent_snapshot_import(
|
||||
if memory_total > 0 {
|
||||
let owner_pubkey = nostr::PublicKey::from_hex(&owner_pubkey_hex)
|
||||
.map_err(|e| format!("failed to parse owner pubkey: {e}"))?;
|
||||
|
||||
// Monotonic timestamp seed: use current time, bumped by 1 per entry
|
||||
// so no two events land at the same second.
|
||||
let base_ts = nostr::Timestamp::now().as_secs();
|
||||
|
||||
for (idx, entry) in snapshot.memory.entries.iter().enumerate() {
|
||||
@@ -714,13 +705,12 @@ pub async fn confirm_agent_snapshot_import(
|
||||
Ok(event) => {
|
||||
let event_json = nostr::JsonUtil::as_json(&event).into_bytes();
|
||||
let url = format!("{}/events", crate::relay::relay_http_base_url(&relay_url));
|
||||
match submit_engram_event(
|
||||
&state,
|
||||
&agent_keys,
|
||||
&event_json,
|
||||
&url,
|
||||
auth_tag.as_deref(),
|
||||
)
|
||||
match submit_memory(MemoryPublish {
|
||||
relay_url: &url,
|
||||
event_json: &event_json,
|
||||
agent_keys: &agent_keys,
|
||||
auth_tag: auth_tag.as_deref(),
|
||||
})
|
||||
.await
|
||||
{
|
||||
Ok(()) => memory_written += 1,
|
||||
@@ -745,6 +735,62 @@ pub async fn confirm_agent_snapshot_import(
|
||||
})
|
||||
}
|
||||
|
||||
/// Import a `buzz-agent-snapshot v1` file as a brand-new agent.
|
||||
///
|
||||
/// Thin Tauri command: no-op boundary hooks, real outbound adapters.
|
||||
/// See [`confirm_agent_snapshot_import_core`] for the testable logic.
|
||||
#[tauri::command]
|
||||
pub async fn confirm_agent_snapshot_import(
|
||||
input: AgentSnapshotImportConfirm,
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<AgentSnapshotImportResult, String> {
|
||||
// Clone `app` for the closures so they can obtain a `'static` state handle
|
||||
// via `app_clone.state::<AppState>()` without borrowing the command's
|
||||
// local `State<'_, AppState>`.
|
||||
let app_for_profile = app.clone();
|
||||
let app_for_memory = app.clone();
|
||||
confirm_agent_snapshot_import_core(
|
||||
input,
|
||||
&app,
|
||||
&state,
|
||||
|| {},
|
||||
|| {},
|
||||
move |p| {
|
||||
let app = app_for_profile.clone();
|
||||
let relay = p.relay_url.to_string();
|
||||
let keys = p.agent_keys.clone();
|
||||
let name = p.display_name.to_string();
|
||||
let avatar = p.avatar_url.map(str::to_string);
|
||||
let auth = p.auth_tag.map(str::to_string);
|
||||
Box::pin(async move {
|
||||
let s = app.state::<AppState>();
|
||||
sync_managed_agent_profile(
|
||||
&s,
|
||||
&relay,
|
||||
&keys,
|
||||
&name,
|
||||
avatar.as_deref(),
|
||||
auth.as_deref(),
|
||||
)
|
||||
.await
|
||||
})
|
||||
},
|
||||
move |m| {
|
||||
let app = app_for_memory.clone();
|
||||
let url = m.relay_url.to_string();
|
||||
let json = m.event_json.to_vec();
|
||||
let keys = m.agent_keys.clone();
|
||||
let auth = m.auth_tag.map(str::to_string);
|
||||
Box::pin(async move {
|
||||
let s = app.state::<AppState>();
|
||||
submit_engram_event(&s, &keys, &json, &url, auth.as_deref()).await
|
||||
})
|
||||
},
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Inline retention for the managed-agent kind:30177 event — mirrors
|
||||
/// `agents::retain_managed_agent_pending` without requiring cross-module
|
||||
/// private function access.
|
||||
@@ -857,139 +903,5 @@ pub(crate) async fn submit_engram_event(
|
||||
// ── NIP-49 egress guard: boundary 7 (persona snapshot engram submit) ─────────
|
||||
|
||||
#[cfg(test)]
|
||||
mod egress_guard_tests {
|
||||
use super::submit_engram_event;
|
||||
|
||||
const NCRYPTSEC: &str = "ncryptsec1qgg9947rlpvqu76pj5ecreduf9jxhselq2nae2kghhvd5g7dgjtcxfqtd67p9m0w57lspw8gsq6yphnm8623nsl8xn9j4jdzz84zm3frztj3z7s35vpzmqf6ksu8r89qk5z2zxfmu5gv8th8wclt0h4p";
|
||||
|
||||
/// An engram body carrying an ncryptsec must be rejected by the guard
|
||||
/// before any network I/O (the target port is a discard address; a guard
|
||||
/// error — not a connection error — proves the abort ordering).
|
||||
#[tokio::test]
|
||||
async fn blocks_ncryptsec_before_network() {
|
||||
let state = crate::app_state::build_app_state();
|
||||
let keys = nostr::Keys::generate();
|
||||
let body = format!("{{\"content\":\"{NCRYPTSEC}\"}}");
|
||||
let err = submit_engram_event(
|
||||
&state,
|
||||
&keys,
|
||||
body.as_bytes(),
|
||||
"http://127.0.0.1:9/events",
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(err.contains("key-backup material"), "{err}");
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod import_avatar_tests {
|
||||
use super::materialize_import_avatar;
|
||||
use std::cell::Cell;
|
||||
|
||||
#[tokio::test]
|
||||
async fn inline_avatar_is_uploaded_and_replaced_with_hosted_url() {
|
||||
let uploaded = Cell::new(false);
|
||||
let result = materialize_import_avatar(
|
||||
Some("data:image/png;base64,iVBORw0KGgo="),
|
||||
Some("https://sender.invalid/avatar.png"),
|
||||
|bytes| {
|
||||
uploaded.set(true);
|
||||
async move {
|
||||
assert_eq!(bytes, b"\x89PNG\r\n\x1a\n");
|
||||
Ok("https://relay.example/media/avatar.png".to_string())
|
||||
}
|
||||
},
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert!(uploaded.get());
|
||||
assert_eq!(
|
||||
result.as_deref(),
|
||||
Some("https://relay.example/media/avatar.png")
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn hosted_avatar_skips_upload() {
|
||||
let result =
|
||||
materialize_import_avatar(None, Some("https://sender.example/avatar.png"), |_| async {
|
||||
panic!("hosted avatars must not be uploaded")
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(result.as_deref(), Some("https://sender.example/avatar.png"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn relay_sized_inline_avatar_becomes_bounded_signed_profile() {
|
||||
use base64::{engine::general_purpose::STANDARD, Engine};
|
||||
use image::ImageEncoder;
|
||||
use nostr::JsonUtil;
|
||||
|
||||
let mut pixels = vec![0_u8; 512 * 512 * 4];
|
||||
let mut seed = 0x1234_5678_u32;
|
||||
for byte in &mut pixels {
|
||||
seed ^= seed << 13;
|
||||
seed ^= seed >> 17;
|
||||
seed ^= seed << 5;
|
||||
*byte = seed as u8;
|
||||
}
|
||||
let mut source = Vec::new();
|
||||
image::codecs::png::PngEncoder::new(&mut source)
|
||||
.write_image(&pixels, 512, 512, image::ExtendedColorType::Rgba8)
|
||||
.unwrap();
|
||||
assert!(source.len() > 256 * 1024);
|
||||
let data_url = format!("data:image/png;base64,{}", STANDARD.encode(&source));
|
||||
assert!(data_url.len() > 256 * 1024);
|
||||
|
||||
let avatar = materialize_import_avatar(Some(&data_url), None, |bytes| async move {
|
||||
let mime = crate::commands::media::detect_and_validate_mime(&bytes)?;
|
||||
assert_eq!(mime, "image/png");
|
||||
let sanitized = crate::commands::media::sanitize_image_for_upload(bytes, &mime)?;
|
||||
image::load_from_memory(&sanitized).map_err(|error| error.to_string())?;
|
||||
Ok("https://relay.example/media/avatar.png".to_string())
|
||||
})
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
|
||||
let event =
|
||||
crate::events::build_profile(Some("Imported agent"), None, Some(&avatar), None, None)
|
||||
.unwrap()
|
||||
.sign_with_keys(&nostr::Keys::generate())
|
||||
.unwrap();
|
||||
assert!(event.content.len() < 64 * 1024);
|
||||
assert!(!event.content.contains("data:image/"));
|
||||
assert!(event
|
||||
.content
|
||||
.contains("https://relay.example/media/avatar.png"));
|
||||
assert!(event.as_json().len() < 256 * 1024);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn upload_failure_aborts_avatar_materialization() {
|
||||
let result = materialize_import_avatar(
|
||||
Some("data:image/png;base64,iVBORw0KGgo="),
|
||||
None,
|
||||
|_| async { Err("relay upload failed".to_string()) },
|
||||
)
|
||||
.await;
|
||||
|
||||
assert_eq!(result.unwrap_err(), "relay upload failed");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn malformed_inline_avatar_fails_before_upload() {
|
||||
let result =
|
||||
materialize_import_avatar(Some("data:image/png;base64,not-base64!"), None, |_| async {
|
||||
panic!("malformed avatars must not be uploaded")
|
||||
})
|
||||
.await;
|
||||
|
||||
assert_eq!(result.unwrap_err(), "Snapshot avatar data is malformed.");
|
||||
}
|
||||
}
|
||||
#[path = "import_tests.rs"]
|
||||
mod tests;
|
||||
|
||||
@@ -0,0 +1,389 @@
|
||||
//! Tests for `confirm_agent_snapshot_import` seams and helpers.
|
||||
//!
|
||||
//! Extracted from `import.rs` to keep that file within the 1000-line gate.
|
||||
//! Included via `#[path]` from `import.rs`.
|
||||
|
||||
use super::{materialize_import_avatar, submit_engram_event};
|
||||
|
||||
const NCRYPTSEC: &str = "ncryptsec1qgg9947rlpvqu76pj5ecreduf9jxhselq2nae2kghhvd5g7dgjtcxfqtd67p9m0w57lspw8gsq6yphnm8623nsl8xn9j4jdzz84zm3frztj3z7s35vpzmqf6ksu8r89qk5z2zxfmu5gv8th8wclt0h4p";
|
||||
|
||||
/// An engram body carrying an ncryptsec must be rejected by the guard
|
||||
/// before any network I/O (the target port is a discard address; a guard
|
||||
/// error — not a connection error — proves the abort ordering).
|
||||
#[tokio::test]
|
||||
async fn blocks_ncryptsec_before_network() {
|
||||
let state = crate::app_state::build_app_state();
|
||||
let keys = nostr::Keys::generate();
|
||||
let body = format!("{{\"content\":\"{NCRYPTSEC}\"}}");
|
||||
let err = submit_engram_event(
|
||||
&state,
|
||||
&keys,
|
||||
body.as_bytes(),
|
||||
"http://127.0.0.1:9/events",
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(err.contains("key-backup material"), "{err}");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn inline_avatar_is_uploaded_and_replaced_with_hosted_url() {
|
||||
let uploaded = std::cell::Cell::new(false);
|
||||
let result = materialize_import_avatar(
|
||||
Some("data:image/png;base64,iVBORw0KGgo="),
|
||||
Some("https://sender.invalid/avatar.png"),
|
||||
|bytes| {
|
||||
uploaded.set(true);
|
||||
async move {
|
||||
assert_eq!(bytes, b"\x89PNG\r\n\x1a\n");
|
||||
Ok("https://relay.example/media/avatar.png".to_string())
|
||||
}
|
||||
},
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert!(uploaded.get());
|
||||
assert_eq!(
|
||||
result.as_deref(),
|
||||
Some("https://relay.example/media/avatar.png")
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn hosted_avatar_skips_upload() {
|
||||
let result =
|
||||
materialize_import_avatar(None, Some("https://sender.example/avatar.png"), |_| async {
|
||||
panic!("hosted avatars must not be uploaded")
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(result.as_deref(), Some("https://sender.example/avatar.png"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn relay_sized_inline_avatar_becomes_bounded_signed_profile() {
|
||||
use base64::{engine::general_purpose::STANDARD, Engine};
|
||||
use image::ImageEncoder;
|
||||
use nostr::JsonUtil;
|
||||
|
||||
let mut pixels = vec![0_u8; 512 * 512 * 4];
|
||||
let mut seed = 0x1234_5678_u32;
|
||||
for byte in &mut pixels {
|
||||
seed ^= seed << 13;
|
||||
seed ^= seed >> 17;
|
||||
seed ^= seed << 5;
|
||||
*byte = seed as u8;
|
||||
}
|
||||
let mut source = Vec::new();
|
||||
image::codecs::png::PngEncoder::new(&mut source)
|
||||
.write_image(&pixels, 512, 512, image::ExtendedColorType::Rgba8)
|
||||
.unwrap();
|
||||
assert!(source.len() > 256 * 1024);
|
||||
let data_url = format!("data:image/png;base64,{}", STANDARD.encode(&source));
|
||||
assert!(data_url.len() > 256 * 1024);
|
||||
|
||||
let avatar = materialize_import_avatar(Some(&data_url), None, |bytes| async move {
|
||||
let mime = crate::commands::media::detect_and_validate_mime(&bytes)?;
|
||||
assert_eq!(mime, "image/png");
|
||||
let sanitized = crate::commands::media::sanitize_image_for_upload(bytes, &mime)?;
|
||||
image::load_from_memory(&sanitized).map_err(|error| error.to_string())?;
|
||||
Ok("https://relay.example/media/avatar.png".to_string())
|
||||
})
|
||||
.await
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
|
||||
let event =
|
||||
crate::events::build_profile(Some("Imported agent"), None, Some(&avatar), None, None)
|
||||
.unwrap()
|
||||
.sign_with_keys(&nostr::Keys::generate())
|
||||
.unwrap();
|
||||
assert!(event.content.len() < 64 * 1024);
|
||||
assert!(!event.content.contains("data:image/"));
|
||||
assert!(event
|
||||
.content
|
||||
.contains("https://relay.example/media/avatar.png"));
|
||||
assert!(event.as_json().len() < 256 * 1024);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn upload_failure_aborts_avatar_materialization() {
|
||||
let result = materialize_import_avatar(
|
||||
Some("data:image/png;base64,iVBORw0KGgo="),
|
||||
None,
|
||||
|_| async { Err("relay upload failed".to_string()) },
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert!(result.unwrap_err().contains("relay upload failed"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn malformed_inline_avatar_fails_before_upload() {
|
||||
let result =
|
||||
materialize_import_avatar(Some("data:image/png;base64,not-base64!"), None, |_| async {
|
||||
panic!("malformed avatars must not be uploaded")
|
||||
})
|
||||
.await;
|
||||
|
||||
assert_eq!(result.unwrap_err(), "Snapshot avatar data is malformed.");
|
||||
}
|
||||
|
||||
// ── Phase-boundary seam tests (Area 3) ───────────────────────────────────────
|
||||
|
||||
/// Shared cross-module serialization lock for generation-sensitive tests.
|
||||
/// See `managed_agents::scope::SCOPE_GENERATION_TEST_LOCK` for full rationale.
|
||||
/// Re-exported here so test functions can reference it without the full path.
|
||||
use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK as GENERATION_TEST_LOCK;
|
||||
|
||||
/// Build a minimal agent snapshot JSON for import tests.
|
||||
fn minimal_agent_snapshot_json(name: &str) -> Vec<u8> {
|
||||
use crate::managed_agents::agent_snapshot::encode_snapshot_json;
|
||||
use crate::managed_agents::agent_snapshot::{
|
||||
AgentSnapshot, AgentSnapshotDefinition, AgentSnapshotMemory, AgentSnapshotProfile,
|
||||
MemoryLevel, FORMAT_DISCRIMINATOR, FORMAT_VERSION,
|
||||
};
|
||||
|
||||
let snap = AgentSnapshot {
|
||||
format: FORMAT_DISCRIMINATOR.to_string(),
|
||||
version: FORMAT_VERSION,
|
||||
definition: AgentSnapshotDefinition {
|
||||
name: name.to_string(),
|
||||
source_is_builtin: false,
|
||||
system_prompt: Some(format!("{name} prompt")),
|
||||
runtime: None,
|
||||
model: None,
|
||||
provider: None,
|
||||
parallelism: None,
|
||||
respond_to: None,
|
||||
respond_to_allowlist: vec![],
|
||||
name_pool: vec![],
|
||||
idle_timeout_seconds: None,
|
||||
max_turn_duration_seconds: None,
|
||||
},
|
||||
profile: AgentSnapshotProfile {
|
||||
display_name: name.to_string(),
|
||||
about: None,
|
||||
avatar_data_url: None,
|
||||
avatar_url: Some(format!("https://example.test/{name}.png")),
|
||||
},
|
||||
memory: AgentSnapshotMemory {
|
||||
level: MemoryLevel::None,
|
||||
entries: vec![],
|
||||
},
|
||||
};
|
||||
encode_snapshot_json(&snap).expect("encode_snapshot_json must succeed for minimal snapshot")
|
||||
}
|
||||
|
||||
/// Set up a mock Tauri `App` with `AppState` managed, an active workspace
|
||||
/// scope, and matching owner keys.
|
||||
///
|
||||
/// Returns the built `App` (keeps state alive for test duration) and the
|
||||
/// generated owner keys. The `App`'s `handle()` is passed as the `app`
|
||||
/// parameter to core functions; `app.state::<AppState>()` gives the state
|
||||
/// reference for setup mutations inside hook closures.
|
||||
///
|
||||
/// Uses `tauri::test::mock_builder().manage(state)` so that `app.state::<AppState>()`
|
||||
/// works inside `try_regenerate_nest` and other AppHandle users called by the core.
|
||||
///
|
||||
/// Uses `next_scope_generation()` to claim the current generation slot so the
|
||||
/// scope's generation matches the global counter at entry, reducing the race
|
||||
/// window vs. tests that call `next_scope_generation()` concurrently.
|
||||
fn setup_import_app_with_scope(
|
||||
tmp: &tempfile::TempDir,
|
||||
) -> (tauri::App<tauri::test::MockRuntime>, nostr::Keys) {
|
||||
use crate::managed_agents::scope::{next_scope_generation, WorkspaceAgentScope};
|
||||
|
||||
let owner_keys = nostr::Keys::generate();
|
||||
let state = crate::app_state::build_app_state();
|
||||
{
|
||||
let mut locked = state.keys.lock().unwrap();
|
||||
*locked = owner_keys.clone();
|
||||
}
|
||||
|
||||
let app = tauri::test::mock_builder()
|
||||
.manage(state)
|
||||
.build(tauri::test::mock_context(tauri::test::noop_assets()))
|
||||
.expect("failed to build mock app for import test");
|
||||
|
||||
{
|
||||
use tauri::Manager;
|
||||
let s = app.state::<crate::app_state::AppState>();
|
||||
// Claim a fresh generation slot: next_scope_generation() increments the
|
||||
// global counter and returns the new value. The active scope uses this
|
||||
// value so capture_agent_snapshot_import_entry sees a matching generation
|
||||
// when it reads current_scope_generation() at entry.
|
||||
let gen = next_scope_generation();
|
||||
let scope = WorkspaceAgentScope {
|
||||
scope_id: "test-scope".to_string(),
|
||||
relay_url: "wss://captured.example".to_string(),
|
||||
owner_pubkey: owner_keys.public_key().to_hex(),
|
||||
definitions_dir: tmp.path().to_path_buf(),
|
||||
generation: gen,
|
||||
};
|
||||
s.commit_active_scope(scope);
|
||||
}
|
||||
|
||||
(app, owner_keys)
|
||||
}
|
||||
|
||||
/// `after_store` hook commits a genuinely different live scope + owner —
|
||||
/// Phase 3b outbound must still use the OLD (captured) relay URL, not the
|
||||
/// new live relay.
|
||||
///
|
||||
/// Thufir requirement: `after_store` must commit a genuinely different live
|
||||
/// scope and owner, not merely increment a counter. We swap the active scope
|
||||
/// to a different relay + fresh owner inside the hook so that if Phase 3b
|
||||
/// ever re-read live state it would see the new relay. The injected profile
|
||||
/// adapter asserts it receives the OLD captured relay, proving Phase 3b is
|
||||
/// scope-independent after Phase 3a completes.
|
||||
///
|
||||
/// The snapshot in this test has no memory entries so the memory adapter is
|
||||
/// never called; the profile-relay assertion is sufficient.
|
||||
#[tokio::test]
|
||||
async fn test_agent_switch_between_store_and_profile_finishes_captured_outbound() {
|
||||
// Serialize against the before_store rejection test to prevent the
|
||||
// concurrent next_scope_generation() bump from causing a spurious
|
||||
// Phase 3a stale-scope failure.
|
||||
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
|
||||
|
||||
use crate::commands::personas::snapshot::import::{
|
||||
confirm_agent_snapshot_import_core, AgentSnapshotImportConfirm, MemoryPublish,
|
||||
ProfilePublish,
|
||||
};
|
||||
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use tauri::Manager;
|
||||
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let (app, _owner_keys) = setup_import_app_with_scope(&tmp);
|
||||
let handle = app.handle();
|
||||
|
||||
let file_bytes = minimal_agent_snapshot_json("TestAgent");
|
||||
let input = AgentSnapshotImportConfirm {
|
||||
file_bytes,
|
||||
keep_allowlist: false,
|
||||
};
|
||||
|
||||
// Relay URL embedded in the captured scope (must appear in profile_sync).
|
||||
let expected_relay = "wss://captured.example".to_string();
|
||||
|
||||
// Track which relay URLs the outbound adapters received.
|
||||
let profile_relay = Arc::new(Mutex::new(None::<String>));
|
||||
let pr = profile_relay.clone();
|
||||
|
||||
// The after_store hook needs to commit a new scope via AppState.
|
||||
// Get the state reference from the app's managed state.
|
||||
let state = app.state::<crate::app_state::AppState>();
|
||||
// Clone the app handle for the hook to use.
|
||||
let handle_for_hook = handle.clone();
|
||||
|
||||
let result = confirm_agent_snapshot_import_core(
|
||||
input,
|
||||
&handle,
|
||||
&state,
|
||||
|| {}, // before_store: no-op
|
||||
move || {
|
||||
// after_store: commit a genuinely DIFFERENT live scope + owner.
|
||||
// This simulates a workspace switch at the Phase-3a→3b boundary.
|
||||
// Phase 3b must still use the OLD captured relay, not this new one.
|
||||
let new_owner = nostr::Keys::generate();
|
||||
let new_scope = WorkspaceAgentScope {
|
||||
scope_id: "switched-scope".to_string(),
|
||||
relay_url: "wss://new-relay-after-switch.example".to_string(),
|
||||
owner_pubkey: new_owner.public_key().to_hex(),
|
||||
definitions_dir: std::path::PathBuf::from("/tmp/switched"),
|
||||
generation: current_scope_generation(),
|
||||
};
|
||||
let s = handle_for_hook.state::<crate::app_state::AppState>();
|
||||
s.commit_active_scope(new_scope);
|
||||
},
|
||||
move |p: ProfilePublish<'_>| {
|
||||
let relay = p.relay_url.to_string();
|
||||
*pr.lock().unwrap() = Some(relay.clone());
|
||||
Box::pin(async move {
|
||||
let _ = relay;
|
||||
Ok(())
|
||||
})
|
||||
},
|
||||
|_m: MemoryPublish<'_>| Box::pin(async { panic!("no memory entries in this snapshot") }),
|
||||
)
|
||||
.await;
|
||||
|
||||
// The import succeeded — agent was written to the captured scope.
|
||||
assert!(result.is_ok(), "import must succeed: {:?}", result.err());
|
||||
|
||||
// Profile adapter received the OLD captured relay URL, not the new live relay.
|
||||
let profile_seen = profile_relay.lock().unwrap().clone();
|
||||
assert_eq!(
|
||||
profile_seen.as_deref(),
|
||||
Some(expected_relay.as_str()),
|
||||
"profile adapter must receive captured relay, got: {profile_seen:?}"
|
||||
);
|
||||
// No memory entries in this snapshot — memory adapter not called.
|
||||
}
|
||||
|
||||
/// `before_store` hook advances scope generation — Phase 3a must reject with
|
||||
/// a generation mismatch error BEFORE any write.
|
||||
///
|
||||
/// This is the identity-switch-before-store-is-rejected test. `before_store`
|
||||
/// fires after entry capture and before lock acquisition — simulating a
|
||||
/// concurrent workspace switch that arrived after `capture_agent_snapshot_import_entry`
|
||||
/// returned but before Phase 3a acquired the store lock.
|
||||
#[tokio::test]
|
||||
async fn test_agent_identity_switch_before_store_is_rejected() {
|
||||
// Serialize against the after_store test to prevent the generation bump
|
||||
// inside before_store from racing Phase 3a of the after_store test.
|
||||
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
|
||||
|
||||
use crate::commands::personas::snapshot::import::{
|
||||
confirm_agent_snapshot_import_core, AgentSnapshotImportConfirm, MemoryPublish,
|
||||
ProfilePublish,
|
||||
};
|
||||
use crate::managed_agents::scope::next_scope_generation;
|
||||
use tauri::Manager;
|
||||
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let (app, _owner_keys) = setup_import_app_with_scope(&tmp);
|
||||
let handle = app.handle();
|
||||
let state = app.state::<crate::app_state::AppState>();
|
||||
|
||||
let file_bytes = minimal_agent_snapshot_json("TestAgent");
|
||||
let input = AgentSnapshotImportConfirm {
|
||||
file_bytes,
|
||||
keep_allowlist: false,
|
||||
};
|
||||
|
||||
let result = confirm_agent_snapshot_import_core(
|
||||
input,
|
||||
&handle,
|
||||
&state,
|
||||
move || {
|
||||
// before_store: advance generation — simulates a workspace switch
|
||||
// that raced the import after entry capture but before Phase 3a lock.
|
||||
next_scope_generation();
|
||||
},
|
||||
|| {},
|
||||
|_p: ProfilePublish<'_>| {
|
||||
Box::pin(async { panic!("profile must not be called: store rejected") })
|
||||
},
|
||||
|_m: MemoryPublish<'_>| {
|
||||
Box::pin(async { panic!("memory must not be called: store rejected") })
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
result.is_err(),
|
||||
"pre-store switch must cause Phase 3a rejection"
|
||||
);
|
||||
let err = result.unwrap_err();
|
||||
assert!(
|
||||
err.contains("stale") || err.contains("generation") || err.contains("mismatch"),
|
||||
"error must describe generation mismatch: {err}"
|
||||
);
|
||||
}
|
||||
@@ -4,13 +4,20 @@
|
||||
//! and `ManagedAgentRecord` for every member plus one `TeamRecord`. Exporting
|
||||
//! optionally includes member memory at the requested level.
|
||||
|
||||
use futures_util::future::BoxFuture;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tauri::{AppHandle, Emitter, State};
|
||||
use tauri::{AppHandle, Emitter, Manager, State};
|
||||
use uuid::Uuid;
|
||||
|
||||
use crate::{
|
||||
app_state::AppState,
|
||||
commands::{export_util::save_bytes_with_dialog, personas::resolve_snapshot_import_behavior},
|
||||
commands::{
|
||||
export_util::save_bytes_with_dialog,
|
||||
personas::{
|
||||
resolve_snapshot_import_behavior,
|
||||
snapshot::import::{MemoryPublish, ProfilePublish},
|
||||
},
|
||||
},
|
||||
managed_agents::team_snapshot::{
|
||||
build_team_snapshot, decode_team_snapshot_json, decode_team_snapshot_png,
|
||||
encode_team_snapshot_json, encode_team_snapshot_png, TeamSnapshot,
|
||||
@@ -502,18 +509,34 @@ pub(crate) use team_snapshot_entry::capture_team_snapshot_import_entry;
|
||||
/// 5. Memory restore — for each member with non-empty snapshot memory,
|
||||
/// publish each entry as a `kind:30174` engram event. Best-effort.
|
||||
///
|
||||
/// Importing the same file twice yields two distinct teams with different
|
||||
/// agent keypairs (same as individual agent import).
|
||||
#[tauri::command]
|
||||
pub async fn confirm_team_snapshot_import(
|
||||
/// Testable core of [`confirm_team_snapshot_import`].
|
||||
///
|
||||
/// `before_store` — called after entry capture, immediately before Phase 3
|
||||
/// acquires `managed_agents_store_lock`. Test-only hook for pre-store switch
|
||||
/// simulation; no-op in production.
|
||||
///
|
||||
/// `after_store` — called after Phase 3 releases `managed_agents_store_lock`,
|
||||
/// immediately before Phase 4 (first outbound call). Used in tests to prove
|
||||
/// Phase 4/5 reads captured variables; no-op in production.
|
||||
///
|
||||
/// `profile_sync` and `submit_memory` are the per-member outbound adapters.
|
||||
pub(crate) async fn confirm_team_snapshot_import_core<R, Before, After, Profile, Memory>(
|
||||
input: TeamSnapshotImportConfirm,
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<TeamSnapshotImportResult, String> {
|
||||
// Capture the active scope and verify owner-key agreement at entry.
|
||||
// `capture_team_snapshot_import_entry` is the production boundary guard;
|
||||
// it is also called directly by unit tests.
|
||||
let entry = capture_team_snapshot_import_entry(&state)?;
|
||||
app: &tauri::AppHandle<R>,
|
||||
state: &AppState,
|
||||
before_store: Before,
|
||||
after_store: After,
|
||||
profile_sync: Profile,
|
||||
submit_memory: Memory,
|
||||
) -> Result<TeamSnapshotImportResult, String>
|
||||
where
|
||||
R: tauri::Runtime,
|
||||
Before: Fn() + Send + Sync,
|
||||
After: Fn() + Send + Sync,
|
||||
Profile: for<'a> Fn(ProfilePublish<'a>) -> BoxFuture<'a, Result<(), String>>,
|
||||
Memory: for<'a> Fn(MemoryPublish<'a>) -> BoxFuture<'a, Result<(), String>>,
|
||||
{
|
||||
let entry = capture_team_snapshot_import_entry(state)?;
|
||||
let captured_scope = entry.captured_scope;
|
||||
let captured_owner_keys = entry.captured_owner_keys;
|
||||
let definitions_dir = captured_scope.definitions_dir.clone();
|
||||
@@ -522,13 +545,11 @@ pub async fn confirm_team_snapshot_import(
|
||||
let snapshot = decode_team_snapshot_from_bytes(&input.file_bytes)?;
|
||||
let now = now_iso();
|
||||
|
||||
// Resolve behavioral defaults for every member before any key generation.
|
||||
let definitions = build_import_definitions(&snapshot, input.keep_allowlist, &now)?;
|
||||
let persona_ids: Vec<String> = definitions.iter().map(|d| d.id.clone()).collect();
|
||||
let imported_team = build_import_team(&snapshot, persona_ids.clone(), &now)?;
|
||||
|
||||
// ── Phase 2: mint keys + auth tags (sync, outside lock) ─────────────────
|
||||
// All mints must succeed before we enter the store. If any fails, zero writes.
|
||||
let owner_pubkey_hex = captured_owner_keys.public_key().to_hex();
|
||||
|
||||
let mut minted: Vec<MintedMember> = Vec::with_capacity(snapshot.members.len());
|
||||
@@ -548,7 +569,6 @@ pub async fn confirm_team_snapshot_import(
|
||||
.to_bech32()
|
||||
.map_err(|e| format!("failed to encode agent private key: {e}"))?
|
||||
};
|
||||
// NIP-OA auth tag: bridge nostr 0.37 → 0.36 (buzz-sdk) via hex round-trip.
|
||||
let compat_owner =
|
||||
nostr::Keys::parse(&captured_owner_keys.secret_key().to_secret_hex())
|
||||
.map_err(|e| format!("failed to bridge owner keys: {e}"))?;
|
||||
@@ -561,7 +581,6 @@ pub async fn confirm_team_snapshot_import(
|
||||
(agent_keys, private_key_nsec, pubkey, auth_tag)
|
||||
};
|
||||
|
||||
// Build the ManagedAgentRecord for this member.
|
||||
let record = ManagedAgentRecord {
|
||||
pubkey: pubkey.clone(),
|
||||
name: display_name.clone(),
|
||||
@@ -638,20 +657,18 @@ pub async fn confirm_team_snapshot_import(
|
||||
}
|
||||
|
||||
// ── Phase 3: store (sync, inside lock) ──────────────────────────────────
|
||||
// `before_store` fires after entry capture and before lock acquisition so
|
||||
// a test-injected workspace switch arrives here — not via stale entry setup.
|
||||
before_store();
|
||||
let team = {
|
||||
let _store_guard = state
|
||||
.managed_agents_store_lock
|
||||
.lock()
|
||||
.map_err(|e| e.to_string())?;
|
||||
|
||||
// Re-validate the captured scope's generation before any write.
|
||||
// If the workspace switched since Phase 1, abort — writing to the
|
||||
// new scope would import into the wrong workspace.
|
||||
crate::managed_agents::scope::validate_scope_generation(&captured_scope)
|
||||
.map_err(|e| format!("confirm_team_snapshot_import: {e}"))?;
|
||||
|
||||
// Re-verify owner key agreement under the store lock so a concurrent
|
||||
// identity import cannot split disk writes across two different owners.
|
||||
if captured_owner_keys.public_key().to_hex() != captured_scope.owner_pubkey {
|
||||
return Err(
|
||||
"confirm_team_snapshot_import: owner key changed before Phase 3 commit".to_string(),
|
||||
@@ -662,7 +679,6 @@ pub async fn confirm_team_snapshot_import(
|
||||
captured_owner_keys.clone(),
|
||||
)?;
|
||||
|
||||
// Guard against duplicate pubkeys (astronomically unlikely).
|
||||
let existing_records = load_managed_agents_at(&definitions_dir)?;
|
||||
for m in &minted {
|
||||
if existing_records.iter().any(|r| r.pubkey == m.pubkey) {
|
||||
@@ -673,10 +689,6 @@ pub async fn confirm_team_snapshot_import(
|
||||
}
|
||||
}
|
||||
|
||||
// Snapshot both store files for rollback on partial write failure.
|
||||
// Distinguish "file exists with content" from "file absent" so rollback
|
||||
// can delete a file created by the import rather than leaving orphaned
|
||||
// records.
|
||||
let agents_store_path = managed_agents_store_path_at(&definitions_dir);
|
||||
let agents_store_snapshot = match std::fs::read(&agents_store_path) {
|
||||
Ok(bytes) => Some(bytes),
|
||||
@@ -690,26 +702,17 @@ pub async fn confirm_team_snapshot_import(
|
||||
Err(e) => return Err(format!("failed to snapshot teams store: {e}")),
|
||||
};
|
||||
|
||||
// Pre-read teams via the read-only loader BEFORE any agent commits.
|
||||
// This avoids load_teams()'s write-on-load side effect (teams.rs:165-166
|
||||
// saves whenever the file is absent or built-ins changed). A failure here
|
||||
// aborts cleanly — zero writes have occurred.
|
||||
let mut teams = load_teams_readonly(&teams_store_path)?;
|
||||
|
||||
// Collect minted pubkeys for keyring cleanup on rollback.
|
||||
let minted_pubkeys: Vec<&str> = minted.iter().map(|m| m.pubkey.as_str()).collect();
|
||||
|
||||
// Restore the agent store to pre-import state and clean minted keyring
|
||||
// entries. Returns the original error, extended with rollback details.
|
||||
let rollback_agents = |original_err: String| -> String {
|
||||
let mut errors = vec![original_err];
|
||||
// Clean minted keyring entries.
|
||||
for pubkey in &minted_pubkeys {
|
||||
if let Err(e) = crate::managed_agents::storage::try_delete_agent_key(pubkey) {
|
||||
errors.push(format!("keyring cleanup {pubkey}: {e}"));
|
||||
}
|
||||
}
|
||||
// Restore agent store file.
|
||||
let restore = match &agents_store_snapshot {
|
||||
Some(bytes) => crate::managed_agents::storage::atomic_write_json_restricted(
|
||||
&agents_store_path,
|
||||
@@ -730,7 +733,6 @@ pub async fn confirm_team_snapshot_import(
|
||||
}
|
||||
};
|
||||
|
||||
// Write all definitions.
|
||||
let mut personas = load_personas_at(&definitions_dir)?;
|
||||
for m in &minted {
|
||||
personas.push(m.definition.clone());
|
||||
@@ -739,7 +741,6 @@ pub async fn confirm_team_snapshot_import(
|
||||
return Err(rollback_agents(e));
|
||||
}
|
||||
|
||||
// Write all managed-agent records.
|
||||
let mut records = existing_records;
|
||||
for m in &minted {
|
||||
records.push(m.record.clone());
|
||||
@@ -748,13 +749,9 @@ pub async fn confirm_team_snapshot_import(
|
||||
return Err(rollback_agents(e));
|
||||
}
|
||||
|
||||
// Write the team record. `teams` was pre-loaded via the read-only
|
||||
// loader before any agent commits, so a read/parse failure already
|
||||
// aborted before any phase-3 write. save_teams_at sorts and persists.
|
||||
teams.push(imported_team.clone());
|
||||
if let Err(e) = save_teams_at(&definitions_dir, &teams) {
|
||||
let err = rollback_agents(e);
|
||||
// Also restore teams store.
|
||||
let teams_restore = match &teams_store_snapshot {
|
||||
Some(bytes) => {
|
||||
crate::managed_agents::storage::atomic_write_json(&teams_store_path, bytes)
|
||||
@@ -770,7 +767,6 @@ pub async fn confirm_team_snapshot_import(
|
||||
});
|
||||
}
|
||||
|
||||
// All writes committed — safe to update in-memory state.
|
||||
for m in &minted {
|
||||
crate::commands::personas::retain_persona_pending_in_scope(
|
||||
&retention_scope,
|
||||
@@ -780,17 +776,20 @@ pub async fn confirm_team_snapshot_import(
|
||||
for m in &minted {
|
||||
retain_agent_pending(&retention_scope, &m.record);
|
||||
}
|
||||
crate::commands::teams::retain_team_pending(&app, &state, &imported_team);
|
||||
// Use the captured retention scope — not the live active scope — so
|
||||
// team retention writes to the correct workspace even after a switch.
|
||||
crate::commands::teams::retain_team_pending_in_scope(&retention_scope, &imported_team);
|
||||
|
||||
crate::managed_agents::try_regenerate_nest(&app).ok();
|
||||
crate::managed_agents::try_regenerate_nest(app).ok();
|
||||
let _ = app.emit("agents-data-changed", ());
|
||||
|
||||
imported_team
|
||||
};
|
||||
// Phase 3 lock released. `after_store` fires before Phase 4 so a test can
|
||||
// advance scope generation and verify outbound still reads captured vars.
|
||||
after_store();
|
||||
|
||||
// ── Phase 4 & 5: profile sync + memory restore (async, outside lock) ────
|
||||
// Use the captured scope's relay URL so profile and memory publication
|
||||
// targets the same workspace where definitions were written in Phase 3.
|
||||
let relay_ws: &str = &captured_scope.relay_url;
|
||||
let mut member_results: Vec<TeamSnapshotImportMemberResult> = Vec::with_capacity(minted.len());
|
||||
|
||||
@@ -798,14 +797,13 @@ pub async fn confirm_team_snapshot_import(
|
||||
let relay_url = effective_agent_relay_url(&m.record.relay_url, relay_ws);
|
||||
|
||||
// Phase 4: profile sync (best-effort).
|
||||
let profile_sync_error = sync_managed_agent_profile(
|
||||
&state,
|
||||
&relay_url,
|
||||
&m.agent_keys,
|
||||
&m.display_name,
|
||||
m.effective_avatar.as_deref(),
|
||||
m.auth_tag.as_deref(),
|
||||
)
|
||||
let profile_sync_error = profile_sync(ProfilePublish {
|
||||
relay_url: &relay_url,
|
||||
agent_keys: &m.agent_keys,
|
||||
display_name: &m.display_name,
|
||||
avatar_url: m.effective_avatar.as_deref(),
|
||||
auth_tag: m.auth_tag.as_deref(),
|
||||
})
|
||||
.await
|
||||
.err();
|
||||
|
||||
@@ -843,13 +841,12 @@ pub async fn confirm_team_snapshot_import(
|
||||
let event_json = event.as_json().into_bytes();
|
||||
let url =
|
||||
format!("{}/events", crate::relay::relay_http_base_url(&relay_url));
|
||||
match submit_engram_event(
|
||||
&state,
|
||||
&m.agent_keys,
|
||||
&event_json,
|
||||
&url,
|
||||
m.auth_tag.as_deref(),
|
||||
)
|
||||
match submit_memory(MemoryPublish {
|
||||
relay_url: &url,
|
||||
event_json: &event_json,
|
||||
agent_keys: &m.agent_keys,
|
||||
auth_tag: m.auth_tag.as_deref(),
|
||||
})
|
||||
.await
|
||||
{
|
||||
Ok(()) => memory_written += 1,
|
||||
@@ -881,6 +878,66 @@ pub async fn confirm_team_snapshot_import(
|
||||
})
|
||||
}
|
||||
|
||||
/// Import a `buzz-team-snapshot v1` file as a brand-new team.
|
||||
///
|
||||
/// Thin Tauri command: no-op boundary hooks, real outbound adapters.
|
||||
/// See [`confirm_team_snapshot_import_core`] for the testable logic.
|
||||
#[tauri::command]
|
||||
pub async fn confirm_team_snapshot_import(
|
||||
input: TeamSnapshotImportConfirm,
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<TeamSnapshotImportResult, String> {
|
||||
let app_for_profile = app.clone();
|
||||
let app_for_memory = app.clone();
|
||||
confirm_team_snapshot_import_core(
|
||||
input,
|
||||
&app,
|
||||
&state,
|
||||
|| {},
|
||||
|| {},
|
||||
move |p| {
|
||||
let app = app_for_profile.clone();
|
||||
let relay = p.relay_url.to_string();
|
||||
let keys = p.agent_keys.clone();
|
||||
let name = p.display_name.to_string();
|
||||
let avatar = p.avatar_url.map(str::to_string);
|
||||
let auth = p.auth_tag.map(str::to_string);
|
||||
Box::pin(async move {
|
||||
let s = app.state::<AppState>();
|
||||
sync_managed_agent_profile(
|
||||
&s,
|
||||
&relay,
|
||||
&keys,
|
||||
&name,
|
||||
avatar.as_deref(),
|
||||
auth.as_deref(),
|
||||
)
|
||||
.await
|
||||
})
|
||||
},
|
||||
move |m| {
|
||||
let app = app_for_memory.clone();
|
||||
let url = m.relay_url.to_string();
|
||||
let json = m.event_json.to_vec();
|
||||
let keys = m.agent_keys.clone();
|
||||
let auth = m.auth_tag.map(str::to_string);
|
||||
Box::pin(async move {
|
||||
let s = app.state::<AppState>();
|
||||
crate::commands::personas::snapshot::import::submit_engram_event(
|
||||
&s,
|
||||
&keys,
|
||||
&json,
|
||||
&url,
|
||||
auth.as_deref(),
|
||||
)
|
||||
.await
|
||||
})
|
||||
},
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
/// Inline retention for the managed-agent kind:30177 event — mirrors
|
||||
/// `commands::personas::snapshot::import::retain_agent_pending`.
|
||||
fn retain_agent_pending(
|
||||
@@ -931,63 +988,5 @@ fn retain_agent_pending(
|
||||
}
|
||||
}
|
||||
|
||||
/// POST a pre-built signed engram event to the relay, authenticating as the
|
||||
/// new agent. Mirrors the same helper in `snapshot::import`.
|
||||
pub(crate) async fn submit_engram_event(
|
||||
state: &AppState,
|
||||
agent_keys: &nostr::Keys,
|
||||
event_json: &[u8],
|
||||
url: &str,
|
||||
auth_tag: Option<&str>,
|
||||
) -> Result<(), String> {
|
||||
use crate::relay::build_nip98_auth_header_for_keys;
|
||||
use reqwest::Method;
|
||||
|
||||
crate::egress_guard::assert_no_key_backup_bytes(event_json, "team 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)?;
|
||||
let mut request = state
|
||||
.http_client
|
||||
.post(url)
|
||||
.header("Authorization", auth)
|
||||
.header("Content-Type", "application/json");
|
||||
if let Some(tag) = auth_tag {
|
||||
request = request.header("x-auth-tag", tag);
|
||||
}
|
||||
let response = request
|
||||
.body(event_json.to_vec())
|
||||
.send()
|
||||
.await
|
||||
.map_err(|e| crate::relay::classify_request_error(&e))?;
|
||||
|
||||
if !response.status().is_success() {
|
||||
let msg = crate::relay::relay_error_message(response).await;
|
||||
return Err(format!("relay rejected engram: {msg}"));
|
||||
}
|
||||
|
||||
let body = response
|
||||
.text()
|
||||
.await
|
||||
.map_err(|e| format!("failed to read relay response: {e}"))?;
|
||||
let parsed: serde_json::Value =
|
||||
serde_json::from_str(&body).map_err(|e| format!("relay response not JSON: {e}"))?;
|
||||
let accepted = parsed
|
||||
.get("accepted")
|
||||
.and_then(|v| v.as_bool())
|
||||
.unwrap_or(false);
|
||||
if !accepted {
|
||||
let message = parsed
|
||||
.get("message")
|
||||
.and_then(|v| v.as_str())
|
||||
.unwrap_or("unknown");
|
||||
return Err(format!("relay rejected engram: {message}"));
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests;
|
||||
|
||||
@@ -738,7 +738,7 @@ fn full_rollback_at_teams_boundary_absent_agents_store() {
|
||||
|
||||
mod egress_guard_boundary {
|
||||
use super::super::capture_team_snapshot_import_entry;
|
||||
use super::super::submit_engram_event;
|
||||
use crate::commands::personas::snapshot::import::submit_engram_event;
|
||||
|
||||
const NCRYPTSEC: &str = "ncryptsec1qgg9947rlpvqu76pj5ecreduf9jxhselq2nae2kghhvd5g7dgjtcxfqtd67p9m0w57lspw8gsq6yphnm8623nsl8xn9j4jdzz84zm3frztj3z7s35vpzmqf6ksu8r89qk5z2zxfmu5gv8th8wclt0h4p";
|
||||
|
||||
@@ -891,3 +891,191 @@ mod egress_guard_boundary {
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// ── Phase-boundary seam tests (Area 3 — team) ─────────────────────────────
|
||||
|
||||
/// Serializes tests that modify the process-global scope generation counter.
|
||||
/// See the equivalent comment in `import_tests.rs` for rationale.
|
||||
use crate::managed_agents::scope::SCOPE_GENERATION_TEST_LOCK as GENERATION_TEST_LOCK;
|
||||
|
||||
fn setup_team_import_app_with_scope(
|
||||
tmp: &tempfile::TempDir,
|
||||
) -> (tauri::App<tauri::test::MockRuntime>, nostr::Keys) {
|
||||
use crate::managed_agents::scope::{next_scope_generation, WorkspaceAgentScope};
|
||||
|
||||
let owner_keys = nostr::Keys::generate();
|
||||
let state = crate::app_state::build_app_state();
|
||||
{
|
||||
let mut locked = state.keys.lock().unwrap();
|
||||
*locked = owner_keys.clone();
|
||||
}
|
||||
|
||||
let app = tauri::test::mock_builder()
|
||||
.manage(state)
|
||||
.build(tauri::test::mock_context(tauri::test::noop_assets()))
|
||||
.expect("failed to build mock app for team import test");
|
||||
|
||||
{
|
||||
use tauri::Manager;
|
||||
let s = app.state::<crate::app_state::AppState>();
|
||||
let gen = next_scope_generation();
|
||||
let scope = WorkspaceAgentScope {
|
||||
scope_id: "ts-test-scope".to_string(),
|
||||
relay_url: "wss://captured.example".to_string(),
|
||||
owner_pubkey: owner_keys.public_key().to_hex(),
|
||||
definitions_dir: tmp.path().to_path_buf(),
|
||||
generation: gen,
|
||||
};
|
||||
s.commit_active_scope(scope);
|
||||
}
|
||||
|
||||
(app, owner_keys)
|
||||
}
|
||||
|
||||
/// `after_store` hook commits a genuinely different live scope + owner —
|
||||
/// Phase 4/5 outbound must use the OLD (captured) relay URL, not the new
|
||||
/// live relay.
|
||||
///
|
||||
/// Thufir requirement: `after_store` must commit a genuinely different live
|
||||
/// scope and owner, not merely increment a counter. We swap the active scope
|
||||
/// to a different relay + fresh owner inside the hook; all per-member profile
|
||||
/// adapters must receive the old captured relay URL.
|
||||
#[tokio::test]
|
||||
async fn test_confirm_team_snapshot_import_switch_between_store_and_profile() {
|
||||
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
|
||||
|
||||
use crate::commands::personas::snapshot::import::{MemoryPublish, ProfilePublish};
|
||||
use crate::commands::team_snapshot::confirm_team_snapshot_import_core;
|
||||
use crate::managed_agents::scope::{current_scope_generation, WorkspaceAgentScope};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use tauri::Manager;
|
||||
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let (app, _owner_keys) = setup_team_import_app_with_scope(&tmp);
|
||||
let handle = app.handle();
|
||||
|
||||
let snap = snapshot(vec![member("Alice"), member("Bob")]);
|
||||
let encoded = crate::managed_agents::team_snapshot::encode_team_snapshot_json(&snap).unwrap();
|
||||
let input = TeamSnapshotImportConfirm {
|
||||
file_bytes: encoded,
|
||||
keep_allowlist: false,
|
||||
};
|
||||
|
||||
let expected_relay = "wss://captured.example".to_string();
|
||||
let profile_relays: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(vec![]));
|
||||
let memory_relays: Arc<Mutex<Vec<String>>> = Arc::new(Mutex::new(vec![]));
|
||||
let pr = profile_relays.clone();
|
||||
let mr = memory_relays.clone();
|
||||
|
||||
let state = app.state::<crate::app_state::AppState>();
|
||||
let handle_for_hook = handle.clone();
|
||||
|
||||
let result = confirm_team_snapshot_import_core(
|
||||
input,
|
||||
&handle,
|
||||
&state,
|
||||
|| {},
|
||||
move || {
|
||||
// after_store: commit a genuinely DIFFERENT live scope + owner.
|
||||
let new_owner = nostr::Keys::generate();
|
||||
let new_scope = WorkspaceAgentScope {
|
||||
scope_id: "switched-scope".to_string(),
|
||||
relay_url: "wss://new-relay-after-switch.example".to_string(),
|
||||
owner_pubkey: new_owner.public_key().to_hex(),
|
||||
definitions_dir: std::path::PathBuf::from("/tmp/switched"),
|
||||
generation: current_scope_generation(),
|
||||
};
|
||||
let s = handle_for_hook.state::<crate::app_state::AppState>();
|
||||
s.commit_active_scope(new_scope);
|
||||
},
|
||||
move |p: ProfilePublish<'_>| {
|
||||
let relay = p.relay_url.to_string();
|
||||
pr.lock().unwrap().push(relay.clone());
|
||||
Box::pin(async move {
|
||||
let _ = relay;
|
||||
Ok(())
|
||||
})
|
||||
},
|
||||
move |m: MemoryPublish<'_>| {
|
||||
let relay = m.relay_url.to_string();
|
||||
mr.lock().unwrap().push(relay.clone());
|
||||
Box::pin(async move {
|
||||
let _ = relay;
|
||||
Ok(())
|
||||
})
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
result.is_ok(),
|
||||
"team import must succeed: {:?}",
|
||||
result.err()
|
||||
);
|
||||
|
||||
// Every member's profile adapter received the captured relay URL.
|
||||
let seen = profile_relays.lock().unwrap();
|
||||
assert_eq!(
|
||||
seen.len(),
|
||||
2,
|
||||
"profile adapter must be called once per member"
|
||||
);
|
||||
for relay in seen.iter() {
|
||||
assert_eq!(
|
||||
relay, &expected_relay,
|
||||
"profile adapter must receive captured relay, got: {relay}"
|
||||
);
|
||||
}
|
||||
// No memory entries in these members — memory adapter not called.
|
||||
}
|
||||
|
||||
/// `before_store` hook advances scope generation — Phase 3 must reject BEFORE
|
||||
/// any write and BEFORE any outbound call.
|
||||
#[tokio::test]
|
||||
async fn test_team_identity_switch_before_store_is_rejected() {
|
||||
let _gen_guard = GENERATION_TEST_LOCK.lock().unwrap();
|
||||
|
||||
use crate::commands::personas::snapshot::import::{MemoryPublish, ProfilePublish};
|
||||
use crate::commands::team_snapshot::confirm_team_snapshot_import_core;
|
||||
use crate::managed_agents::scope::next_scope_generation;
|
||||
use tauri::Manager;
|
||||
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let (app, _owner_keys) = setup_team_import_app_with_scope(&tmp);
|
||||
let handle = app.handle();
|
||||
let state = app.state::<crate::app_state::AppState>();
|
||||
|
||||
let snap = snapshot(vec![member("Alice")]);
|
||||
let encoded = crate::managed_agents::team_snapshot::encode_team_snapshot_json(&snap).unwrap();
|
||||
let input = TeamSnapshotImportConfirm {
|
||||
file_bytes: encoded,
|
||||
keep_allowlist: false,
|
||||
};
|
||||
|
||||
let result = confirm_team_snapshot_import_core(
|
||||
input,
|
||||
&handle,
|
||||
&state,
|
||||
move || {
|
||||
next_scope_generation();
|
||||
},
|
||||
|| {},
|
||||
|_p: ProfilePublish<'_>| {
|
||||
Box::pin(async { panic!("profile must not be called: store rejected") })
|
||||
},
|
||||
|_m: MemoryPublish<'_>| {
|
||||
Box::pin(async { panic!("memory must not be called: store rejected") })
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
assert!(
|
||||
result.is_err(),
|
||||
"pre-store switch must cause Phase 3 rejection"
|
||||
);
|
||||
let err = result.unwrap_err();
|
||||
assert!(
|
||||
err.contains("stale") || err.contains("generation") || err.contains("mismatch"),
|
||||
"error must describe generation mismatch: {err}"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -75,6 +75,55 @@ pub(super) fn retain_team_pending(app: &AppHandle, state: &AppState, team: &Team
|
||||
}
|
||||
}
|
||||
|
||||
/// Captured-scope sibling of [`retain_team_pending`].
|
||||
///
|
||||
/// Accepts a pre-built [`RetentionScope`] instead of resolving the live active
|
||||
/// scope. Used by Phase 3a of `confirm_team_snapshot_import` where the
|
||||
/// retention scope has already been built from the captured entry; calling the
|
||||
/// live `retain_team_pending` there would re-resolve the active scope, which
|
||||
/// may have diverged after a workspace switch that occurred after Phase-2.
|
||||
///
|
||||
/// Like its live counterpart, this function is best-effort: errors are logged
|
||||
/// and swallowed so a retention failure never blocks the calling phase.
|
||||
pub(crate) fn retain_team_pending_in_scope(
|
||||
scope: &crate::managed_agents::retention::RetentionScope,
|
||||
team: &TeamRecord,
|
||||
) {
|
||||
use crate::managed_agents::{
|
||||
persona_events::monotonic_created_at,
|
||||
retention::{get_retained_event, open_retention_db, retain_event, RetainedEvent},
|
||||
team_events::build_team_event,
|
||||
};
|
||||
use buzz_core_pkg::kind::KIND_TEAM;
|
||||
use nostr::JsonUtil;
|
||||
|
||||
let result = (|| -> Result<(), String> {
|
||||
let conn = open_retention_db(&scope.db_path)?;
|
||||
let pubkey = scope.owner_keys.public_key().to_hex();
|
||||
let prior =
|
||||
get_retained_event(&conn, KIND_TEAM, &pubkey, &team.id)?.map(|row| row.created_at);
|
||||
let event = build_team_event(team)?
|
||||
.custom_created_at(monotonic_created_at(prior))
|
||||
.sign_with_keys(&scope.owner_keys)
|
||||
.map_err(|e| format!("failed to sign team event: {e}"))?;
|
||||
retain_event(
|
||||
&conn,
|
||||
&RetainedEvent {
|
||||
kind: KIND_TEAM,
|
||||
pubkey,
|
||||
d_tag: team.id.clone(),
|
||||
content: event.content.to_string(),
|
||||
created_at: event.created_at.as_secs() as i64,
|
||||
raw_event: event.as_json(),
|
||||
pending_sync: true,
|
||||
},
|
||||
)
|
||||
})();
|
||||
if let Err(e) = result {
|
||||
eprintln!("buzz-desktop: team-retain-in-scope: {e}");
|
||||
}
|
||||
}
|
||||
|
||||
/// Purge a deleted team's pending row and enqueue a NIP-09 tombstone, both
|
||||
/// inside the `managed_agents_store_lock`-held delete body.
|
||||
///
|
||||
|
||||
@@ -11,7 +11,7 @@
|
||||
//! | 3 | pre-signed path into the boundary-1 funnel | `relay/submit.rs` |
|
||||
//! | 4 | `submit_signed_event_with_keys` | `relay.rs` |
|
||||
//! | 5 | huddle STT publisher | `huddle/pipeline.rs` |
|
||||
//! | 6 | `submit_engram_event` (team snapshot) | `commands/team_snapshot.rs` |
|
||||
//! | 6 | `submit_engram_event` (team snapshot) | `commands/personas/snapshot/import.rs` (shared with boundary 7) |
|
||||
//! | 7 | `submit_engram_event` (persona import) | `commands/personas/snapshot/import.rs` |
|
||||
//! | 8 | native websocket send loop (all webview relay WS) | `native_websocket.rs` |
|
||||
//!
|
||||
|
||||
@@ -243,14 +243,15 @@ const EVENTS_INVENTORY: &[(&str, usize, usize)] = &[
|
||||
("src/relay.rs", 2, 2), // boundaries 2, 4
|
||||
("src/relay/submit.rs", 1, 1), // boundaries 1 + 3 (shared funnel)
|
||||
("src/huddle/pipeline.rs", 1, 1), // boundary 5
|
||||
("src/commands/team_snapshot.rs", 1, 1), // boundary 6
|
||||
("src/commands/personas/snapshot/import.rs", 2, 1), // boundary 7 + its in-file injection-test fixture URL
|
||||
("src/commands/team_snapshot.rs", 1, 0), // boundary 6 guard in import.rs submit_engram_event (shared)
|
||||
("src/commands/personas/snapshot/import.rs", 1, 1), // boundary 6+7 (shared submit_engram_event) — injection test moved to import_tests.rs
|
||||
("src/native_websocket.rs", 0, 2), // boundary 8 (WS frames; no events URL)
|
||||
// Test-only fixtures — no production egress, no guard:
|
||||
("src/relay_admission.rs", 1, 0),
|
||||
("src/archive/mod_tests.rs", 1, 0),
|
||||
("src/managed_agents/persona_events/tests.rs", 1, 0),
|
||||
("src/commands/team_snapshot/tests.rs", 1, 0),
|
||||
("src/commands/personas/snapshot/import_tests.rs", 1, 0), // ncryptsec guard injection test (discard addr)
|
||||
// Mock-relay route in its in-file tests; production publish goes through
|
||||
// the guarded boundary-1 funnel (`submit_signed_event_at_with_keys`).
|
||||
("src/commands/personas/sharing.rs", 1, 0),
|
||||
@@ -417,6 +418,7 @@ fn ncryptsec_handling_is_confined_to_allowlisted_files() {
|
||||
"src/commands/team_snapshot.rs",
|
||||
"src/commands/team_snapshot/tests.rs",
|
||||
"src/commands/personas/snapshot/import.rs",
|
||||
"src/commands/personas/snapshot/import_tests.rs",
|
||||
"src/native_websocket.rs",
|
||||
];
|
||||
|
||||
|
||||
@@ -15,7 +15,7 @@ use crate::relay::relay_ws_url_with_override;
|
||||
use std::fs;
|
||||
use std::io;
|
||||
use std::path::{Path, PathBuf};
|
||||
use tauri::{AppHandle, Manager};
|
||||
use tauri::Manager;
|
||||
|
||||
use crate::managed_agents::discovery::known_skill_dirs;
|
||||
#[cfg(unix)]
|
||||
@@ -670,7 +670,7 @@ pub fn upsert_managed_section(file_path: &Path, new_section_content: &str) -> io
|
||||
Ok(())
|
||||
}
|
||||
|
||||
pub fn regenerate_nest_context(app: &AppHandle) -> Result<(), String> {
|
||||
pub fn regenerate_nest_context<R: tauri::Runtime>(app: &tauri::AppHandle<R>) -> Result<(), String> {
|
||||
let nest = nest_dir().ok_or("cannot resolve home directory for nest")?;
|
||||
let agents_md = nest.join("AGENTS.md");
|
||||
|
||||
@@ -694,7 +694,7 @@ pub fn regenerate_nest_context(app: &AppHandle) -> Result<(), String> {
|
||||
/// All call sites treat regeneration as fire-and-forget — agents run fine with
|
||||
/// a stale AGENTS.md. Returns `Err` when regeneration fails so callers can
|
||||
/// report it as degradation in the workspace-apply result.
|
||||
pub fn try_regenerate_nest(app: &AppHandle) -> Result<(), String> {
|
||||
pub fn try_regenerate_nest<R: tauri::Runtime>(app: &tauri::AppHandle<R>) -> Result<(), String> {
|
||||
regenerate_nest_context(app).map_err(|error| {
|
||||
eprintln!("buzz-desktop: nest context regeneration failed: {error}");
|
||||
error
|
||||
|
||||
@@ -442,3 +442,15 @@ mod tests {
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Process-global mutex that serializes tests touching the process-global
|
||||
/// scope generation counter.
|
||||
///
|
||||
/// Any test that (a) captures a generation and requires it to be stable
|
||||
/// through Phase 3a or (b) calls `next_scope_generation()` inside a hook
|
||||
/// must hold this guard for its entire duration. Tests across modules share
|
||||
/// the same counter so they must share the same serialization primitive.
|
||||
///
|
||||
/// Exposed only under `#[cfg(test)]` to avoid polluting the production API.
|
||||
#[cfg(test)]
|
||||
pub(crate) static SCOPE_GENERATION_TEST_LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
|
||||
|
||||
Reference in New Issue
Block a user