mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
feat(desktop): publish persona create/edit/delete to the relay
Persona create, edit, and delete now reach the relay so personas persist across reboots and reflect in other Buzz clients. #939 built the local retention primitives but left them inert; this wires every persona writer to a single publish path. pending_sync is the seam: each writer (UI create/edit, delete tombstone, launch reconcile) retains a signed event locally with pending_sync = 1 while holding the managed-agents store lock, and one background flush loop is the sole relay publisher. Delete purges the persona's retention row inside the lock-held command body (closing the same-second resurrect race) and enqueues a NIP-09 a-tag-only kind:5 tombstone for the flush loop to publish out of band. mark_synced is compare-and-clear on created_at+content so an edit landing mid-flush is never falsely cleared. Ownership is intrinsic to the signing key, so the only skip anywhere is is_builtin. 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
a47bb15254
commit
1747d2ec56
@@ -29,6 +29,105 @@ fn trim_optional(value: Option<String>) -> Option<String> {
|
||||
})
|
||||
}
|
||||
|
||||
/// 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;
|
||||
|
||||
@@ -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::<AppState>();
|
||||
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![
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -63,6 +63,18 @@ pub fn build_persona_event(record: &PersonaRecord) -> Result<EventBuilder, Strin
|
||||
Ok(EventBuilder::new(Kind::Custom(KIND_PERSONA as u16), content_json).tags(tags))
|
||||
}
|
||||
|
||||
/// Build a NIP-09 deletion (kind:5) targeting a persona's kind:30175 event.
|
||||
///
|
||||
/// Carries a single `a`-tag with the NIP-33 coordinate `30175:<owner>:<d_tag>`
|
||||
/// 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<EventBuilder, String> {
|
||||
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<PersonaRecord, String>
|
||||
})
|
||||
}
|
||||
|
||||
/// 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<String, String> {
|
||||
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<Vec<nostr::Event>, String> {
|
||||
let pubkey = {
|
||||
let keys = state.keys.lock().map_err(|e| e.to_string())?;
|
||||
keys.public_key().to_hex()
|
||||
) -> Result<u32, String> {
|
||||
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 {
|
||||
|
||||
@@ -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<Connection, String> {
|
||||
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<Vec<RetainedEvent>, 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();
|
||||
|
||||
@@ -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<u32, String> {
|
||||
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<u32, S
|
||||
|
||||
let pubkey = keys.public_key().to_hex();
|
||||
|
||||
// Read personas.json fresh at migration time. Nothing to migrate if the
|
||||
// file is absent.
|
||||
// Read personas.json fresh at reconcile time. Nothing to do if absent.
|
||||
let personas_path = base_dir.join("personas.json");
|
||||
if !personas_path.exists() {
|
||||
return Ok(0);
|
||||
@@ -991,12 +993,6 @@ fn migrate_personas_in_dir(base_dir: &Path, keys: &nostr::Keys) -> Result<u32, S
|
||||
let conn =
|
||||
open_retention_db(&db_path).map_err(|e| format!("failed to open retention db: {e}"))?;
|
||||
|
||||
// Idempotency: the retention rows themselves are the sentinel. If the
|
||||
// owner already has retained personas, migration ran on a prior launch.
|
||||
if has_retained_personas(&conn, &pubkey)? {
|
||||
return Ok(0);
|
||||
}
|
||||
|
||||
let mut migrated = 0u32;
|
||||
|
||||
for record in &records {
|
||||
@@ -1014,11 +1010,20 @@ fn migrate_personas_in_dir(base_dir: &Path, keys: &nostr::Keys) -> Result<u32, S
|
||||
.sign_with_keys(keys)
|
||||
.map_err(|e| format!("failed to sign event for '{}': {e}", record.display_name))?;
|
||||
|
||||
// Per-coordinate reconcile: skip when an identical body is already
|
||||
// retained, so an unchanged persona doesn't reset `pending_sync`.
|
||||
let event_content = event.content.to_string();
|
||||
if let Some(existing) = get_retained_event(&conn, KIND_PERSONA, &pubkey, &d_tag)? {
|
||||
if existing.content == event_content {
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
let retained = RetainedEvent {
|
||||
kind: KIND_PERSONA,
|
||||
pubkey: pubkey.clone(),
|
||||
d_tag,
|
||||
content: event.content.to_string(),
|
||||
content: event_content,
|
||||
// Safety: nostr timestamps are seconds and stay below i64::MAX
|
||||
// until year 2262.
|
||||
created_at: event.created_at.as_secs() as i64,
|
||||
|
||||
@@ -793,18 +793,99 @@ fn migrate_personas_skips_builtins() {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn migrate_personas_skips_when_retention_already_populated() {
|
||||
fn migrate_personas_unchanged_second_run_is_noop() {
|
||||
let base = tempfile::tempdir().unwrap();
|
||||
write_base_personas(base.path(), &one_persona());
|
||||
let keys = nostr::Keys::generate();
|
||||
|
||||
// First run migrates; second run is a no-op (retention rows are the
|
||||
// sentinel — no separate sentinel file).
|
||||
// First run retains; second run with identical personas re-retains
|
||||
// nothing — the per-coordinate content matches, so `pending_sync` is
|
||||
// not churned.
|
||||
assert_eq!(migrate_personas_in_dir(base.path(), &keys).unwrap(), 1);
|
||||
assert_eq!(migrate_personas_in_dir(base.path(), &keys).unwrap(), 0);
|
||||
assert!(!base.path().join("migration_state.json").exists());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn migrate_personas_new_persona_after_first_run_gets_retained() {
|
||||
use crate::managed_agents::retention::{get_retained_personas, open_retention_db};
|
||||
|
||||
let base = tempfile::tempdir().unwrap();
|
||||
write_base_personas(base.path(), &one_persona());
|
||||
let keys = nostr::Keys::generate();
|
||||
let pubkey = keys.public_key().to_hex();
|
||||
|
||||
assert_eq!(migrate_personas_in_dir(base.path(), &keys).unwrap(), 1);
|
||||
|
||||
// A persona added to personas.json after the first reconcile must be
|
||||
// picked up — the whole-store sentinel that previously short-circuited
|
||||
// this is gone.
|
||||
let mut two = one_persona();
|
||||
two.as_array_mut().unwrap().push(serde_json::json!({
|
||||
"id": "test-writer",
|
||||
"display_name": "Test Writer",
|
||||
"system_prompt": "You write tests.",
|
||||
"is_builtin": false,
|
||||
"is_active": true,
|
||||
"name_pool": [],
|
||||
"env_vars": {},
|
||||
"created_at": "2025-01-02T00:00:00Z",
|
||||
"updated_at": "2025-01-02T00:00:00Z"
|
||||
}));
|
||||
write_base_personas(base.path(), &two);
|
||||
|
||||
assert_eq!(migrate_personas_in_dir(base.path(), &keys).unwrap(), 1);
|
||||
|
||||
let conn = open_retention_db(&base.path().join("retention.db")).unwrap();
|
||||
let rows = get_retained_personas(&conn, &pubkey).unwrap();
|
||||
assert_eq!(rows.len(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn migrate_personas_edited_persona_re_retains_pending() {
|
||||
use crate::managed_agents::retention::{get_retained_event, mark_synced, open_retention_db};
|
||||
use buzz_core_pkg::kind::KIND_PERSONA;
|
||||
|
||||
let base = tempfile::tempdir().unwrap();
|
||||
write_base_personas(base.path(), &one_persona());
|
||||
let keys = nostr::Keys::generate();
|
||||
let pubkey = keys.public_key().to_hex();
|
||||
|
||||
assert_eq!(migrate_personas_in_dir(base.path(), &keys).unwrap(), 1);
|
||||
|
||||
// Simulate the flush loop confirming the first publish.
|
||||
let conn = open_retention_db(&base.path().join("retention.db")).unwrap();
|
||||
let row = get_retained_event(&conn, KIND_PERSONA, &pubkey, "code-reviewer")
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
mark_synced(
|
||||
&conn,
|
||||
KIND_PERSONA,
|
||||
&pubkey,
|
||||
"code-reviewer",
|
||||
row.created_at,
|
||||
&row.content,
|
||||
)
|
||||
.unwrap();
|
||||
drop(conn);
|
||||
|
||||
// Editing the persona on disk must re-retain it as pending so the edit
|
||||
// reaches the relay on the next flush.
|
||||
let mut edited = one_persona();
|
||||
edited.as_array_mut().unwrap()[0]["system_prompt"] =
|
||||
serde_json::json!("You review code carefully.");
|
||||
write_base_personas(base.path(), &edited);
|
||||
|
||||
assert_eq!(migrate_personas_in_dir(base.path(), &keys).unwrap(), 1);
|
||||
|
||||
let conn = open_retention_db(&base.path().join("retention.db")).unwrap();
|
||||
let row = get_retained_event(&conn, KIND_PERSONA, &pubkey, "code-reviewer")
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert!(row.pending_sync);
|
||||
assert!(row.content.contains("carefully"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn migrate_personas_no_file_is_noop() {
|
||||
let base = tempfile::tempdir().unwrap();
|
||||
|
||||
@@ -503,6 +503,47 @@ pub async fn submit_event(
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
/// POST an already-signed event to `/events` with NIP-98 auth.
|
||||
///
|
||||
/// The persona flush loop drains pre-signed events from the retention store,
|
||||
/// so it must publish them verbatim — re-signing through `submit_event` would
|
||||
/// mint a new `created_at`/signature and break the compare-and-clear that
|
||||
/// `mark_synced` relies on. Only the NIP-98 request auth is signed here (with
|
||||
/// the owner keys), and that lock is dropped before the `.await`.
|
||||
pub async fn submit_signed_event(
|
||||
event: &nostr::Event,
|
||||
state: &AppState,
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
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
|
||||
|
||||
Reference in New Issue
Block a user