fix(desktop): migrate the legacy global retention db into the active scope

Scoped retention databases (retention/<relay+owner hash>.db) read a
different file than the pre-scoping global retention.db, so an upgrade
abandoned every row the previous release left pending — including signed
kind:5 tombstones and NIP-IA archive requests queued while offline, which
no reconcile can rebuild. The relay kept deleted heads and another device
could resurrect them.

The copy is claimed by exactly one relay scope (the legacy rows were
queued for one relay, so fanning them out to every community would
reintroduce the cross-community leak scoping exists to close) and commits
its completion marker in the same transaction as the rows, so a crash
mid-copy neither double-applies nor drops.

Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
Will Pfleger
2026-07-27 13:05:15 -04:00
parent 9f082be395
commit fbc9d515be
4 changed files with 438 additions and 5 deletions
+37 -5
View File
@@ -10,6 +10,31 @@ use crate::managed_agents::{
};
use crate::relay;
/// Adopt the pre-scoping global retention database's pending rows into `scope`.
///
/// Best-effort: a failure is logged and the boot proceeds. The migration's own
/// crash-safety guards make the next launch retry safely, and blocking the
/// workspace apply on it would be worse than a delayed publish.
fn migrate_legacy_retention_into(
app: &AppHandle,
scope: &crate::managed_agents::retention::RetentionScope,
) {
let Ok(base_dir) = crate::managed_agents::managed_agents_base_dir(app) else {
return;
};
match crate::managed_agents::retention::migrate_legacy_retention_db(
&base_dir,
&scope.db_path,
&scope.owner_keys.public_key().to_hex(),
) {
Ok(0) => {}
Ok(copied) => {
eprintln!("buzz-desktop: adopted {copied} legacy retained event(s) into this community")
}
Err(error) => eprintln!("buzz-desktop: legacy retention migration failed: {error}"),
}
}
#[derive(Deserialize)]
struct RelayInfoIcon {
#[serde(default)]
@@ -191,11 +216,18 @@ pub async fn apply_workspace(
// applied. Running at process boot would target the fallback relay and
// collapse every community into one pending-event store.
match crate::managed_agents::retention::active_retention_scope(&restore_app, &state) {
Ok(scope) => crate::event_sync::spawn_event_sync(
restore_app.clone(),
scope.owner_keys,
scope.db_path,
),
Ok(scope) => {
// Adopt whatever the pre-scoping release left queued in the global
// retention database BEFORE the scoped reconcile and flush run, so
// stranded tombstones and archive requests publish on this boot
// instead of being abandoned by the storage cutover.
migrate_legacy_retention_into(&restore_app, &scope);
crate::event_sync::spawn_event_sync(
restore_app.clone(),
scope.owner_keys,
scope.db_path,
)
}
Err(error) => {
eprintln!("buzz-desktop: scoped event-sync unavailable after workspace apply: {error}");
}
@@ -14,6 +14,9 @@ use tauri::AppHandle;
use crate::app_state::AppState;
mod legacy_migration;
pub use legacy_migration::migrate_legacy_retention_db;
/// Durable event-retention scope for one community relay and owner identity.
///
/// Persona, team, and managed-agent definitions are workspace-global, but
@@ -0,0 +1,212 @@
//! One-time migration of the pre-scoping global retention database into the
//! active relay+owner scope.
//!
//! Before community scoping, every durable event lived in one
//! `<managed-agents-base>/retention.db`. Scoped storage
//! ([`super::scoped_retention_db_path`]) reads a different file, so an upgrade
//! would otherwise abandon whatever the previous release left pending —
//! including signed kind:5 tombstones and NIP-IA archive requests queued while
//! offline, which no reconcile can reconstruct (boot reconcile rebuilds upserts
//! from records still on disk, and deletions have no reconcile at all).
//!
//! # Crash safety
//!
//! Two guards, each written transactionally, make the migration exactly-once
//! without a completion file:
//!
//! 1. A **claim** in the legacy database naming the scope that owns its rows.
//! Legacy rows were queued for whichever single relay the old build had
//! active, so exactly one scope may take them; every other scope skips. This
//! is what keeps the migration from fanning one community's pending events
//! out to all of them — the leak class scoping exists to close.
//! 2. A **marker** in the scoped database, committed in the same transaction as
//! the copied rows. A crash mid-copy therefore leaves neither rows nor
//! marker, and the next boot copies from scratch; once the marker is there
//! the copy never repeats.
//!
//! The relay dimension is not recoverable from the legacy file — only the owner
//! pubkey is — so the claiming scope is the first one this owner activates after
//! upgrading. That is the workspace the app restores at launch, i.e. the same
//! relay the stranded rows were queued for in all but a contrived
//! switch-before-first-flush case.
use std::path::{Path, PathBuf};
use rusqlite::{params, Connection, OptionalExtension};
use super::{open_retention_db, RetainedEvent};
/// Marker/claim identifier for this migration.
const MIGRATION_NAME: &str = "legacy_global_retention_db";
/// The pre-scoping global retention database path.
pub fn legacy_retention_db_path(base_dir: &Path) -> PathBuf {
base_dir.join("retention.db")
}
/// Copy the legacy global database's rows for `owner_pubkey` into the scoped
/// database at `scope_db_path`.
///
/// Returns the number of rows copied — `0` both when there is nothing to do and
/// when another scope already claimed the legacy rows. Best-effort by design:
/// the caller logs a failure and proceeds, and the guards make a later retry
/// safe.
pub fn migrate_legacy_retention_db(
base_dir: &Path,
scope_db_path: &Path,
owner_pubkey: &str,
) -> Result<usize, String> {
let legacy_path = legacy_retention_db_path(base_dir);
if !legacy_path.exists() || legacy_path == scope_db_path {
return Ok(0);
}
let scope_id = scope_identifier(scope_db_path);
let mut scope_conn = open_retention_db(scope_db_path)?;
if migration_marker_present(&scope_conn)? {
return Ok(0);
}
let legacy_conn = open_retention_db(&legacy_path)?;
if !claim_legacy_rows(&legacy_conn, &scope_id)? {
return Ok(0); // another scope owns these rows
}
let rows = legacy_rows_for_owner(&legacy_conn, owner_pubkey)?;
let copied = rows.len();
let transaction = scope_conn
.transaction()
.map_err(|e| format!("failed to open retention migration transaction: {e}"))?;
for row in &rows {
// The scoped database is authoritative for any coordinate it already
// holds: those rows were written after the upgrade, so they are newer
// than anything legacy by construction. Legacy rows only fill gaps.
transaction
.execute(
"INSERT INTO persona_events
(kind, pubkey, d_tag, content, created_at, raw_event, pending_sync)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)
ON CONFLICT (kind, pubkey, d_tag) DO NOTHING",
params![
row.kind,
row.pubkey,
row.d_tag,
row.content,
row.created_at,
row.raw_event,
row.pending_sync as i32,
],
)
.map_err(|e| format!("failed to copy legacy retained event: {e}"))?;
}
write_migration_marker(&transaction, &scope_id)?;
transaction
.commit()
.map_err(|e| format!("failed to commit retention migration: {e}"))?;
Ok(copied)
}
/// Read every retained row authored by `owner_pubkey` from the legacy database.
///
/// Owner-filtered because the flush loop only publishes rows matching the
/// active owner anyway; a different identity's rows belong to that identity's
/// scope, not this one.
fn legacy_rows_for_owner(
conn: &Connection,
owner_pubkey: &str,
) -> Result<Vec<RetainedEvent>, String> {
let mut stmt = conn
.prepare(
"SELECT kind, pubkey, d_tag, content, created_at, raw_event, pending_sync
FROM persona_events
WHERE pubkey = ?1
ORDER BY (kind != 5), created_at ASC",
)
.map_err(|e| format!("failed to prepare legacy retention query: {e}"))?;
let rows = stmt
.query_map(params![owner_pubkey], |row| {
Ok(RetainedEvent {
kind: row.get(0)?,
pubkey: row.get(1)?,
d_tag: row.get(2)?,
content: row.get(3)?,
created_at: row.get(4)?,
raw_event: row.get(5)?,
pending_sync: row.get::<_, i32>(6)? != 0,
})
})
.map_err(|e| format!("failed to query legacy retained events: {e}"))?;
rows.collect::<Result<Vec<_>, _>>()
.map_err(|e| format!("failed to read legacy retained row: {e}"))
}
/// Identify a scope by its database file stem — the relay+owner hash
/// [`super::scoped_retention_db_path`] already computes.
fn scope_identifier(scope_db_path: &Path) -> String {
scope_db_path
.file_stem()
.map(|stem| stem.to_string_lossy().to_string())
.unwrap_or_default()
}
fn ensure_migration_table(conn: &Connection) -> Result<(), String> {
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS retention_migrations (
name TEXT PRIMARY KEY,
scope_id TEXT NOT NULL
);",
)
.map_err(|e| format!("failed to create retention migration table: {e}"))
}
fn migration_marker_present(conn: &Connection) -> Result<bool, String> {
ensure_migration_table(conn)?;
conn.query_row(
"SELECT EXISTS(SELECT 1 FROM retention_migrations WHERE name = ?1)",
params![MIGRATION_NAME],
|row| row.get(0),
)
.map_err(|e| format!("failed to read retention migration marker: {e}"))
}
fn write_migration_marker(conn: &Connection, scope_id: &str) -> Result<(), String> {
ensure_migration_table(conn)?;
conn.execute(
"INSERT OR REPLACE INTO retention_migrations (name, scope_id) VALUES (?1, ?2)",
params![MIGRATION_NAME, scope_id],
)
.map_err(|e| format!("failed to write retention migration marker: {e}"))?;
Ok(())
}
/// Record `scope_id` as the owner of the legacy rows, or confirm it already is.
///
/// `INSERT OR IGNORE` then read-back is atomic enough for this purpose: the
/// loser of a race reads the winner's scope id and returns `false`.
fn claim_legacy_rows(legacy_conn: &Connection, scope_id: &str) -> Result<bool, String> {
ensure_migration_table(legacy_conn)?;
legacy_conn
.execute(
"INSERT OR IGNORE INTO retention_migrations (name, scope_id) VALUES (?1, ?2)",
params![MIGRATION_NAME, scope_id],
)
.map_err(|e| format!("failed to claim legacy retention rows: {e}"))?;
let claimed_by: Option<String> = legacy_conn
.query_row(
"SELECT scope_id FROM retention_migrations WHERE name = ?1",
params![MIGRATION_NAME],
|row| row.get(0),
)
.optional()
.map_err(|e| format!("failed to read legacy retention claim: {e}"))?;
Ok(claimed_by.as_deref() == Some(scope_id))
}
#[cfg(test)]
mod tests;
@@ -0,0 +1,186 @@
use super::*;
use crate::managed_agents::retention::{
get_pending_sync, get_retained_event, retain_event, scoped_retention_db_path,
tombstone_retention_d_tag,
};
use buzz_core_pkg::kind::KIND_PERSONA;
const KIND_DELETE: u32 = 5;
const OWNER: &str = "a1b2c3";
fn pending_tombstone(d_tag: &str) -> RetainedEvent {
RetainedEvent {
kind: KIND_DELETE,
pubkey: OWNER.to_string(),
d_tag: tombstone_retention_d_tag(KIND_PERSONA, d_tag),
content: String::new(),
created_at: 1_700_000_000,
raw_event: format!(r#"{{"kind":5,"d":"{d_tag}"}}"#),
pending_sync: true,
}
}
fn seed_legacy(base_dir: &Path, events: &[RetainedEvent]) {
let conn = open_retention_db(&legacy_retention_db_path(base_dir)).unwrap();
for event in events {
retain_event(&conn, event).unwrap();
}
}
fn scope_path(base_dir: &Path, relay: &str) -> PathBuf {
let path = scoped_retention_db_path(base_dir, relay, OWNER);
std::fs::create_dir_all(path.parent().unwrap()).unwrap();
path
}
#[test]
fn test_pending_legacy_tombstone_migrates_into_the_active_scope_and_stays_pending() {
let dir = tempfile::tempdir().unwrap();
seed_legacy(dir.path(), &[pending_tombstone("retired-agent")]);
let scope = scope_path(dir.path(), "wss://a.example");
let copied = migrate_legacy_retention_db(dir.path(), &scope, OWNER).unwrap();
assert_eq!(copied, 1);
let conn = open_retention_db(&scope).unwrap();
let migrated = get_retained_event(
&conn,
KIND_DELETE,
OWNER,
&tombstone_retention_d_tag(KIND_PERSONA, "retired-agent"),
)
.unwrap()
.expect("legacy tombstone lands in the scoped db");
assert!(
migrated.pending_sync,
"the tombstone must still be queued for the flush loop"
);
assert_eq!(
migrated.raw_event,
pending_tombstone("retired-agent").raw_event
);
assert_eq!(get_pending_sync(&conn).unwrap().len(), 1);
}
#[test]
fn test_repeat_migration_of_the_same_scope_copies_nothing_further() {
let dir = tempfile::tempdir().unwrap();
seed_legacy(dir.path(), &[pending_tombstone("retired-agent")]);
let scope = scope_path(dir.path(), "wss://a.example");
assert_eq!(
migrate_legacy_retention_db(dir.path(), &scope, OWNER).unwrap(),
1
);
// Simulate the flush loop clearing the row, then boot again: the marker
// must stop the legacy row from being resurrected as pending.
let conn = open_retention_db(&scope).unwrap();
conn.execute("UPDATE persona_events SET pending_sync = 0", [])
.unwrap();
drop(conn);
assert_eq!(
migrate_legacy_retention_db(dir.path(), &scope, OWNER).unwrap(),
0
);
let conn = open_retention_db(&scope).unwrap();
assert!(
get_pending_sync(&conn).unwrap().is_empty(),
"a published row must not be re-queued by a second migration pass"
);
}
#[test]
fn test_second_community_does_not_receive_another_communitys_legacy_rows() {
let dir = tempfile::tempdir().unwrap();
seed_legacy(dir.path(), &[pending_tombstone("retired-agent")]);
let first = scope_path(dir.path(), "wss://a.example");
let second = scope_path(dir.path(), "wss://b.example");
assert_eq!(
migrate_legacy_retention_db(dir.path(), &first, OWNER).unwrap(),
1
);
assert_eq!(
migrate_legacy_retention_db(dir.path(), &second, OWNER).unwrap(),
0,
"legacy rows belong to exactly one relay scope"
);
let conn = open_retention_db(&second).unwrap();
assert!(get_pending_sync(&conn).unwrap().is_empty());
}
#[test]
fn test_rows_authored_by_another_identity_are_left_behind() {
let dir = tempfile::tempdir().unwrap();
let mut foreign = pending_tombstone("someone-elses");
foreign.pubkey = "ffffff".to_string();
seed_legacy(dir.path(), &[pending_tombstone("mine"), foreign]);
let scope = scope_path(dir.path(), "wss://a.example");
assert_eq!(
migrate_legacy_retention_db(dir.path(), &scope, OWNER).unwrap(),
1
);
let conn = open_retention_db(&scope).unwrap();
let pending = get_pending_sync(&conn).unwrap();
assert_eq!(pending.len(), 1);
assert_eq!(pending[0].pubkey, OWNER);
}
#[test]
fn test_post_upgrade_scoped_row_is_not_overwritten_by_its_legacy_ancestor() {
let dir = tempfile::tempdir().unwrap();
let legacy_head = RetainedEvent {
kind: KIND_PERSONA,
pubkey: OWNER.to_string(),
d_tag: "reviewer".to_string(),
content: r#"{"display_name":"Old"}"#.to_string(),
created_at: 1_700_000_000,
raw_event: r#"{"content":"old"}"#.to_string(),
pending_sync: true,
};
seed_legacy(dir.path(), &[legacy_head]);
let scope = scope_path(dir.path(), "wss://a.example");
// An edit made after the upgrade already occupies the coordinate.
let conn = open_retention_db(&scope).unwrap();
retain_event(
&conn,
&RetainedEvent {
kind: KIND_PERSONA,
pubkey: OWNER.to_string(),
d_tag: "reviewer".to_string(),
content: r#"{"display_name":"New"}"#.to_string(),
created_at: 1_700_000_500,
raw_event: r#"{"content":"new"}"#.to_string(),
pending_sync: true,
},
)
.unwrap();
drop(conn);
migrate_legacy_retention_db(dir.path(), &scope, OWNER).unwrap();
let conn = open_retention_db(&scope).unwrap();
let row = get_retained_event(&conn, KIND_PERSONA, OWNER, "reviewer")
.unwrap()
.unwrap();
assert_eq!(row.created_at, 1_700_000_500);
assert_eq!(row.raw_event, r#"{"content":"new"}"#);
}
#[test]
fn test_absent_legacy_database_is_a_no_op() {
let dir = tempfile::tempdir().unwrap();
let scope = scope_path(dir.path(), "wss://a.example");
assert_eq!(
migrate_legacy_retention_db(dir.path(), &scope, OWNER).unwrap(),
0
);
assert!(!legacy_retention_db_path(dir.path()).exists());
}