Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
138 changes: 98 additions & 40 deletions crates/buzz-relay/src/handlers/ingest.rs
Original file line number Diff line number Diff line change
Expand Up @@ -195,6 +195,38 @@ fn emit_product_feedback_success(
);
}

/// Build the abstract write [`TraceAction`] for a persist outcome.
///
/// `channel_id.is_some()` distinguishes channel-bearing writes
/// (Insert/Duplicate) from channel-less writes (InsertGlobal);
/// `was_inserted` distinguishes accepted-new writes from no-op-on-conflict
/// duplicates. This mirrors the spec's three-way split at the `WriteInsert`
/// (line 514), `WriteInsertGlobal` (line 559), and `WriteDuplicate`
/// (line 606) seams.
///
/// Shared by the reaction path and the general message-write path so the
/// two seams cannot diverge — see #4936 for the bug that motivated
/// extracting this helper.
fn write_trace_action(event: &Event, channel_id: Option<Uuid>, was_inserted: bool) -> TraceAction {
let claimed = claimed_community_from_event(event);
match (channel_id, was_inserted) {
(Some(ch), true) => TraceAction::WriteInsert {
msg_id: msg_id_label(event.id.as_bytes()),
channel: channel_label(ch),
claimed_community: claimed,
},
(Some(ch), false) => TraceAction::WriteDuplicate {
msg_id: msg_id_label(event.id.as_bytes()),
channel: channel_label(ch),
claimed_community: claimed,
},
(None, _) => TraceAction::WriteInsertGlobal {
msg_id: msg_id_label(event.id.as_bytes()),
claimed_community: claimed,
},
}
}

/// Increment the rejection counter with a bounded reason and transport label.
///
/// Shared by the WS `EVENT` handler and the HTTP `POST /events` handler so
Expand Down Expand Up @@ -2764,25 +2796,14 @@ async fn ingest_event_inner(
};

let pubkey_hex = auth.pubkey().to_hex();
// Spec WriteInsert (line 514) / WriteDuplicate (line 606): emit
// the abstract write action. The persist API returns
// `was_inserted` (true → Insert, false → Duplicate). This branch
// is the reaction path; channel_id is always Some here, so
// WriteInsertGlobal does not apply.
let claimed = claimed_community_from_event(&event);
let action = if was_inserted {
TraceAction::WriteInsert {
msg_id: msg_id_label(event.id.as_bytes()),
channel: channel_label(channel_id.expect("reaction path has channel")),
claimed_community: claimed,
}
} else {
TraceAction::WriteDuplicate {
msg_id: msg_id_label(event.id.as_bytes()),
channel: channel_label(channel_id.expect("reaction path has channel")),
claimed_community: claimed,
}
};
// Spec WriteInsert (line 514) / WriteInsertGlobal (line 559) /
// WriteDuplicate (line 606): emit the abstract write action. A
// reaction on a project event (kind 1621 issue, 1618 PR, or a
// kind-1 comment on one) carries no `h` tag, so
// `derive_reaction_channel` returns `NoChannel` and `channel_id`
// is `None`. The shared `write_trace_action` helper mirrors the
// message write path's three-way split — see #4936.
let action = write_trace_action(&event, channel_id, was_inserted);
emit(tracer, action, state_for_request(tenant, auth.pubkey()));
dispatch_persistent_event(
tenant,
Expand Down Expand Up @@ -2900,31 +2921,12 @@ async fn ingest_event_inner(
let pubkey_hex = auth.pubkey().to_hex();
// Spec WriteInsert (line 514) / WriteInsertGlobal (line 559) /
// WriteDuplicate (line 606): emit the abstract write at the trailing
// dispatch site. `channel_id.is_some()` distinguishes channel-bearing
// (Insert/Duplicate) from channel-less (InsertGlobal); `was_inserted`
// distinguishes accepted-new (Insert/Global) from no-op-on-conflict
// (Duplicate). The WriteInsertGlobal duplicate case is not modeled
// dispatch site. The WriteInsertGlobal duplicate case is not modeled
// separately in the spec (channel-less duplicates collapse to the
// same observation shape as channel-less inserts at this seam);
// see docs/spec/MultiTenantRelay.tla lines 559-595.
{
let claimed = claimed_community_from_event(&event);
let action = match (channel_id, was_inserted) {
(Some(ch), true) => TraceAction::WriteInsert {
msg_id: msg_id_label(event.id.as_bytes()),
channel: channel_label(ch),
claimed_community: claimed,
},
(Some(ch), false) => TraceAction::WriteDuplicate {
msg_id: msg_id_label(event.id.as_bytes()),
channel: channel_label(ch),
claimed_community: claimed,
},
(None, _) => TraceAction::WriteInsertGlobal {
msg_id: msg_id_label(event.id.as_bytes()),
claimed_community: claimed,
},
};
let action = write_trace_action(&event, channel_id, was_inserted);
emit(tracer, action, state_for_request(tenant, auth.pubkey()));
}
dispatch_persistent_event(
Expand Down Expand Up @@ -4919,4 +4921,60 @@ mod tests {
Some(&1)
);
}

// --- write_trace_action regression tests (#4936) ---
// A reaction on a project event (kind 1621 issue, 1618 PR, or a
// kind-1 comment on one) carries no `h` tag, so `derive_reaction_channel`
// returns `NoChannel` and `channel_id` is `None`. The old code
// panicked on `channel_id.expect("reaction path has channel")`; the
// shared `write_trace_action` helper must emit `WriteInsertGlobal`
// instead.

fn make_reaction_event() -> Event {
EventBuilder::new(Kind::Custom(KIND_REACTION as u16), "+")
.sign_with_keys(&nostr::Keys::generate())
.expect("sign reaction")
}

#[test]
fn write_trace_action_channel_less_insert_emits_write_insert_global() {
let event = make_reaction_event();
let action = write_trace_action(&event, None, true);
assert!(
matches!(action, TraceAction::WriteInsertGlobal { .. }),
"channel-less insert must emit WriteInsertGlobal, got {action:?}"
);
}

#[test]
fn write_trace_action_channel_less_duplicate_emits_write_insert_global() {
let event = make_reaction_event();
let action = write_trace_action(&event, None, false);
assert!(
matches!(action, TraceAction::WriteInsertGlobal { .. }),
"channel-less duplicate must emit WriteInsertGlobal, got {action:?}"
);
}

#[test]
fn write_trace_action_channel_bearing_insert_emits_write_insert() {
let event = make_reaction_event();
let ch = Uuid::new_v4();
let action = write_trace_action(&event, Some(ch), true);
assert!(
matches!(action, TraceAction::WriteInsert { .. }),
"channel-bearing insert must emit WriteInsert, got {action:?}"
);
}

#[test]
fn write_trace_action_channel_bearing_duplicate_emits_write_duplicate() {
let event = make_reaction_event();
let ch = Uuid::new_v4();
let action = write_trace_action(&event, Some(ch), false);
assert!(
matches!(action, TraceAction::WriteDuplicate { .. }),
"channel-bearing duplicate must emit WriteDuplicate, got {action:?}"
);
}
}
Loading