mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): make inbound 30177 frozen-linkage convergence observable and retryable
The §2.8 canonical-linkage freeze path retained the inbound event as the head before attempting the corrective library-authoritative re-retain, then swallowed a failed re-retain to stderr and returned Ok. That left the non-authoritative head retained with no recovery: replay is dead because the same event re-arriving is Skipped at the equal-created_at guard before the convergence branch runs. Propagate the corrective re-retain failure (converge_frozen_linkage) so the command cannot report success over a divergent head. The durable retry owner is the boot-time reconcile_agents_to_events pass, which re-diffs the still-authoritative on-disk record against the retained head every launch and re-queues the corrective row at a monotonic bump. Surface the freeze reason as a typed InboundReconcileOutcome (reconcile_inbound_persona_event now returns it) instead of stderr-only, and consume it in usePersonaSync. Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
@@ -56,7 +56,7 @@ pub async fn reconcile_inbound_persona_event(
|
||||
event_json: String,
|
||||
arrival_relay_url: String,
|
||||
app: AppHandle,
|
||||
) -> Result<(), String> {
|
||||
) -> Result<InboundReconcileOutcome, String> {
|
||||
tokio::task::spawn_blocking(move || {
|
||||
reconcile_inbound_persona_event_blocking(event_json, arrival_relay_url, app)
|
||||
})
|
||||
@@ -68,7 +68,7 @@ fn reconcile_inbound_persona_event_blocking(
|
||||
event_json: String,
|
||||
arrival_relay_url: String,
|
||||
app: AppHandle,
|
||||
) -> Result<(), String> {
|
||||
) -> Result<InboundReconcileOutcome, String> {
|
||||
use crate::managed_agents::{
|
||||
agent_events::managed_agent_content_from_event,
|
||||
load_managed_agents, load_teams,
|
||||
@@ -93,11 +93,12 @@ fn reconcile_inbound_persona_event_blocking(
|
||||
// in its `a` tag (`<target_kind>:<owner>:<d_tag>`). Handled before the
|
||||
// upsert dispatch because its coordinate and retention key differ.
|
||||
if kind == KIND_DELETION {
|
||||
return reconcile_inbound_tombstone(&event, &arrival_relay_url, &app, &state);
|
||||
return reconcile_inbound_tombstone(&event, &arrival_relay_url, &app, &state)
|
||||
.map(|()| InboundReconcileOutcome::default());
|
||||
}
|
||||
|
||||
if !matches!(kind, KIND_PERSONA | KIND_TEAM | KIND_MANAGED_AGENT) {
|
||||
return Ok(());
|
||||
return Ok(InboundReconcileOutcome::default());
|
||||
}
|
||||
|
||||
// The d-tag identifies the record within its kind. Persona derives it from
|
||||
@@ -132,7 +133,7 @@ fn reconcile_inbound_persona_event_blocking(
|
||||
&arrival_owner_pubkey,
|
||||
)?
|
||||
else {
|
||||
return Ok(());
|
||||
return Ok(InboundReconcileOutcome::default());
|
||||
};
|
||||
|
||||
// Library-projection preflight (§2.7): a projected persona is
|
||||
@@ -149,7 +150,7 @@ fn reconcile_inbound_persona_event_blocking(
|
||||
if MutationRoute::for_persona_d_tag(&raw_definitions, &d_tag)
|
||||
== MutationRoute::LibraryProjected
|
||||
{
|
||||
return Ok(());
|
||||
return Ok(InboundReconcileOutcome::default());
|
||||
}
|
||||
}
|
||||
|
||||
@@ -167,9 +168,10 @@ fn reconcile_inbound_persona_event_blocking(
|
||||
},
|
||||
)?;
|
||||
if outcome == InboundOutcome::Skipped {
|
||||
return Ok(());
|
||||
return Ok(InboundReconcileOutcome::default());
|
||||
}
|
||||
|
||||
let mut result = InboundReconcileOutcome::default();
|
||||
match kind {
|
||||
KIND_PERSONA => {
|
||||
let mut personas = load_personas(&app)?;
|
||||
@@ -206,19 +208,19 @@ fn reconcile_inbound_persona_event_blocking(
|
||||
// mechanism §2.7's 30175 rule uses. The frozen record still exists
|
||||
// (only its linkage was left intact), so it always has a projection
|
||||
// to reassert.
|
||||
//
|
||||
// The corrective retain is NOT best-effort: if it fails, the
|
||||
// non-authoritative inbound head is left retained and ordinary
|
||||
// replay cannot repair it (the same event re-arriving is `Skipped`
|
||||
// at the equal-`created_at` guard above). Propagating the error is
|
||||
// what keeps the command from reporting success over a divergent
|
||||
// head — the durable retry owner is the boot-time
|
||||
// `reconcile_agents_to_events` pass, which re-diffs the
|
||||
// still-authoritative on-disk record against the retained head every
|
||||
// launch and re-queues this same corrective row.
|
||||
if let InboundAgentLinkage::Frozen(reason) = linkage {
|
||||
if let Some(local) = agents.iter().find(|record| record.pubkey == d_tag) {
|
||||
if let Err(e) = crate::managed_agents::reconcile::retain_agent_record(
|
||||
&conn,
|
||||
&scope.owner_keys,
|
||||
local,
|
||||
) {
|
||||
eprintln!(
|
||||
"buzz-desktop: inbound 30177 convergence re-retain failed for \
|
||||
{d_tag}: {e}"
|
||||
);
|
||||
}
|
||||
}
|
||||
converge_frozen_linkage(&conn, &scope.owner_keys, &agents, &d_tag)?;
|
||||
result.degradation = Some(LinkageDegradation::new(reason, &d_tag));
|
||||
eprintln!("buzz-desktop: {}", reason.note(&d_tag));
|
||||
}
|
||||
}
|
||||
@@ -230,7 +232,7 @@ fn reconcile_inbound_persona_event_blocking(
|
||||
// land on disk silently, leaving the Agents tab stale until restart.
|
||||
let _ = app.emit("agents-data-changed", ());
|
||||
|
||||
Ok(())
|
||||
Ok(result)
|
||||
}
|
||||
|
||||
/// Parse an inbound wire event and enforce the signature gate. Everything
|
||||
@@ -438,6 +440,45 @@ fn apply_inbound_persona(personas: &mut Vec<AgentDefinition>, inbound: AgentDefi
|
||||
}
|
||||
}
|
||||
|
||||
/// The typed result of an inbound persona/team/managed-agent reconcile.
|
||||
///
|
||||
/// Non-agent kinds, no-match, and admissible-linkage reconciles return the
|
||||
/// default (`degradation: None`). Only a §2.8 frozen-linkage reconcile sets
|
||||
/// `degradation`, so the frontend can observe the interim posture rather than
|
||||
/// having it discarded to stderr. No new UI is wired in this phase — the field
|
||||
/// is a surface the frontend CAN consume.
|
||||
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct InboundReconcileOutcome {
|
||||
/// Set when an inbound kind:30177 event's linkage authorship was rejected
|
||||
/// under the §2.8 canonical-linkage rule and the local head was re-retained
|
||||
/// to converge the relay back.
|
||||
pub degradation: Option<LinkageDegradation>,
|
||||
}
|
||||
|
||||
/// A typed, frontend-consumable description of a frozen-linkage degradation
|
||||
/// (§2.8). Carries the machine-readable reason and the affected agent pubkey.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct LinkageDegradation {
|
||||
/// Machine-readable freeze reason.
|
||||
pub reason: LinkageFreezeReason,
|
||||
/// The affected managed-agent pubkey (the event's d-tag).
|
||||
pub agent_pubkey: String,
|
||||
/// Human-readable note, identical to the log line.
|
||||
pub message: String,
|
||||
}
|
||||
|
||||
impl LinkageDegradation {
|
||||
fn new(reason: LinkageFreezeReason, agent_pubkey: &str) -> Self {
|
||||
Self {
|
||||
reason,
|
||||
agent_pubkey: agent_pubkey.to_string(),
|
||||
message: reason.note(agent_pubkey),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The outcome of applying an inbound kind:30177 event under the §2.8
|
||||
/// canonical-linkage rule.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
@@ -457,8 +498,9 @@ enum InboundAgentLinkage {
|
||||
/// Why an inbound kind:30177 event's linkage authorship was rejected (§2.8).
|
||||
/// Typed so the convergence + degradation surfacing at the call site is not a
|
||||
/// bare string.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum LinkageFreezeReason {
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub enum LinkageFreezeReason {
|
||||
/// The matched instance is currently linked to a library-projected
|
||||
/// definition and the inbound event would clear or re-point its
|
||||
/// `persona_id`. Linkage is owned by the library machinery — the same
|
||||
@@ -583,6 +625,39 @@ fn apply_inbound_managed_agent(
|
||||
InboundAgentLinkage::Applied
|
||||
}
|
||||
|
||||
/// §2.8 convergence: re-retain the local managed-agent record's projection at a
|
||||
/// monotonically newer `created_at` so the relay head converges back to the
|
||||
/// library-authoritative linkage after a frozen inbound event was retained.
|
||||
///
|
||||
/// This is the corrective step that must NOT be best-effort. The inbound event
|
||||
/// is already the retained head (retained before the store apply); if this
|
||||
/// re-retain fails, that non-authoritative head is left in place and ordinary
|
||||
/// replay cannot repair it — the same event re-arriving is `Skipped` at the
|
||||
/// equal-`created_at` guard before the convergence branch runs. Propagating the
|
||||
/// error keeps the command from reporting success over a divergent head; the
|
||||
/// durable retry owner is the boot-time `reconcile_agents_to_events` pass, which
|
||||
/// re-diffs the still-authoritative on-disk record and re-queues this same row.
|
||||
///
|
||||
/// The frozen record still exists (only its linkage was left intact), so it
|
||||
/// always has a projection to reassert; a missing local record would mean the
|
||||
/// freeze classification and the store diverged and is treated as an error.
|
||||
fn converge_frozen_linkage(
|
||||
conn: &rusqlite::Connection,
|
||||
owner_keys: &nostr::Keys,
|
||||
agents: &[ManagedAgentRecord],
|
||||
d_tag: &str,
|
||||
) -> Result<(), String> {
|
||||
let local = agents
|
||||
.iter()
|
||||
.find(|record| record.pubkey == d_tag)
|
||||
.ok_or_else(|| {
|
||||
format!("inbound 30177 convergence: frozen record {d_tag} not found locally")
|
||||
})?;
|
||||
crate::managed_agents::reconcile::retain_agent_record(conn, owner_keys, local)
|
||||
.map_err(|e| format!("inbound 30177 convergence re-retain failed for {d_tag}: {e}"))?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Merge an inbound kind:30176 team projection into the local set.
|
||||
///
|
||||
/// Matches the local record whose `id` equals the event's d-tag (the d-tag IS
|
||||
|
||||
@@ -557,6 +557,74 @@ fn inbound_30177_no_linkage_change_applies_even_when_projected() {
|
||||
);
|
||||
}
|
||||
|
||||
// ── §2.8 convergence: corrective re-retain is not best-effort (P1-I1) ─────
|
||||
|
||||
/// The healthy convergence path re-retains the LOCAL record's authoritative
|
||||
/// projection under its coordinate, queued for publish — this is the row that
|
||||
/// makes the relay head converge back after a frozen inbound event was
|
||||
/// retained.
|
||||
#[test]
|
||||
fn converge_frozen_linkage_re_retains_local_head_pending() {
|
||||
use crate::managed_agents::retention::{get_retained_event, open_retention_db};
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
let conn = open_retention_db(&dir.path().join("retention.db")).unwrap();
|
||||
let keys = nostr::Keys::generate();
|
||||
let agents = vec![local_agent()];
|
||||
|
||||
converge_frozen_linkage(&conn, &keys, &agents, AGENT_PUBKEY).unwrap();
|
||||
|
||||
let row = get_retained_event(
|
||||
&conn,
|
||||
buzz_core_pkg::kind::KIND_MANAGED_AGENT,
|
||||
&keys.public_key().to_hex(),
|
||||
AGENT_PUBKEY,
|
||||
)
|
||||
.unwrap()
|
||||
.expect("convergence must retain the local head");
|
||||
assert!(row.pending_sync, "corrective row must queue for publish");
|
||||
assert!(
|
||||
row.content.contains("Local Agent"),
|
||||
"retained row must be the local authoritative projection",
|
||||
);
|
||||
}
|
||||
|
||||
/// A failed corrective re-retain MUST propagate as `Err`, never be swallowed.
|
||||
/// The inbound event is already the retained head when this runs; if it were
|
||||
/// swallowed the command would report success while a non-authoritative head
|
||||
/// stayed retained (and replay is dead — the same event re-arriving is
|
||||
/// `Skipped` at the equal-`created_at` guard). A connection with no
|
||||
/// `persona_events` table models the retention write failing after the inbound
|
||||
/// retain: `retain_agent_record`'s first query fails.
|
||||
#[test]
|
||||
fn converge_frozen_linkage_errs_when_retain_fails() {
|
||||
let poisoned = rusqlite::Connection::open_in_memory().unwrap();
|
||||
let keys = nostr::Keys::generate();
|
||||
let agents = vec![local_agent()];
|
||||
|
||||
let err = converge_frozen_linkage(&poisoned, &keys, &agents, AGENT_PUBKEY).unwrap_err();
|
||||
assert!(
|
||||
err.contains("convergence"),
|
||||
"a failed corrective re-retain must surface as an error: {err}"
|
||||
);
|
||||
}
|
||||
|
||||
/// A frozen classification with no matching local record is an internal
|
||||
/// inconsistency (the freeze was decided against a record that must exist) and
|
||||
/// fails closed rather than silently no-oping the convergence.
|
||||
#[test]
|
||||
fn converge_frozen_linkage_errs_when_record_missing() {
|
||||
use crate::managed_agents::retention::open_retention_db;
|
||||
let dir = tempfile::TempDir::new().unwrap();
|
||||
let conn = open_retention_db(&dir.path().join("retention.db")).unwrap();
|
||||
let keys = nostr::Keys::generate();
|
||||
|
||||
let err = converge_frozen_linkage(&conn, &keys, &[], AGENT_PUBKEY).unwrap_err();
|
||||
assert!(
|
||||
err.contains("not found"),
|
||||
"a missing frozen record must fail closed: {err}"
|
||||
);
|
||||
}
|
||||
|
||||
// ── Team (30176) inbound ─────────────────────────────────────────────────
|
||||
|
||||
const TEAM_ID: &str = "team-local-id";
|
||||
|
||||
@@ -308,6 +308,69 @@ fn slimming_republish_wave_is_one_time() {
|
||||
// engine both the boot reconcile and the interactive edit paths
|
||||
// (`retain_managed_agent_pending`, persona-rename propagation) run on.
|
||||
|
||||
/// Durable retry owner (P1-I1): the boot-time reconcile re-queues the
|
||||
/// authoritative row over a hostile retained head. This is what makes the
|
||||
/// inbound §2.8 convergence recoverable even if its corrective re-retain failed
|
||||
/// — the on-disk record is authoritative, and every launch re-diffs it against
|
||||
/// the retained head. Models a frozen inbound event that landed a foreign
|
||||
/// projection at a future `created_at`: the next boot re-retains the local
|
||||
/// record's real projection at a monotonic bump past it, queued for publish.
|
||||
#[test]
|
||||
fn boot_reconcile_requeues_authoritative_row_over_hostile_head() {
|
||||
let dir = TempDir::new().unwrap();
|
||||
let keys = nostr::Keys::generate();
|
||||
let owner = keys.public_key().to_hex();
|
||||
let pubkey = "a".repeat(64);
|
||||
let record = sample_record(&pubkey, "Authoritative");
|
||||
write_store(&dir, &[record]);
|
||||
|
||||
// Seed a hostile retained head: a foreign projection at a far-future
|
||||
// created_at, exactly what a frozen inbound 30177 leaves behind when its
|
||||
// corrective re-retain failed.
|
||||
let hostile_created_at = (nostr::Timestamp::now().as_secs() as i64) + 86_400;
|
||||
{
|
||||
let conn = open_retention_db(&dir.path().join("retention.db")).unwrap();
|
||||
retain_event(
|
||||
&conn,
|
||||
&RetainedEvent {
|
||||
kind: KIND_MANAGED_AGENT,
|
||||
pubkey: owner.clone(),
|
||||
d_tag: pubkey.clone(),
|
||||
content: serde_json::json!({ "name": "Hostile Repoint" }).to_string(),
|
||||
created_at: hostile_created_at,
|
||||
raw_event: String::new(),
|
||||
pending_sync: false,
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
// Boot reconcile: the on-disk record's projection differs from the hostile
|
||||
// head, so it re-queues the authoritative row.
|
||||
assert_eq!(reconcile_agents_in_dir(dir.path(), &keys).unwrap(), 1);
|
||||
|
||||
let conn = open_retention_db(&dir.path().join("retention.db")).unwrap();
|
||||
let row = get_retained_event(&conn, KIND_MANAGED_AGENT, &owner, &pubkey)
|
||||
.unwrap()
|
||||
.unwrap();
|
||||
assert!(
|
||||
row.content.contains("Authoritative"),
|
||||
"boot reconcile must restore the authoritative projection"
|
||||
);
|
||||
assert!(
|
||||
!row.content.contains("Hostile"),
|
||||
"the hostile head must be superseded"
|
||||
);
|
||||
assert!(
|
||||
row.pending_sync,
|
||||
"the corrective row must queue for publish"
|
||||
);
|
||||
assert!(
|
||||
row.created_at > hostile_created_at,
|
||||
"created_at must bump past the hostile head so the relay accepts it"
|
||||
);
|
||||
}
|
||||
|
||||
/// A rename re-retains the identity record under the SAME coordinate (the
|
||||
/// agent pubkey) with the new name, queued for publish, with a created_at
|
||||
/// strictly past the retained head so the relay's replaceable-event rule
|
||||
|
||||
@@ -37,11 +37,24 @@ export function startPersonaSync(
|
||||
): () => Promise<void> {
|
||||
const reconcile = (event: RelayEvent) => {
|
||||
if (event.pubkey !== pubkey) return;
|
||||
void reconcileInboundPersonaEvent(JSON.stringify(event), relayUrl).catch(
|
||||
(error) => {
|
||||
void reconcileInboundPersonaEvent(JSON.stringify(event), relayUrl)
|
||||
.then((outcome) => {
|
||||
// §2.8 interim posture: a frozen-linkage reconcile re-retained the
|
||||
// local head to converge the relay back. Surface the typed degradation
|
||||
// rather than letting it vanish — no UI is wired in this phase, so a
|
||||
// console warning is the minimal observable surface. A rejected invoke
|
||||
// (below) means the corrective re-retain FAILED and the non-authoritative
|
||||
// head is still retained; the boot-time reconcile is the durable retry.
|
||||
if (outcome.degradation) {
|
||||
console.warn(
|
||||
"[usePersonaSync] inbound linkage frozen (§2.8):",
|
||||
outcome.degradation.message,
|
||||
);
|
||||
}
|
||||
})
|
||||
.catch((error) => {
|
||||
console.warn("[usePersonaSync] reconcile failed:", error);
|
||||
},
|
||||
);
|
||||
});
|
||||
};
|
||||
|
||||
// One-shot backfill of existing heads + tombstones (closes the fresh-start
|
||||
|
||||
@@ -443,17 +443,45 @@ export async function confirmAgentSnapshotImport(
|
||||
);
|
||||
}
|
||||
|
||||
// Machine-readable reason an inbound kind:30177 event's linkage authorship was
|
||||
// rejected under the §2.8 canonical-linkage rule. Mirrors the Rust
|
||||
// `LinkageFreezeReason` (camelCase-serialized).
|
||||
export type LinkageFreezeReason = "ownedByLibrary" | "inadmissibleNewLink";
|
||||
|
||||
// A frozen-linkage degradation surfaced by `reconcileInboundPersonaEvent`.
|
||||
// Mirrors the Rust `LinkageDegradation`.
|
||||
export type LinkageDegradation = {
|
||||
reason: LinkageFreezeReason;
|
||||
agentPubkey: string;
|
||||
message: string;
|
||||
};
|
||||
|
||||
// Typed result of `reconcileInboundPersonaEvent`. `degradation` is null unless
|
||||
// the reconcile froze a §2.8 linkage change. Mirrors the Rust
|
||||
// `InboundReconcileOutcome`.
|
||||
export type InboundReconcileOutcome = {
|
||||
degradation: LinkageDegradation | null;
|
||||
};
|
||||
|
||||
// Patches a single inbound persona/team/agent projection event into the local
|
||||
// store (personas.json). The backend resolves the match key and the
|
||||
// pending-edit race; the frontend forwards the raw Nostr event JSON plus the
|
||||
// relay it arrived on, so a workspace switch mid-flight cannot retain the event
|
||||
// into the newly active community's scoped store.
|
||||
//
|
||||
// Returns a typed outcome. Its `degradation` is set only when an inbound
|
||||
// kind:30177 event's linkage authorship was rejected under the §2.8
|
||||
// canonical-linkage rule (interim fail-closed posture) and the local head was
|
||||
// re-retained to converge the relay back; otherwise it is null.
|
||||
export async function reconcileInboundPersonaEvent(
|
||||
eventJson: string,
|
||||
arrivalRelayUrl: string,
|
||||
): Promise<void> {
|
||||
await invokeTauri("reconcile_inbound_persona_event", {
|
||||
eventJson,
|
||||
arrivalRelayUrl,
|
||||
});
|
||||
): Promise<InboundReconcileOutcome> {
|
||||
return invokeTauri<InboundReconcileOutcome>(
|
||||
"reconcile_inbound_persona_event",
|
||||
{
|
||||
eventJson,
|
||||
arrivalRelayUrl,
|
||||
},
|
||||
);
|
||||
}
|
||||
|
||||
@@ -11853,7 +11853,9 @@ export function maybeInstallE2eTauriMocks() {
|
||||
}
|
||||
// Mirror the real Rust backend: emit "agents-data-changed" after reconcile.
|
||||
await emit("agents-data-changed");
|
||||
return undefined;
|
||||
// Mirror the typed InboundReconcileOutcome — the e2e mock never
|
||||
// exercises the §2.8 frozen-linkage path, so degradation is always null.
|
||||
return { degradation: null };
|
||||
}
|
||||
case "set_persona_active":
|
||||
return handleSetPersonaActive(
|
||||
|
||||
Reference in New Issue
Block a user