mirror of
https://github.com/block/buzz.git
synced 2026-08-18 06:50:31 +02:00
fix(desktop): bind nest regen in cores, gate on highest-requested
Two review findings on the archive/unarchive nest-regen wiring. Wrapper callback was unprotected. The Tauri command wrappers each constructed the regeneration closure themselves (`|| try_regenerate_nest`) before delegating, so the core tests proved a core forwards its callback but never that either production command supplies one — no-oping either wrapper's closure left the whole suite green. Bind regeneration to a `NestRegenTrigger` type the cores invoke instead: the wrapper hands the core its `AppHandle` as the trigger with no closure to construct, so the 'regenerate on success' selection lives inside the cores where the tests traverse it. Each core's `|| regen.trigger()` is now RED-on-revert. Regen gate permitted stale rollback. `NestRegenGate::commit` gated on the highest *written* generation, which advances only when a newer generation successfully writes. A newer request that failed during relay work therefore left an older, stale render free to publish its obsolete roster. Gate on the highest *requested* generation instead, advanced by `claim` under the same lock `commit` reads — so once a newer generation is requested no older one can publish, even if the newer one fails, and a claim cannot race between an older task's eligibility check and its write. Semantic delta: after a newer request fails, nothing publishes until the next trigger. Co-authored-by: Will Pfleger <pfleger.will@gmail.com> Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
@@ -172,6 +172,22 @@ pub struct UnarchiveRequest {
|
||||
pub reason: Option<String>,
|
||||
}
|
||||
|
||||
/// Roster refresh a successful archive/unarchive triggers. Binding the action
|
||||
/// to a *type* rather than a closure selected at each call site is what closes
|
||||
/// the regression Thufir found: the command wrapper passes a value (`&app`)
|
||||
/// with no callback to construct, so the "regenerate on success" selection
|
||||
/// lives entirely inside the cores below — where the tests traverse it. The
|
||||
/// production binding is the single, irreducible `AppHandle` adapter.
|
||||
pub(crate) trait NestRegenTrigger {
|
||||
fn trigger(&self);
|
||||
}
|
||||
|
||||
impl NestRegenTrigger for AppHandle {
|
||||
fn trigger(&self) {
|
||||
try_regenerate_nest(self);
|
||||
}
|
||||
}
|
||||
|
||||
/// Submit `builder` to the active workspace relay, then trigger `on_success`
|
||||
/// exactly once iff the relay accepted the event.
|
||||
///
|
||||
@@ -190,19 +206,21 @@ async fn submit_then_regenerate(
|
||||
}
|
||||
|
||||
/// `AppHandle`-free core of [`archive_identity`]: resolve the owner-of-agent
|
||||
/// `auth` tag, build the real `kind:9035` request, submit it, and forward
|
||||
/// `on_success` so a successful archive refreshes the roster.
|
||||
/// `auth` tag, build the real `kind:9035` request, submit it, and trigger
|
||||
/// `regen` so a successful archive refreshes the roster.
|
||||
///
|
||||
/// The command wrapper is untestable (it needs a live Tauri runtime for its
|
||||
/// `AppHandle`), so extracting this core is what lets a test drive the exact
|
||||
/// archive wiring — including the forwarded regeneration callback — through a
|
||||
/// loopback relay. RED-on-revert: pass `|| {}` here instead of `on_success` and
|
||||
/// `AppHandle`), so this core owns the whole orchestration — including *binding*
|
||||
/// the regeneration trigger onto the successful-submit path. The wrapper only
|
||||
/// hands it the `AppHandle` as the trigger; a test drives the exact archive
|
||||
/// wiring with a counting trigger over a loopback relay. RED-on-revert: change
|
||||
/// `|| regen.trigger()` to `|| {}` here and
|
||||
/// `archive_core_fires_regen_only_on_accepted_submit` fails while the unarchive
|
||||
/// core test stays green.
|
||||
async fn archive_identity_core(
|
||||
req: &ArchiveRequest,
|
||||
state: &AppState,
|
||||
on_success: impl FnOnce(),
|
||||
regen: &impl NestRegenTrigger,
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
let auth_tag = maybe_owner_auth_tag(state, &req.target_pubkey).await?;
|
||||
let builder = events::build_archive_identity_request(
|
||||
@@ -212,18 +230,18 @@ async fn archive_identity_core(
|
||||
req.replaced_by.as_deref(),
|
||||
auth_tag.as_ref(),
|
||||
)?;
|
||||
submit_then_regenerate(builder, state, on_success).await
|
||||
submit_then_regenerate(builder, state, || regen.trigger()).await
|
||||
}
|
||||
|
||||
/// `AppHandle`-free core of [`unarchive_identity`]: builds the real `kind:9036`
|
||||
/// request and forwards `on_success` on acceptance. See [`archive_identity_core`]
|
||||
/// for why this seam is extracted. RED-on-revert: pass `|| {}` here and
|
||||
/// `unarchive_core_fires_regen_only_on_accepted_submit` fails while the archive
|
||||
/// core test stays green.
|
||||
/// request and triggers `regen` on acceptance. See [`archive_identity_core`]
|
||||
/// for why this seam is extracted. RED-on-revert: change `|| regen.trigger()`
|
||||
/// to `|| {}` here and `unarchive_core_fires_regen_only_on_accepted_submit`
|
||||
/// fails while the archive core test stays green.
|
||||
async fn unarchive_identity_core(
|
||||
req: &UnarchiveRequest,
|
||||
state: &AppState,
|
||||
on_success: impl FnOnce(),
|
||||
regen: &impl NestRegenTrigger,
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
let auth_tag = maybe_owner_auth_tag(state, &req.target_pubkey).await?;
|
||||
let builder = events::build_unarchive_identity_request(
|
||||
@@ -232,7 +250,7 @@ async fn unarchive_identity_core(
|
||||
req.reason.as_deref(),
|
||||
auth_tag.as_ref(),
|
||||
)?;
|
||||
submit_then_regenerate(builder, state, on_success).await
|
||||
submit_then_regenerate(builder, state, || regen.trigger()).await
|
||||
}
|
||||
|
||||
/// Submit a `kind:9035` archive request to the relay. Consent path is selected
|
||||
@@ -251,7 +269,7 @@ pub async fn archive_identity(
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
archive_identity_core(&req, &state, || try_regenerate_nest(&app)).await
|
||||
archive_identity_core(&req, &state, &app).await
|
||||
}
|
||||
|
||||
/// Submit a `kind:9036` unarchive request to the relay. See
|
||||
@@ -263,7 +281,7 @@ pub async fn unarchive_identity(
|
||||
app: AppHandle,
|
||||
state: State<'_, AppState>,
|
||||
) -> Result<SubmitEventResponse, String> {
|
||||
unarchive_identity_core(&req, &state, || try_regenerate_nest(&app)).await
|
||||
unarchive_identity_core(&req, &state, &app).await
|
||||
}
|
||||
|
||||
/// If the current user is the verified NIP-OA owner of `target`, return the
|
||||
@@ -458,6 +476,29 @@ pub async fn get_relay_self(state: State<'_, AppState>) -> Result<Option<String>
|
||||
mod tests {
|
||||
use super::*;
|
||||
use nostr::{EventBuilder, Keys, Kind, Tag};
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
/// Counting [`NestRegenTrigger`] double: records how many times the core
|
||||
/// fires regeneration on the successful-submit path, standing in for the
|
||||
/// production `AppHandle` binding without a live Tauri runtime.
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
#[derive(Default)]
|
||||
struct CountingRegen(AtomicUsize);
|
||||
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
impl CountingRegen {
|
||||
fn count(&self) -> usize {
|
||||
self.0.load(Ordering::SeqCst)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(not(target_os = "windows"))]
|
||||
impl NestRegenTrigger for CountingRegen {
|
||||
fn trigger(&self) {
|
||||
self.0.fetch_add(1, Ordering::SeqCst);
|
||||
}
|
||||
}
|
||||
|
||||
/// Build a fake `kind:0` with a valid NIP-OA auth tag for a fresh owner.
|
||||
fn kind0_with_auth(agent: &Keys, owner: &Keys) -> nostr::Event {
|
||||
@@ -732,7 +773,6 @@ mod tests {
|
||||
async fn archive_core_fires_regen_only_on_accepted_submit() {
|
||||
use crate::app_state::build_app_state;
|
||||
use crate::relay_admission::{reset_rate_limit_gate, TEST_SERIAL};
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
let _serial = TEST_SERIAL.lock().await;
|
||||
reset_rate_limit_gate();
|
||||
@@ -747,29 +787,24 @@ mod tests {
|
||||
|
||||
// Accepted archive → hook fires exactly once.
|
||||
*state.relay_url_override.lock().unwrap() = Some(spawn_submit_relay(true).await);
|
||||
let fired = AtomicUsize::new(0);
|
||||
let response = archive_identity_core(&req, &state, || {
|
||||
fired.fetch_add(1, Ordering::SeqCst);
|
||||
})
|
||||
.await
|
||||
.expect("accepted archive returns Ok");
|
||||
let regen = CountingRegen::default();
|
||||
let response = archive_identity_core(&req, &state, ®en)
|
||||
.await
|
||||
.expect("accepted archive returns Ok");
|
||||
assert!(response.accepted);
|
||||
assert_eq!(
|
||||
fired.load(Ordering::SeqCst),
|
||||
regen.count(),
|
||||
1,
|
||||
"an accepted archive must trigger regeneration exactly once"
|
||||
);
|
||||
|
||||
// Rejected submit → error propagates, hook never fires.
|
||||
*state.relay_url_override.lock().unwrap() = Some(spawn_submit_relay(false).await);
|
||||
let fired = AtomicUsize::new(0);
|
||||
let result = archive_identity_core(&req, &state, || {
|
||||
fired.fetch_add(1, Ordering::SeqCst);
|
||||
})
|
||||
.await;
|
||||
let regen = CountingRegen::default();
|
||||
let result = archive_identity_core(&req, &state, ®en).await;
|
||||
assert!(result.is_err(), "a rejected archive must return an error");
|
||||
assert_eq!(
|
||||
fired.load(Ordering::SeqCst),
|
||||
regen.count(),
|
||||
0,
|
||||
"a rejected archive changed nothing, so regeneration must not fire"
|
||||
);
|
||||
@@ -788,7 +823,6 @@ mod tests {
|
||||
async fn unarchive_core_fires_regen_only_on_accepted_submit() {
|
||||
use crate::app_state::build_app_state;
|
||||
use crate::relay_admission::{reset_rate_limit_gate, TEST_SERIAL};
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
|
||||
let _serial = TEST_SERIAL.lock().await;
|
||||
reset_rate_limit_gate();
|
||||
@@ -802,29 +836,24 @@ mod tests {
|
||||
|
||||
// Accepted unarchive → hook fires exactly once.
|
||||
*state.relay_url_override.lock().unwrap() = Some(spawn_submit_relay(true).await);
|
||||
let fired = AtomicUsize::new(0);
|
||||
let response = unarchive_identity_core(&req, &state, || {
|
||||
fired.fetch_add(1, Ordering::SeqCst);
|
||||
})
|
||||
.await
|
||||
.expect("accepted unarchive returns Ok");
|
||||
let regen = CountingRegen::default();
|
||||
let response = unarchive_identity_core(&req, &state, ®en)
|
||||
.await
|
||||
.expect("accepted unarchive returns Ok");
|
||||
assert!(response.accepted);
|
||||
assert_eq!(
|
||||
fired.load(Ordering::SeqCst),
|
||||
regen.count(),
|
||||
1,
|
||||
"an accepted unarchive must trigger regeneration exactly once"
|
||||
);
|
||||
|
||||
// Rejected submit → error propagates, hook never fires.
|
||||
*state.relay_url_override.lock().unwrap() = Some(spawn_submit_relay(false).await);
|
||||
let fired = AtomicUsize::new(0);
|
||||
let result = unarchive_identity_core(&req, &state, || {
|
||||
fired.fetch_add(1, Ordering::SeqCst);
|
||||
})
|
||||
.await;
|
||||
let regen = CountingRegen::default();
|
||||
let result = unarchive_identity_core(&req, &state, ®en).await;
|
||||
assert!(result.is_err(), "a rejected unarchive must return an error");
|
||||
assert_eq!(
|
||||
fired.load(Ordering::SeqCst),
|
||||
regen.count(),
|
||||
0,
|
||||
"a rejected unarchive changed nothing, so regeneration must not fire"
|
||||
);
|
||||
|
||||
@@ -16,7 +16,6 @@ use std::collections::HashSet;
|
||||
use std::fs;
|
||||
use std::io;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::Mutex;
|
||||
use tauri::{AppHandle, Manager};
|
||||
|
||||
@@ -665,58 +664,71 @@ pub fn upsert_managed_section(file_path: &Path, new_section_content: &str) -> io
|
||||
}
|
||||
|
||||
/// Serializes nest-context writes so a slow, stale regeneration cannot roll the
|
||||
/// file back over a newer one. This is an ordered, latest-write-wins gate — not
|
||||
/// a work coalescer: every superseded generation still performs its relay reads,
|
||||
/// then drops its result at commit time. Adding a true dirty-loop owner would be
|
||||
/// a larger change and is unwarranted at this user-driven trigger rate.
|
||||
/// file back over a newer one. This is an ordered, latest-request-wins gate —
|
||||
/// not a work coalescer: every superseded generation still performs its relay
|
||||
/// reads, then drops its result at commit time. Adding a true dirty-loop owner
|
||||
/// would be a larger change and is unwarranted at this user-driven trigger rate.
|
||||
///
|
||||
/// Each regeneration request claims a monotonic generation *synchronously* at
|
||||
/// request time (see [`NestRegenGate::claim`]), so the generation encodes
|
||||
/// program order: boot's regen is claimed before `apply_workspace`'s, an edit's
|
||||
/// regen before the next edit's. The claimed generation travels with the
|
||||
/// spawned task and gates its write in [`NestRegenGate::commit`]: a task
|
||||
/// whose generation is below the high-water mark drops its result instead of
|
||||
/// overwriting the newer file. The commit lock is held across the compare and
|
||||
/// the synchronous file write, so two tasks cannot both read a stale mark and
|
||||
/// both write — a bare mutex acquired around only the write would still permit
|
||||
/// that stale-last rollback.
|
||||
/// spawned task and gates its write in [`NestRegenGate::commit`]: a task drops
|
||||
/// its result once a *newer generation has been requested*, even if that newer
|
||||
/// generation later fails before it writes. Gating on the highest *requested*
|
||||
/// generation — not the highest *written* one — is what stops a slow, stale
|
||||
/// pre-edit render from publishing after a newer post-edit render was claimed
|
||||
/// and then failed during its relay work (which would otherwise leave the
|
||||
/// obsolete roster authoritative until the next unrelated trigger). Declared
|
||||
/// semantic: once a newer regeneration is requested, no older one publishes;
|
||||
/// if that newer one fails, the file simply waits for the next trigger.
|
||||
///
|
||||
/// `claim` and `commit` share one lock, so the "is this still the newest
|
||||
/// request?" compare is atomic with the synchronous file write. A bare atomic
|
||||
/// watermark checked separately from the write would let a new claim slip
|
||||
/// between an older task's eligibility check and its write; holding the lock
|
||||
/// across both closes that window (no `await` occurs while it is held).
|
||||
struct NestRegenGate {
|
||||
/// Next generation to hand out. `fetch_add` under a single atomic preserves
|
||||
/// the program order of `claim` calls regardless of memory ordering.
|
||||
next_gen: AtomicU64,
|
||||
/// Highest generation already written. `0` means nothing written yet.
|
||||
last_written: Mutex<u64>,
|
||||
/// Highest generation *requested* so far (`0` = none yet). Advanced by
|
||||
/// [`claim`] and read by [`commit`]; guarding both under this single lock
|
||||
/// keeps the eligibility compare atomic with the file write.
|
||||
highest_requested: Mutex<u64>,
|
||||
}
|
||||
|
||||
impl NestRegenGate {
|
||||
const fn new() -> Self {
|
||||
Self {
|
||||
next_gen: AtomicU64::new(1),
|
||||
last_written: Mutex::new(0),
|
||||
highest_requested: Mutex::new(0),
|
||||
}
|
||||
}
|
||||
|
||||
/// Claim the next generation. Call synchronously at request time so the
|
||||
/// value reflects when the regeneration was requested, not when its task
|
||||
/// happens to run.
|
||||
/// happens to run. Advancing the shared watermark here is what lets a later
|
||||
/// [`commit`] recognize — and drop — any older generation's stale render.
|
||||
fn claim(&self) -> u64 {
|
||||
self.next_gen.fetch_add(1, Ordering::SeqCst)
|
||||
let mut requested = self
|
||||
.highest_requested
|
||||
.lock()
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||
*requested += 1;
|
||||
*requested
|
||||
}
|
||||
|
||||
/// Commit `content` for `generation`, dropping the write when a newer
|
||||
/// generation has already committed. Returns whether the file was written.
|
||||
/// The lock spans the compare and the write so the check-and-write is
|
||||
/// atomic and no await occurs while it is held.
|
||||
/// Commit `content` for `generation`, dropping the write once a newer
|
||||
/// generation has been *requested* (regardless of whether that newer
|
||||
/// generation has written or ever will). Returns whether the file was
|
||||
/// written. The lock spans the compare and the write so the check-and-write
|
||||
/// is atomic and no await occurs while it is held.
|
||||
fn commit(&self, agents_md: &Path, content: &str, generation: u64) -> io::Result<bool> {
|
||||
let mut last = self
|
||||
.last_written
|
||||
let requested = self
|
||||
.highest_requested
|
||||
.lock()
|
||||
.map_err(|_| io::Error::other("nest regen gate lock poisoned"))?;
|
||||
if generation < *last {
|
||||
if generation < *requested {
|
||||
return Ok(false);
|
||||
}
|
||||
upsert_managed_section(agents_md, content)?;
|
||||
*last = generation;
|
||||
Ok(true)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -579,6 +579,68 @@ fn commit_boot_fallback_relay_cannot_bury_apply_workspace_relay() {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn commit_failed_newer_request_still_supersedes_older_snapshot() {
|
||||
// Carl 4954831197, case 1: a newer request that never writes must still
|
||||
// permanently supersede an older snapshot. gen1 (pre-edit) is claimed and
|
||||
// its relay work is slow; an edit claims gen2 (post-edit); gen2 then FAILS
|
||||
// during its relay work, so it never commits. When gen1 finally finishes,
|
||||
// it must NOT publish its obsolete roster — gating on highest-*requested*
|
||||
// (advanced by gen2's claim) drops it, whereas gating on highest-*written*
|
||||
// (0, since gen2 never wrote) would wrongly let gen1 publish.
|
||||
let gate = NestRegenGate::new();
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let file = agents_md_with_markers(tmp.path());
|
||||
|
||||
let gen1 = gate.claim(); // pre-edit request
|
||||
let gen2 = gate.claim(); // post-edit request
|
||||
assert!(gen1 < gen2);
|
||||
|
||||
// gen2 fails during relay work and never reaches commit — nothing written.
|
||||
|
||||
// gen1 finishes last; its stale render must be dropped.
|
||||
assert!(
|
||||
!gate.commit(&file, "pre-edit roster", gen1).unwrap(),
|
||||
"an older snapshot must not publish once a newer generation was requested, \
|
||||
even if that newer generation failed before writing"
|
||||
);
|
||||
|
||||
let content = fs::read_to_string(&file).unwrap();
|
||||
assert!(
|
||||
!content.contains("pre-edit roster"),
|
||||
"the obsolete pre-edit roster must never become authoritative"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn commit_claim_at_the_older_tasks_cutover_supersedes_it() {
|
||||
// Carl 4954831197, case 2: a claim arriving at the older task's commit
|
||||
// cutover must supersede it. Model the exact interleaving the shared lock
|
||||
// must forbid: gen1 becomes eligible (it is the highest request), then gen2
|
||||
// is claimed *before* gen1 actually writes. Because claim advances the same
|
||||
// watermark commit reads under one lock, gen1's write is rejected — a bare
|
||||
// atomic checked separately from the write would have let gen1 slip through.
|
||||
let gate = NestRegenGate::new();
|
||||
let tmp = tempfile::tempdir().unwrap();
|
||||
let file = agents_md_with_markers(tmp.path());
|
||||
|
||||
let gen1 = gate.claim();
|
||||
// A new request lands at the cutover, before gen1's commit runs.
|
||||
let gen2 = gate.claim();
|
||||
assert!(gen1 < gen2);
|
||||
|
||||
assert!(
|
||||
!gate.commit(&file, "gen1 roster", gen1).unwrap(),
|
||||
"a claim arriving before the older task's write must supersede it"
|
||||
);
|
||||
// gen2 is the surviving request and may publish.
|
||||
assert!(gate.commit(&file, "gen2 roster", gen2).unwrap());
|
||||
|
||||
let content = fs::read_to_string(&file).unwrap();
|
||||
assert!(content.contains("gen2 roster"));
|
||||
assert!(!content.contains("gen1 roster"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn commit_equal_generation_is_allowed() {
|
||||
// The gate rejects only strictly-lower generations. Re-committing the same
|
||||
@@ -612,7 +674,7 @@ fn commit_poisoned_lock_returns_error_instead_of_panicking() {
|
||||
|
||||
let poisoner = gate.clone();
|
||||
let _ = std::thread::spawn(move || {
|
||||
let _guard = poisoner.last_written.lock().unwrap();
|
||||
let _guard = poisoner.highest_requested.lock().unwrap();
|
||||
panic!("poison the gate lock");
|
||||
})
|
||||
.join();
|
||||
|
||||
Reference in New Issue
Block a user