mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): resolve P3 compliance gaps in Areas 3 and 4
Area 4 — capture scope BEFORE discovery in ensure_relay_mesh_for_record. Previously capture_active_scope was called after resolve_mesh_bootstrap_target; a workspace switch between discovery and capture would let the helper install a client against an endpoint discovered under the old scope while validating against the new scope. Move capture_active_scope above the discovery call so the captured scope and discovered endpoint are always from the same workspace. Area 4 — make with_workspace_transition_preflight production-callable. Change the transition body parameter from FnOnce() -> Result<T> (sync) to FnOnce() -> BoxFuture<'static, Result<T>> (async) so production callers that dispatch spawn_blocking and await the result keep the workspace_transition guard alive across the full async body. Extract apply_workspace_body as a standalone async fn; apply_workspace (mesh-llm feature) routes through with_workspace_transition_preflight so the helper is no longer test-only dead code. The cfg_attr(not(test), allow(dead_code)) annotation is removed. Update the two test call sites to pass BoxFuture-returning closures. Area 3 — add memory-bearing fixture and owner-coordinate assertion. Add minimal_agent_snapshot_json_with_memory which carries one core memory entry (MemoryLevel::Core). Rewrite test_agent_switch_between_store_and_profile_finishes_captured_outbound to use this fixture so submit_memory actually fires in Phase 4. The injected MemoryPublish adapter now asserts both that relay_url contains the captured relay URL and that the built engram event's p tag (owner counterpart/coordinate) equals the captured owner's pubkey hex, not the post-switch owner committed in after_store. Add extract_p_tag_from_event_json helper. 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
c29cd651db
commit
48015b838c
@@ -751,6 +751,17 @@ pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
return wait_for_mesh_inference(model_id).await;
|
||||
}
|
||||
|
||||
// No serving configuration exists — genuine consumer-only start.
|
||||
// Capture scope BEFORE discovery so a concurrent workspace switch can be
|
||||
// detected under the install lock. Route through
|
||||
// `install_client_under_workspace_transition` which acquires
|
||||
// `workspace_transition`, validates full scope identity
|
||||
// (scope_id, relay, owner, generation), then calls the install closure —
|
||||
// serialized against apply_workspace and live identity import.
|
||||
let captured_scope = state
|
||||
.capture_active_scope()
|
||||
.ok_or("mesh client install: no active workspace scope")?;
|
||||
|
||||
let target = match resolve_mesh_bootstrap_target(&state, model_id).await {
|
||||
Ok(Some(target)) => target,
|
||||
Ok(None) => {
|
||||
@@ -765,17 +776,6 @@ pub(crate) async fn ensure_relay_mesh_for_record(
|
||||
));
|
||||
}
|
||||
};
|
||||
|
||||
// No serving configuration exists — genuine consumer-only start.
|
||||
// Capture scope BEFORE discovery so a concurrent workspace switch can be
|
||||
// detected under the install lock. Route through
|
||||
// `install_client_under_workspace_transition` which acquires
|
||||
// `workspace_transition`, validates full scope identity
|
||||
// (scope_id, relay, owner, generation), then calls the install closure —
|
||||
// serialized against apply_workspace and live identity import.
|
||||
let captured_scope = state
|
||||
.capture_active_scope()
|
||||
.ok_or("mesh client install: no active workspace scope")?;
|
||||
let model_id_owned = model_id.to_string();
|
||||
let endpoint = target.endpoint_addr;
|
||||
scope_impl::install_client_under_workspace_transition(app, &captured_scope, move || {
|
||||
|
||||
@@ -6,6 +6,7 @@
|
||||
|
||||
use std::future::Future;
|
||||
|
||||
use futures_util::future::BoxFuture;
|
||||
use tauri::{AppHandle, Manager, State};
|
||||
|
||||
use crate::app_state::AppState;
|
||||
@@ -104,26 +105,19 @@ pub(crate) async fn fail_if_client_mesh_active<R: tauri::Runtime>(
|
||||
/// Used by `apply_workspace`, live identity import, and tests 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.
|
||||
/// `transition_body` is an async closure (returns a `BoxFuture`) so that
|
||||
/// production callers that dispatch `tokio::task::spawn_blocking` and await the
|
||||
/// result keep the guard alive across the dispatch. Tests may pass a simple
|
||||
/// sync-compatible closure via `|| Box::pin(async { Ok("done") })`.
|
||||
///
|
||||
/// Production callers that must `.await` after dispatch use this helper with
|
||||
/// a body that dispatches `tokio::task::spawn_blocking` and immediately returns
|
||||
/// the join handle's `.await` — the guard is held across the dispatch and the
|
||||
/// `.await` is done at the call site after this helper returns.
|
||||
///
|
||||
/// Alternatively, callers that cannot fit their body into a sync `FnOnce` call
|
||||
/// `run_mesh_transition_preflight` (which they invoke after acquiring the lock
|
||||
/// themselves) to share the preflight logic without lifetime constraints.
|
||||
#[cfg_attr(not(test), allow(dead_code))]
|
||||
/// 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>,
|
||||
F: FnOnce() -> BoxFuture<'static, Result<T, String>>,
|
||||
{
|
||||
let state = app.state::<AppState>();
|
||||
let _transition_guard = state.workspace_transition.lock().await;
|
||||
@@ -132,7 +126,7 @@ where
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
fail_if_client_mesh_active(app).await?;
|
||||
|
||||
transition_body()
|
||||
transition_body().await
|
||||
}
|
||||
|
||||
/// Run only the Mesh-preflight portion of the workspace transition check.
|
||||
|
||||
@@ -70,7 +70,7 @@ async fn test_active_client_stop_then_transition_preflight_succeeds() {
|
||||
// Step 3 — production transition preflight must succeed (no client active).
|
||||
// The no-op body proves the lock was acquired and the preflight passed.
|
||||
let result = super::scope_impl::with_workspace_transition_preflight(&app_handle, || {
|
||||
Ok::<&str, String>("body ran")
|
||||
Box::pin(async { Ok::<&str, String>("body ran") })
|
||||
})
|
||||
.await;
|
||||
|
||||
@@ -253,7 +253,7 @@ async fn test_install_held_transition_preflight_observes_client() {
|
||||
let app_handle_clone = app_handle.clone();
|
||||
let transition_task = tokio::task::spawn(async move {
|
||||
super::scope_impl::with_workspace_transition_preflight(&app_handle_clone, || {
|
||||
Ok::<&str, String>("body ran")
|
||||
Box::pin(async { Ok::<&str, String>("body ran") })
|
||||
})
|
||||
.await
|
||||
});
|
||||
|
||||
@@ -179,6 +179,51 @@ fn minimal_agent_snapshot_json(name: &str) -> Vec<u8> {
|
||||
encode_snapshot_json(&snap).expect("encode_snapshot_json must succeed for minimal snapshot")
|
||||
}
|
||||
|
||||
/// Build a minimal agent snapshot JSON that includes one core memory entry.
|
||||
///
|
||||
/// Used by tests that must exercise the Phase-4 memory-publish path and assert
|
||||
/// that `MemoryPublish` carries the captured (pre-switch) relay and owner.
|
||||
fn minimal_agent_snapshot_json_with_memory(name: &str) -> Vec<u8> {
|
||||
use crate::managed_agents::agent_snapshot::encode_snapshot_json;
|
||||
use crate::managed_agents::agent_snapshot::{
|
||||
AgentSnapshot, AgentSnapshotDefinition, AgentSnapshotMemory, AgentSnapshotMemoryEntry,
|
||||
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::Core,
|
||||
entries: vec![AgentSnapshotMemoryEntry {
|
||||
slug: buzz_core_pkg::engram::CORE_SLUG.to_string(),
|
||||
body: format!("# {name}\nTest memory body."),
|
||||
}],
|
||||
},
|
||||
};
|
||||
encode_snapshot_json(&snap).expect("encode_snapshot_json must succeed for memory snapshot")
|
||||
}
|
||||
|
||||
/// Set up a mock Tauri `App` with `AppState` managed, an active workspace
|
||||
/// scope, and matching owner keys.
|
||||
///
|
||||
@@ -242,8 +287,10 @@ fn setup_import_app_with_scope(
|
||||
/// 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.
|
||||
/// The snapshot carries one core memory entry so `submit_memory` actually
|
||||
/// fires. The memory adapter asserts BOTH the captured relay URL AND that the
|
||||
/// built engram event's `p` tag (owner counterpart/coordinate) matches the
|
||||
/// CAPTURED owner key — not the new owner committed in `after_store`.
|
||||
#[tokio::test]
|
||||
// SAFETY: `#[tokio::test]` uses a single-threaded runtime by default, so
|
||||
// holding `std::sync::Mutex` across `.await` points cannot deadlock.
|
||||
@@ -266,21 +313,29 @@ async fn test_agent_switch_between_store_and_profile_finishes_captured_outbound(
|
||||
use tauri::Manager;
|
||||
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let (app, _owner_keys) = setup_import_app_with_scope(&tmp);
|
||||
let (app, owner_keys) = setup_import_app_with_scope(&tmp);
|
||||
let handle = app.handle();
|
||||
|
||||
let file_bytes = minimal_agent_snapshot_json("TestAgent");
|
||||
// Snapshot with one core memory entry so Phase 4 actually calls submit_memory.
|
||||
let file_bytes = minimal_agent_snapshot_json_with_memory("TestAgent");
|
||||
let input = AgentSnapshotImportConfirm {
|
||||
file_bytes,
|
||||
keep_allowlist: false,
|
||||
};
|
||||
|
||||
// Relay URL embedded in the captured scope (must appear in profile_sync).
|
||||
// Relay URL embedded in the captured scope (must appear in profile_sync and
|
||||
// in the relay URL passed to submit_memory).
|
||||
let expected_relay = "wss://captured.example".to_string();
|
||||
// Captured owner pubkey — must appear in the `p` tag of the built engram event.
|
||||
let captured_owner_pubkey_hex = owner_keys.public_key().to_hex();
|
||||
|
||||
// Track which relay URLs the outbound adapters received.
|
||||
// Track what the outbound adapters received.
|
||||
let profile_relay = Arc::new(Mutex::new(None::<String>));
|
||||
let memory_relay = Arc::new(Mutex::new(None::<String>));
|
||||
let memory_owner_p_tag = Arc::new(Mutex::new(None::<String>));
|
||||
let pr = profile_relay.clone();
|
||||
let mr = memory_relay.clone();
|
||||
let mop = memory_owner_p_tag.clone();
|
||||
|
||||
// The after_store hook needs to commit a new scope via AppState.
|
||||
// Get the state reference from the app's managed state.
|
||||
@@ -316,7 +371,20 @@ async fn test_agent_switch_between_store_and_profile_finishes_captured_outbound(
|
||||
Ok(())
|
||||
})
|
||||
},
|
||||
|_m: MemoryPublish<'_>| Box::pin(async { panic!("no memory entries in this snapshot") }),
|
||||
move |m: MemoryPublish<'_>| {
|
||||
// Assert: relay URL contains the captured base relay, not the switched one.
|
||||
let relay = m.relay_url.to_string();
|
||||
*mr.lock().unwrap() = Some(relay.clone());
|
||||
|
||||
// Extract the `p` tag from the built engram event JSON.
|
||||
// The `p` tag must carry the CAPTURED owner's pubkey hex — not the
|
||||
// post-switch owner committed in after_store.
|
||||
let event_bytes = m.event_json.to_vec();
|
||||
let p_tag_hex = extract_p_tag_from_event_json(&event_bytes);
|
||||
*mop.lock().unwrap() = p_tag_hex;
|
||||
|
||||
Box::pin(async move { Ok(()) })
|
||||
},
|
||||
)
|
||||
.await;
|
||||
|
||||
@@ -330,7 +398,44 @@ async fn test_agent_switch_between_store_and_profile_finishes_captured_outbound(
|
||||
Some(expected_relay.as_str()),
|
||||
"profile adapter must receive captured relay, got: {profile_seen:?}"
|
||||
);
|
||||
// No memory entries in this snapshot — memory adapter not called.
|
||||
|
||||
// Memory adapter received the OLD captured relay URL in its relay field.
|
||||
let memory_relay_seen = memory_relay.lock().unwrap().clone();
|
||||
assert!(
|
||||
memory_relay_seen
|
||||
.as_deref()
|
||||
.map_or(false, |r| r.contains("captured.example")),
|
||||
"memory adapter relay_url must contain captured relay 'captured.example', \
|
||||
got: {memory_relay_seen:?}"
|
||||
);
|
||||
|
||||
// Memory event's `p` tag must equal the CAPTURED owner's pubkey hex.
|
||||
let p_tag_seen = memory_owner_p_tag.lock().unwrap().clone();
|
||||
assert_eq!(
|
||||
p_tag_seen.as_deref(),
|
||||
Some(captured_owner_pubkey_hex.as_str()),
|
||||
"engram event p-tag (owner counterpart/coordinate) must equal the captured \
|
||||
owner's pubkey, not the post-switch owner; got: {p_tag_seen:?}"
|
||||
);
|
||||
}
|
||||
|
||||
/// Extract the first `p` tag value from a nostr event JSON byte slice.
|
||||
///
|
||||
/// Returns `Some(hex_pubkey)` if a `["p", "<hex>"]` tag entry is found,
|
||||
/// `None` if the JSON cannot be parsed or has no `p` tag.
|
||||
fn extract_p_tag_from_event_json(event_json: &[u8]) -> Option<String> {
|
||||
let val: serde_json::Value = serde_json::from_slice(event_json).ok()?;
|
||||
let tags = val.get("tags")?.as_array()?;
|
||||
for tag in tags {
|
||||
if let Some(arr) = tag.as_array() {
|
||||
if arr.first().and_then(|v| v.as_str()) == Some("p") {
|
||||
if let Some(hex) = arr.get(1).and_then(|v| v.as_str()) {
|
||||
return Some(hex.to_string());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
/// `before_store` hook advances scope generation — Phase 3a must reject with
|
||||
|
||||
@@ -112,22 +112,47 @@ pub async fn apply_workspace(
|
||||
agent_managed_profiles: Option<bool>,
|
||||
app: AppHandle,
|
||||
) -> Result<crate::managed_agents::scope::WorkspaceApplyResult, String> {
|
||||
use crate::managed_agents::scope::WorkspaceApplyResult;
|
||||
|
||||
// ── Layer 1: async serialization lock + Mesh preflight ──────────────────
|
||||
// workspace_transition serializes apply_workspace and live identity import
|
||||
// so scope transitions are never concurrent. We acquire via a clone so the
|
||||
// borrow does not prevent moving `app` into spawn_blocking below.
|
||||
let lock_app = app.clone();
|
||||
let lock_state = lock_app.state::<AppState>();
|
||||
let _transition_guard = lock_state.workspace_transition.lock().await;
|
||||
|
||||
// Fail closed if a client-mode Mesh runtime is active — no client-mode
|
||||
// runtime may be active when a workspace switch begins (Option A ruling).
|
||||
// Called via `run_mesh_transition_preflight` so the shared preflight logic
|
||||
// is not duplicated inline; the guard above remains held throughout.
|
||||
// so scope transitions are never concurrent.
|
||||
//
|
||||
// When the `mesh-llm` feature is active, `with_workspace_transition_preflight`
|
||||
// acquires the lock AND runs `fail_if_client_mesh_active` under a single guard,
|
||||
// with the guard held across the entire async body. When the feature is off,
|
||||
// we acquire the lock inline (no preflight needed).
|
||||
#[cfg(feature = "mesh-llm")]
|
||||
crate::commands::mesh_llm::scope_impl::run_mesh_transition_preflight(&app).await?;
|
||||
{
|
||||
let app_for_preflight = app.clone();
|
||||
return crate::commands::mesh_llm::scope_impl::with_workspace_transition_preflight(
|
||||
&app_for_preflight,
|
||||
move || {
|
||||
Box::pin(apply_workspace_body(
|
||||
relay_url,
|
||||
nsec,
|
||||
repos_dir,
|
||||
agent_managed_profiles,
|
||||
app,
|
||||
))
|
||||
},
|
||||
)
|
||||
.await;
|
||||
}
|
||||
#[cfg(not(feature = "mesh-llm"))]
|
||||
{
|
||||
let lock_app = app.clone();
|
||||
let lock_state = lock_app.state::<AppState>();
|
||||
let _transition_guard = lock_state.workspace_transition.lock().await;
|
||||
apply_workspace_body(relay_url, nsec, repos_dir, agent_managed_profiles, app).await
|
||||
}
|
||||
}
|
||||
async fn apply_workspace_body(
|
||||
relay_url: String,
|
||||
nsec: Option<String>,
|
||||
repos_dir: Option<String>,
|
||||
agent_managed_profiles: Option<bool>,
|
||||
app: AppHandle,
|
||||
) -> Result<crate::managed_agents::scope::WorkspaceApplyResult, String> {
|
||||
use crate::managed_agents::scope::WorkspaceApplyResult;
|
||||
|
||||
let restore_app = app.clone();
|
||||
let blocking_result: Result<WorkspaceApplyResult, String> =
|
||||
@@ -366,7 +391,7 @@ pub async fn apply_workspace(
|
||||
return Ok(apply_result);
|
||||
}
|
||||
|
||||
// ── Post-commit (non-rollback) ────────────────────────────────────────────
|
||||
// ── Post-commit (non-rollback) ──────────────────────────────────────
|
||||
// The workspace HAS switched. Post-commit failures surface as degradation
|
||||
// on the applied result — we never pretend the old scope survived.
|
||||
let mut degraded: Vec<String> = Vec::new();
|
||||
|
||||
Reference in New Issue
Block a user