diff --git a/desktop/src-tauri/src/commands/personas.rs b/desktop/src-tauri/src/commands/personas.rs index 56f6742f4..7dee34839 100644 --- a/desktop/src-tauri/src/commands/personas.rs +++ b/desktop/src-tauri/src/commands/personas.rs @@ -29,6 +29,105 @@ fn trim_optional(value: Option) -> Option { }) } +/// Retain a freshly authored persona event in the local store, flagged for +/// relay sync. Called inside a command's `managed_agents_store_lock`-held body +/// after `save_personas`; the background flush loop publishes it out-of-band. +/// +/// The event is signed with the owner keys at call time, so its `created_at` +/// is `now` — newer than any prior retained row, clearing the upsert's +/// newer-or-equal guard. `pending_sync = 1` enqueues it for the flush loop, +/// which is the sole publisher. Best-effort: a failure here is logged and +/// swallowed so a retention hiccup never blocks the disk-authoritative write. +fn retain_persona_pending(app: &AppHandle, state: &AppState, persona: &PersonaRecord) { + use crate::managed_agents::{ + managed_agents_base_dir, + persona_events::{build_persona_event, persona_d_tag}, + retention::{open_retention_db, retain_event, RetainedEvent}, + }; + use buzz_core_pkg::kind::KIND_PERSONA; + use nostr::JsonUtil; + + let result = (|| -> Result<(), String> { + let (pubkey, event) = { + let keys = state.keys.lock().map_err(|e| e.to_string())?; + let event = build_persona_event(persona)? + .sign_with_keys(&keys) + .map_err(|e| format!("failed to sign persona event: {e}"))?; + (keys.public_key().to_hex(), event) + }; + let conn = open_retention_db(&managed_agents_base_dir(app)?.join("retention.db"))?; + retain_event( + &conn, + &RetainedEvent { + kind: KIND_PERSONA, + pubkey, + d_tag: persona_d_tag(persona), + 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: persona-retain: {e}"); + } +} + +/// Purge a deleted persona's pending row and enqueue a NIP-09 tombstone, both +/// inside the `managed_agents_store_lock`-held delete body. +/// +/// PURGE IN: `delete_retained_event` removes the persona's `(30175, pubkey, +/// d_tag)` row. Running it under the same lock that serializes `retain_event` +/// closes the same-second resurrect race — a concurrent edit can't re-insert a +/// pending persona row after the tombstone is queued. +/// +/// PUBLISH OUT: the kind:5 tombstone is retained at its own coordinate `(5, +/// pubkey, d_tag)` (distinct from the purged persona row) with `pending_sync = +/// 1`; the flush loop publishes it. Best-effort: a failure is logged and +/// swallowed so a retention hiccup never blocks the disk-authoritative delete. +fn tombstone_persona_pending(app: &AppHandle, state: &AppState, d_tag: &str) { + use crate::managed_agents::{ + managed_agents_base_dir, + persona_events::build_persona_delete, + retention::{delete_retained_event, open_retention_db, retain_event, RetainedEvent}, + }; + use buzz_core_pkg::kind::KIND_PERSONA; + use nostr::JsonUtil; + + const KIND_DELETE: u32 = 5; + + let result = (|| -> Result<(), String> { + let (pubkey, event) = { + let keys = state.keys.lock().map_err(|e| e.to_string())?; + let pubkey = keys.public_key().to_hex(); + let event = build_persona_delete(d_tag, &pubkey)? + .sign_with_keys(&keys) + .map_err(|e| format!("failed to sign persona tombstone: {e}"))?; + (pubkey, event) + }; + let conn = open_retention_db(&managed_agents_base_dir(app)?.join("retention.db"))?; + // Purge the persona row first so an unpublished edit can never resurrect + // it after the tombstone publishes. + delete_retained_event(&conn, KIND_PERSONA, &pubkey, d_tag)?; + retain_event( + &conn, + &RetainedEvent { + kind: KIND_DELETE, + pubkey, + d_tag: d_tag.to_string(), + 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: persona-tombstone: {e}"); + } +} + #[tauri::command] pub fn list_personas( app: AppHandle, @@ -86,6 +185,7 @@ pub fn create_persona( }; personas.push(persona.clone()); save_personas(&app, &personas)?; + retain_persona_pending(&app, &state, &persona); try_regenerate_nest(&app); Ok(persona) } @@ -139,6 +239,7 @@ pub fn update_persona( .into_iter() .find(|record| record.id == input.id) .ok_or_else(|| format!("persona {} disappeared unexpectedly", input.id))?; + retain_persona_pending(&app, &state, &result); try_regenerate_nest(&app); Ok(result) } @@ -164,6 +265,10 @@ pub fn delete_persona( .any(|persona_id| persona_id == id.as_str()) }); validate_persona_deletion(persona, referenced_by_team)?; + // Capture the coordinate before the record leaves the list. Only reached + // for non-builtin, non-team personas (validate_persona_deletion rejects + // both), so every deleted persona here is one this owner published. + let d_tag = crate::managed_agents::persona_events::persona_d_tag(persona); let original_len = personas.len(); personas.retain(|record| record.id != id); @@ -171,6 +276,7 @@ pub fn delete_persona( return Err(format!("persona {id} not found")); } save_personas(&app, &personas)?; + tombstone_persona_pending(&app, &state, &d_tag); let mut agents = load_managed_agents(&app)?; let mut changed_agents = false; diff --git a/desktop/src-tauri/src/lib.rs b/desktop/src-tauri/src/lib.rs index bc0c997ce..9a0a5e383 100644 --- a/desktop/src-tauri/src/lib.rs +++ b/desktop/src-tauri/src/lib.rs @@ -695,6 +695,33 @@ pub fn run() { } }); + // Drain persona events the retention store flagged `pending_sync` + // (UI create/edit, delete tombstones, launch reconcile) to the + // relay. One loop is the sole publisher for all four writers; a + // relay-unreachable tick leaves rows pending for the next sweep. + let flush_handle = app.handle().clone(); + tauri::async_runtime::spawn(async move { + use std::time::Duration; + use tauri::Manager; + let Ok(db_path) = managed_agents::managed_agents_base_dir(&flush_handle) + .map(|d| d.join("retention.db")) + else { + eprintln!("buzz-desktop: persona-flush: cannot resolve retention db path"); + return; + }; + loop { + let state = flush_handle.state::(); + if let Err(e) = managed_agents::persona_events::flush_pending_persona_events( + &db_path, &state, + ) + .await + { + eprintln!("buzz-desktop: persona-flush: {e}"); + } + tokio::time::sleep(Duration::from_secs(30)).await; + } + }); + Ok(()) }) .invoke_handler(tauri::generate_handler![ diff --git a/desktop/src-tauri/src/managed_agents/mod.rs b/desktop/src-tauri/src/managed_agents/mod.rs index 084991f6b..1d6c0fd28 100644 --- a/desktop/src-tauri/src/managed_agents/mod.rs +++ b/desktop/src-tauri/src/managed_agents/mod.rs @@ -4,9 +4,6 @@ mod env_vars; mod nest; mod persona_avatars; mod persona_card; -// `publish_persona_event` / `fetch_persona_events` are #939 publishing -// primitives not yet wired to a call site; keep them without a dead-code warn. -#[allow(dead_code)] pub(crate) mod persona_events; mod personas; #[cfg(windows)] @@ -14,7 +11,6 @@ mod process_lifecycle; #[cfg(feature = "mesh-llm")] mod relay_mesh; mod restore; -#[allow(dead_code)] pub mod retention; mod runtime; mod storage; diff --git a/desktop/src-tauri/src/managed_agents/persona_events.rs b/desktop/src-tauri/src/managed_agents/persona_events.rs index bc9292101..0afd9c5a2 100644 --- a/desktop/src-tauri/src/managed_agents/persona_events.rs +++ b/desktop/src-tauri/src/managed_agents/persona_events.rs @@ -63,6 +63,18 @@ pub fn build_persona_event(record: &PersonaRecord) -> Result:` +/// and no `e`-tag: an `e`-tag routes the relay to the event-id deletion path, +/// which leaves the parameterized-replaceable coordinate live. The coordinate +/// delete removes the persona for every client and across reboots. +pub fn build_persona_delete(d_tag: &str, owner_pubkey_hex: &str) -> Result { + let coord = format!("{KIND_PERSONA}:{owner_pubkey_hex}:{d_tag}"); + let tag = Tag::parse(["a", coord.as_str()]).map_err(|e| format!("invalid a-tag: {e}"))?; + Ok(EventBuilder::new(Kind::Custom(5), "").tags(vec![tag])) +} + /// Parse a kind:30175 event back into a `PersonaRecord`. /// /// The event's d-tag becomes the persona ID and slug. @@ -104,29 +116,77 @@ pub fn persona_from_event(event: &nostr::Event) -> Result }) } -/// Publish a persona event to the relay. -pub async fn publish_persona_event( - record: &PersonaRecord, +/// Drain every `pending_sync` event from the retention store to the relay. +/// +/// Each writer (UI create/edit, delete tombstone, launch reconcile) retains a +/// signed event with `pending_sync = 1`; this loop is the sole publisher. +/// +/// Per row, the last synchronous read before the network `.await` is a fresh +/// `get_retained_event` re-check — the connection holds no `Mutex` across the +/// await, so a concurrent edit or delete is observed here: +/// - gone (deleted): skip, nothing to publish. +/// - newer `created_at` or different `content`: skip; the newer row is itself +/// `pending_sync` and publishes on its own pass. +/// +/// Only a row that still matches what we read is published, then cleared via +/// `mark_synced` on the exact `created_at`+`content` the relay accepted — so an +/// edit landing between publish and clear is never falsely marked synced. +/// +/// Returns the number of events the relay accepted. Best-effort: a relay +/// failure on one row leaves it pending for the next sweep and does not abort +/// the remaining rows. +pub async fn flush_pending_persona_events( + db_path: &std::path::Path, state: &AppState, -) -> Result { - let builder = build_persona_event(record)?; - let response = crate::relay::submit_event(builder, state).await?; - Ok(response.event_id) -} - -/// Fetch all persona events authored by the current user from the relay. -pub async fn fetch_persona_events(state: &AppState) -> Result, String> { - let pubkey = { - let keys = state.keys.lock().map_err(|e| e.to_string())?; - keys.public_key().to_hex() +) -> Result { + use crate::managed_agents::retention::{ + get_pending_sync, get_retained_event, mark_synced, open_retention_db, }; + use nostr::JsonUtil; - let filter = serde_json::json!({ - "kinds": [KIND_PERSONA], - "authors": [pubkey] - }); + let pending = { + let conn = open_retention_db(db_path)?; + get_pending_sync(&conn)? + }; // connection dropped before any .await - crate::relay::query_relay(state, &[filter]).await + let mut flushed = 0u32; + for row in pending { + // Re-read immediately before publishing; the row may have been edited + // or deleted since the pending snapshot above. + let current = { + let conn = open_retention_db(db_path)?; + get_retained_event(&conn, row.kind, &row.pubkey, &row.d_tag)? + }; + let Some(current) = current else { + continue; // deleted out from under us + }; + if current.created_at != row.created_at || current.content != row.content { + continue; // superseded by a newer edit; that row publishes itself + } + + let event = nostr::Event::from_json(¤t.raw_event) + .map_err(|e| format!("failed to parse retained event '{}': {e}", current.d_tag))?; + + if crate::relay::submit_signed_event(&event, state) + .await + .is_err() + { + continue; // relay unreachable — stays pending for the next sweep + } + + let conn = open_retention_db(db_path)?; + mark_synced( + &conn, + current.kind, + ¤t.pubkey, + ¤t.d_tag, + current.created_at, + ¤t.content, + )?; + flushed += 1; + } + + Ok(flushed) } /// SHA-256 (lowercase hex) of a persona's canonical content JSON. @@ -311,6 +371,32 @@ mod tests { assert!(restored.is_active); } + #[test] + fn build_persona_delete_has_single_a_tag_no_e_tag() { + const OWNER: &str = "79be667ef9dcbbac55a06295ce870b07029bfcdb2dce28d959f2815b16f81798"; + let builder = build_persona_delete("test-slug", OWNER).unwrap(); + let keys = nostr::Keys::generate(); + let event = builder.sign_with_keys(&keys).unwrap(); + + assert_eq!(event.kind, Kind::Custom(5)); + + let a_tags: Vec<&[String]> = event + .tags + .iter() + .map(|t| t.as_slice()) + .filter(|v| v.first().map(String::as_str) == Some("a")) + .collect(); + assert_eq!(a_tags.len(), 1); + assert_eq!(a_tags[0][1], format!("{KIND_PERSONA}:{OWNER}:test-slug")); + + // An e-tag would route to the event-id deletion path and leave the + // replaceable coordinate live — the tombstone must carry none. + assert!(event + .tags + .iter() + .all(|t| t.as_slice().first().map(String::as_str) != Some("e"))); + } + #[test] fn persona_content_hash_is_deterministic() { let content = PersonaEventContent { diff --git a/desktop/src-tauri/src/managed_agents/retention.rs b/desktop/src-tauri/src/managed_agents/retention.rs index 051c4d917..85a633358 100644 --- a/desktop/src-tauri/src/managed_agents/retention.rs +++ b/desktop/src-tauri/src/managed_agents/retention.rs @@ -1,8 +1,9 @@ //! Local SQLite retention store for persona events. //! //! Provides durable client-side storage for persona events, enabling offline -//! boot when the relay is unreachable. Uses `INSERT OR REPLACE` keyed on -//! `(kind, pubkey, d_tag)` for NIP-33 latest-wins semantics. +//! boot when the relay is unreachable. Upserts via `ON CONFLICT DO UPDATE` +//! keyed on `(kind, pubkey, d_tag)`, replacing only on a newer-or-equal +//! `created_at` for NIP-33 latest-wins semantics. use std::path::Path; @@ -21,9 +22,18 @@ pub struct RetainedEvent { } /// Open (or create) the retention database at the given path. +/// +/// Sets WAL journaling and a `busy_timeout` on every connection so the +/// flush-loop connection and command-path connections can write concurrently +/// without spurious `SQLITE_BUSY` errors. pub fn open_retention_db(path: &Path) -> Result { let conn = Connection::open(path).map_err(|e| format!("failed to open retention db: {e}"))?; + conn.pragma_update(None, "journal_mode", "WAL") + .map_err(|e| format!("failed to set WAL mode: {e}"))?; + conn.pragma_update(None, "busy_timeout", 5000) + .map_err(|e| format!("failed to set busy_timeout: {e}"))?; + conn.execute_batch( "CREATE TABLE IF NOT EXISTS persona_events ( kind INTEGER NOT NULL, @@ -129,14 +139,50 @@ pub fn get_pending_sync(conn: &Connection) -> Result, String> .map_err(|e| format!("failed to read pending sync row: {e}")) } -/// Clear the pending_sync flag for a specific event (after relay confirms). -pub fn mark_synced(conn: &Connection, kind: u32, pubkey: &str, d_tag: &str) -> Result<(), String> { +/// Clear the `pending_sync` flag for an event the relay just confirmed. +/// +/// Compare-and-clear: only clears the row if its `created_at` and `content` +/// still match what was published. A concurrent edit that upserted a newer +/// version at the same coordinate between the flush loop's read and this call +/// leaves `pending_sync` set, so the newer edit publishes on the next pass +/// instead of being silently dropped. +pub fn mark_synced( + conn: &Connection, + kind: u32, + pubkey: &str, + d_tag: &str, + created_at: i64, + content: &str, +) -> Result<(), String> { conn.execute( "UPDATE persona_events SET pending_sync = 0 + WHERE kind = ?1 AND pubkey = ?2 AND d_tag = ?3 + AND created_at = ?4 AND content = ?5", + params![kind, pubkey, d_tag, created_at, content], + ) + .map_err(|e| format!("failed to mark event synced: {e}"))?; + + Ok(()) +} + +/// Delete a retained event by its coordinate. +/// +/// Called from the synchronous, lock-held delete-persona command body so the +/// purge serializes against `retain_event` upserts at the same coordinate — +/// closing the same-second resurrect race where a pending edit would otherwise +/// publish after the deletion tombstone. +pub fn delete_retained_event( + conn: &Connection, + kind: u32, + pubkey: &str, + d_tag: &str, +) -> Result<(), String> { + conn.execute( + "DELETE FROM persona_events WHERE kind = ?1 AND pubkey = ?2 AND d_tag = ?3", params![kind, pubkey, d_tag], ) - .map_err(|e| format!("failed to mark event synced: {e}"))?; + .map_err(|e| format!("failed to delete retained event: {e}"))?; Ok(()) } @@ -263,12 +309,12 @@ mod tests { } #[test] - fn mark_synced_clears_flag() { + fn test_mark_synced_matching_row_clears_flag() { let conn = test_db(); let event = sample_event(); retain_event(&conn, &event).unwrap(); - mark_synced(&conn, 30175, "abc123", "test-persona").unwrap(); + mark_synced(&conn, 30175, "abc123", "test-persona", 1000, &event.content).unwrap(); let pending = get_pending_sync(&conn).unwrap(); assert!(pending.is_empty()); @@ -278,6 +324,53 @@ mod tests { assert!(!results[0].pending_sync); } + #[test] + fn test_mark_synced_stale_version_leaves_flag_set() { + let conn = test_db(); + let published = sample_event(); + retain_event(&conn, &published).unwrap(); + + // A newer edit lands at the same coordinate before the flush loop + // clears the version it published. + let mut newer = sample_event(); + newer.content = r#"{"display_name":"Edited"}"#.to_string(); + newer.created_at = 2000; + retain_event(&conn, &newer).unwrap(); + + // Clearing against the OLD version must not touch the newer pending row. + mark_synced( + &conn, + 30175, + "abc123", + "test-persona", + 1000, + &published.content, + ) + .unwrap(); + + let pending = get_pending_sync(&conn).unwrap(); + assert_eq!(pending.len(), 1); + assert_eq!(pending[0].created_at, 2000); + } + + #[test] + fn test_delete_retained_event_removes_row() { + let conn = test_db(); + retain_event(&conn, &sample_event()).unwrap(); + + delete_retained_event(&conn, 30175, "abc123", "test-persona").unwrap(); + + assert!(get_retained_event(&conn, 30175, "abc123", "test-persona") + .unwrap() + .is_none()); + } + + #[test] + fn test_delete_retained_event_missing_row_is_noop() { + let conn = test_db(); + delete_retained_event(&conn, 30175, "abc123", "nonexistent").unwrap(); + } + #[test] fn has_retained_personas_works() { let conn = test_db(); diff --git a/desktop/src-tauri/src/migration.rs b/desktop/src-tauri/src/migration.rs index 60d2d5a39..daa15216b 100644 --- a/desktop/src-tauri/src/migration.rs +++ b/desktop/src-tauri/src/migration.rs @@ -917,17 +917,20 @@ pub fn migrate_persona_provider_to_runtime(app: &tauri::AppHandle) { rename_provider_to_runtime_in_personas(&path); } -/// Migrate existing `personas.json` entries to persona events in the local -/// retention store. +/// Reconcile `personas.json` into the persona-event retention store. /// /// Must run AFTER `migrate_packs_to_teams` (depends on field renames being /// complete) and AFTER the persisted identity is resolved (it signs every /// retained event with the owner's keys). /// -/// Idempotent: skips when the retention store already holds events for the -/// owner pubkey — the data is the sentinel, so no separate sentinel file is -/// needed. This avoids re-running (and resetting `pending_sync`) on every -/// launch when a sentinel write silently fails. +/// Per-record reconcile: for each non-builtin persona it compares the freshly +/// serialized event content against the retained row at the same coordinate +/// and re-retains (marking `pending_sync = 1`) only when the row is absent or +/// its content differs. An unchanged persona is left untouched, so a launch +/// after a no-op edit does not churn `pending_sync`; a persona added or edited +/// on disk between launches is picked up and republished. There is no +/// whole-store sentinel — comparing per coordinate is what lets newly added +/// personas reach the relay. /// /// Strategy: write to local SQLite retention first (durable copy), mark as /// `pending_sync = 1` for later relay publish. Migration succeeds on local @@ -953,15 +956,15 @@ pub fn migrate_personas_to_events(app: &tauri::AppHandle, keys: &nostr::Keys) { } } -/// Core migration logic, decoupled from the Tauri `AppHandle` for testing. +/// Core reconcile logic, decoupled from the Tauri `AppHandle` for testing. /// -/// Returns the number of personas written to the retention store. Returns -/// `Ok(0)` when migration has already run (retention store has rows for the -/// owner pubkey) or when there are no non-builtin personas to migrate. +/// Returns the number of personas (re)written to the retention store. Returns +/// `Ok(0)` when every non-builtin persona already has a matching retained row +/// (or there are none to reconcile). fn migrate_personas_in_dir(base_dir: &Path, keys: &nostr::Keys) -> Result { use crate::managed_agents::{ persona_events::{build_persona_event, persona_d_tag}, - retention::{has_retained_personas, open_retention_db, retain_event, RetainedEvent}, + retention::{get_retained_event, open_retention_db, retain_event, RetainedEvent}, PersonaRecord, }; use buzz_core_pkg::kind::KIND_PERSONA; @@ -969,8 +972,7 @@ fn migrate_personas_in_dir(base_dir: &Path, keys: &nostr::Keys) -> Result Result Result Result { + let url = format!("{}/events", relay_api_base_url_with_override(state)); + let body_bytes = event.as_json().into_bytes(); + let auth_header = { + let keys = state.keys.lock().map_err(|e| e.to_string())?; + build_nip98_auth_header_for_keys(&keys, &Method::POST, &url, &body_bytes)? + }; // keys lock dropped here + + let response = state + .http_client + .post(&url) + .header("Authorization", auth_header) + .header("Content-Type", "application/json") + .body(body_bytes) + .send() + .await + .map_err(|e| classify_request_error(&e))?; + + if !response.status().is_success() { + return Err(relay_error_message(response).await); + } + + let result: SubmitEventResponse = parse_json_response(response).await?; + + if !result.accepted { + return Err(format!("relay rejected event: {}", result.message)); + } + + Ok(result) +} + /// Sign an event with explicit keys and POST it to `/events` with NIP-98 auth. /// /// Managed-agent flows use this to publish as the agent itself while still