fix(relay): close channel-window bypass and fix e2e draft-expiry tests

Route channel-window rows through reader_can_receive_event so expired
drafts are suppressed on the top_level bridge path — previously this
surface serialized rows directly, bypassing the canonical read gate.

Rewrite e2e fixtures to use a short future expiration tag (now + 2s)
instead of a 31-day-old created_at that ingest rejects (±15 min drift
limit). Add a third e2e test covering the channel-window surface. All
three tests poll rather than fixed-sleep.

Move the canonical-gate rustdoc block back onto reader_can_receive_event
(it was incorrectly placed on draft_expired).

Accepted tradeoff: self-authored draft COUNTs now route through the
5000-candidate fallback rather than raw COUNT(*). A user with >5000
drafts matching one filter hits the "restricted: narrower constraints"
rejection — not a real scenario.

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-14 00:21:25 -04:00
parent a6330254fc
commit cf593a4f5a
3 changed files with 174 additions and 88 deletions
+10 -11
View File
@@ -42,17 +42,6 @@ pub fn is_author_only_event(event: &nostr::Event, requester_pubkey_bytes: &[u8])
&& event.pubkey.to_bytes() != requester_pubkey_bytes
}
/// Canonical per-event read-authorization gate: combines `reader_authorized_for_event`
/// (p-gated/result-gated kinds) and `is_author_only_event` (author-private kinds)
/// into a single predicate.
///
/// Every delivery surface — WS historical pull, WS fan-out, HTTP bridge (all
/// branches: feed, thread, search, catchall, channel-window), and COUNT fallback
/// — must pass each event through this function before serializing it to the wire.
/// Using one canonical gate instead of composing the two predicates at each call
/// site prevents future read surfaces from accidentally omitting half the privacy
/// model.
///
/// Returns `true` if a draft event should be suppressed at read time due to expiry.
///
/// Only applies to `KIND_DRAFT` events; all other kinds return `false`.
@@ -88,6 +77,16 @@ pub fn draft_expired(event: &nostr::Event, now: nostr::Timestamp) -> bool {
now.as_secs() >= effective_expiry
}
/// Canonical per-event read-authorization gate: combines `reader_authorized_for_event`
/// (p-gated/result-gated kinds), `is_author_only_event` (author-private kinds), and
/// `draft_expired` (time-based draft suppression) into a single predicate.
///
/// Every delivery surface — WS historical pull, WS fan-out, HTTP bridge (all
/// branches: feed, thread, search, catchall, channel-window), and COUNT fallback
/// — must pass each event through this function before serializing it to the wire.
/// Using one canonical gate instead of composing the predicates at each call site
/// prevents future read surfaces from accidentally omitting half the privacy model.
///
/// Returns `true` if `reader` MAY receive the event.
pub fn reader_can_receive_event(
event: &nostr::Event,
+13 -1
View File
@@ -367,6 +367,7 @@ const WINDOW_AUX_DELETE_KINDS: [u32; 2] = [
/// Validation errors (missing `#h`, half a cursor) are deterministic client
/// mistakes and return `400`; an inaccessible channel is an access-scope skip
/// that still emits nothing, matching every other read path here.
#[allow(clippy::too_many_arguments)]
async fn handle_channel_window_filter(
state: &AppState,
tenant: &buzz_core::TenantContext,
@@ -375,6 +376,7 @@ async fn handle_channel_window_filter(
accessible_channels: &[uuid::Uuid],
events: &mut Vec<Value>,
pubkey_bytes: &[u8],
reader_pubkey_hex: &str,
) -> Result<(), (StatusCode, Json<Value>)> {
use buzz_core::kind::{KIND_THREAD_SUMMARY, KIND_WINDOW_BOUNDS};
@@ -442,9 +444,18 @@ async fn handle_channel_window_filter(
.await
.map_err(|e| internal_error(&format!("channel window error: {e}")))?;
// 1. Rows, in keyset order.
// 1. Rows, in keyset order. Gate each row through the canonical read-
// authorization predicate so draft-expiry (and any future per-event
// predicates) apply consistently with every other read surface.
let mut row_ids_hex = Vec::with_capacity(window.rows.len());
for row in &window.rows {
if !buzz_core::filter::reader_can_receive_event(
&row.stored_event.event,
reader_pubkey_hex,
pubkey_bytes,
) {
continue;
}
row_ids_hex.push(row.stored_event.event.id.to_hex());
let v = serde_json::to_value(&row.stored_event.event)
.map_err(|e| internal_error(&format!("window row serialize: {e}")))?;
@@ -771,6 +782,7 @@ pub async fn query_events(
&accessible_channels,
&mut events,
&pubkey_bytes,
&authed_pubkey_hex,
)
.await?;
handled.insert(idx);
+151 -76
View File
@@ -30,8 +30,6 @@ const KIND_DRAFT: u16 = 31234;
const KIND_CREATE_CHANNEL: u16 = 9007;
const KIND_PUT_USER: u16 = 9000;
const KIND_REMOVE_USER: u16 = 9001;
/// Must match `buzz_core::kind::DRAFT_MAX_TTL_SECS`.
const DRAFT_MAX_TTL_SECS: u64 = 30 * 24 * 3600;
fn relay_url() -> String {
std::env::var("RELAY_URL").unwrap_or_else(|_| "ws://localhost:3000".to_string())
@@ -3594,60 +3592,71 @@ async fn test_reminder_target_reaction_oracle_closed() {
);
}
// ─── Read-time expiry suppression (30-day server TTL) ─────────────────────────
// ─── Read-time expiry suppression (short-lived expiration tag) ────────────────
#[tokio::test]
#[ignore]
async fn test_draft_expired_by_server_ttl_suppressed_on_http_query() {
// An expired draft (created_at > 30d ago, no expiration tag) must be absent
// from a self-authored HTTP /query while a fresh draft on the same filter is
// present. This verifies `draft_expired` is live on every read surface that
// routes through `reader_can_receive_event`.
// A draft with a short future `expiration` tag (now + 2s) must be served
// immediately after ingest, then suppressed on /query once the tag lapses.
// A fresh draft (no expiration) must remain visible throughout.
//
// Bite check: if draft_expired were removed from reader_can_receive_event, the
// expired draft would appear alongside the fresh one — the final `ids` assertion
// would fail.
// Bite check: if `draft_expired` were removed from `reader_can_receive_event`,
// the expired draft would persist alongside the fresh one — the final `ids`
// assertion would fail.
let client = http_client();
let owner = Keys::generate();
let ch_id = create_open_channel(&owner).await;
let d_expired = uuid::Uuid::new_v4().to_string();
let d_expiring = uuid::Uuid::new_v4().to_string();
let d_fresh = uuid::Uuid::new_v4().to_string();
// The relay accepts events with any created_at (no ingest freshness check).
// draft_expired then suppresses at read because created_at + 30d is in the past.
let expired = build_draft_at(
&owner,
&d_expired,
"9",
&ch_id,
&fake_nip44_v2(),
Timestamp::from(Timestamp::now().as_secs() - DRAFT_MAX_TTL_SECS - 86400),
);
let (ok_e, msg_e) = submit_event_http(&client, &owner, &expired).await;
assert!(ok_e, "expired draft must be accepted at ingest: {msg_e}");
// Short-lived draft: expiration = now + 2s. Passes ingest's strictly-future
// check, then `draft_expired` suppresses once 2s elapse.
let exp_ts = Timestamp::now().as_secs() + 2;
let expiring = EventBuilder::new(Kind::Custom(KIND_DRAFT), fake_nip44_v2())
.tags([
Tag::parse(["d", &d_expiring]).unwrap(),
Tag::parse(["k", "9"]).unwrap(),
Tag::parse(["h", &ch_id]).unwrap(),
Tag::parse(["expiration", &exp_ts.to_string()]).unwrap(),
])
.sign_with_keys(&owner)
.unwrap();
let expiring_id = expiring.id;
let (ok_e, msg_e) = submit_event_http(&client, &owner, &expiring).await;
assert!(ok_e, "expiring draft must be accepted at ingest: {msg_e}");
let fresh = build_draft(&owner, &d_fresh, "9", &ch_id, &fake_nip44_v2());
let fresh_id = fresh.id;
let (ok_f, msg_f) = submit_event_http(&client, &owner, &fresh).await;
assert!(ok_f, "fresh draft must be accepted: {msg_f}");
// Poll until the expiring draft drops from /query (timeout 10s).
let filter = Filter::new()
.kind(nostr::Kind::Custom(KIND_DRAFT))
.author(owner.public_key());
let events = query_events_http(&client, &owner.public_key().to_hex(), vec![filter]).await;
let ids: Vec<String> = events
.iter()
.filter_map(|e| e["id"].as_str().map(String::from))
.collect();
assert!(
ids.contains(&fresh_id.to_hex()),
"fresh draft must appear in self-authored query; got: {ids:?}"
);
assert!(
!ids.contains(&expired.id.to_hex()),
"expired draft (created_at > 30d ago) must be suppressed; got: {ids:?}"
);
let deadline = std::time::Instant::now() + Duration::from_secs(10);
loop {
let events =
query_events_http(&client, &owner.public_key().to_hex(), vec![filter.clone()]).await;
let ids: Vec<String> = events
.iter()
.filter_map(|e| e["id"].as_str().map(String::from))
.collect();
if !ids.contains(&expiring_id.to_hex()) {
// Suppressed — verify the fresh draft is still present.
assert!(
ids.contains(&fresh_id.to_hex()),
"fresh draft must appear in self-authored /query; got: {ids:?}"
);
break;
}
assert!(
std::time::Instant::now() < deadline,
"timed out waiting for expiring draft to be suppressed on /query"
);
tokio::time::sleep(Duration::from_millis(500)).await;
}
}
#[tokio::test]
@@ -3655,28 +3664,29 @@ async fn test_draft_expired_by_server_ttl_suppressed_on_http_query() {
async fn test_draft_expired_not_counted_on_count_surface() {
// COUNT of self-authored drafts must exclude expired drafts.
//
// This test validates the COUNT fast-path bypass: `filter_can_match_draft`
// forces the per-event fallback path (which runs `reader_can_receive_event`
// including `draft_expired`) instead of the raw SQL `count_events()` that
// cannot see per-event expiry. Without the bypass, count_events() would return
// 2 (both drafts in storage) — the assertion `count == 1` would fail, proving
// the test bites against the un-patched fast-path.
// This validates the COUNT fast-path bypass: `filter_can_match_draft` forces
// the per-event fallback path (which runs `reader_can_receive_event` including
// `draft_expired`) instead of the raw SQL `count_events()` that cannot see
// per-event expiry. Without the bypass, `count_events()` would return 2 (both
// drafts in storage) — the `count == 1` assertion proves the test bites.
let client = http_client();
let owner = Keys::generate();
let ch_id = create_open_channel(&owner).await;
let d_expired = uuid::Uuid::new_v4().to_string();
let d_expiring = uuid::Uuid::new_v4().to_string();
let d_fresh = uuid::Uuid::new_v4().to_string();
let expired = build_draft_at(
&owner,
&d_expired,
"9",
&ch_id,
&fake_nip44_v2(),
Timestamp::from(Timestamp::now().as_secs() - DRAFT_MAX_TTL_SECS - 86400),
);
let (ok_e, msg_e) = submit_event_http(&client, &owner, &expired).await;
assert!(ok_e, "expired draft must be accepted at ingest: {msg_e}");
let exp_ts = Timestamp::now().as_secs() + 2;
let expiring = EventBuilder::new(Kind::Custom(KIND_DRAFT), fake_nip44_v2())
.tags([
Tag::parse(["d", &d_expiring]).unwrap(),
Tag::parse(["k", "9"]).unwrap(),
Tag::parse(["h", &ch_id]).unwrap(),
Tag::parse(["expiration", &exp_ts.to_string()]).unwrap(),
])
.sign_with_keys(&owner)
.unwrap();
let (ok_e, msg_e) = submit_event_http(&client, &owner, &expiring).await;
assert!(ok_e, "expiring draft must be accepted at ingest: {msg_e}");
let fresh = build_draft(&owner, &d_fresh, "9", &ch_id, &fake_nip44_v2());
let (ok_f, msg_f) = submit_event_http(&client, &owner, &fresh).await;
@@ -3685,26 +3695,91 @@ async fn test_draft_expired_not_counted_on_count_surface() {
let filter = Filter::new()
.kind(nostr::Kind::Custom(KIND_DRAFT))
.author(owner.public_key());
let resp = client
.post(format!("{}/count", relay_http_url()))
.header("X-Pubkey", &owner.public_key().to_hex())
.header("Content-Type", "application/json")
.json(&vec![filter])
.send()
.await
.expect("count request");
assert!(
resp.status().is_success(),
"author COUNT must succeed, got: {}",
resp.status()
);
let body: Value = resp.json().await.expect("parse count response");
let count = body["count"].as_u64().unwrap_or(u64::MAX);
// 2 drafts in storage; 1 expired → COUNT must return 1.
// If filter_can_match_draft fast-path bypass is missing, count_events() returns
// 2 and this assertion fails — proving the test bites against the old fast-path.
assert_eq!(
count, 1,
"COUNT must return 1 (expired draft excluded); got: {count}"
);
// Poll until COUNT drops to 1 (timeout 10s).
let deadline = std::time::Instant::now() + Duration::from_secs(10);
loop {
let resp = client
.post(format!("{}/count", relay_http_url()))
.header("X-Pubkey", &owner.public_key().to_hex())
.header("Content-Type", "application/json")
.json(&vec![filter.clone()])
.send()
.await
.expect("count request");
assert!(
resp.status().is_success(),
"author COUNT must succeed, got: {}",
resp.status()
);
let body: Value = resp.json().await.expect("parse count response");
let count = body["count"].as_u64().unwrap_or(u64::MAX);
if count == 1 {
break;
}
assert!(
std::time::Instant::now() < deadline,
"timed out waiting for COUNT to exclude expired draft; last count: {count}"
);
tokio::time::sleep(Duration::from_millis(500)).await;
}
}
#[tokio::test]
#[ignore]
async fn test_draft_expired_suppressed_on_channel_window() {
// An expired draft must also be absent from the channel-window (`top_level:true`)
// bridge path. This surface previously bypassed `reader_can_receive_event` —
// this test proves the fix bites.
//
// Bite check: if the `reader_can_receive_event` gate in `handle_channel_window_filter`
// were removed, the expired draft would appear in the window response and this
// test would fail.
let client = http_client();
let owner = Keys::generate();
let ch_id = create_open_channel(&owner).await;
let d_expiring = uuid::Uuid::new_v4().to_string();
let d_fresh = uuid::Uuid::new_v4().to_string();
let exp_ts = Timestamp::now().as_secs() + 2;
let expiring = EventBuilder::new(Kind::Custom(KIND_DRAFT), fake_nip44_v2())
.tags([
Tag::parse(["d", &d_expiring]).unwrap(),
Tag::parse(["k", "9"]).unwrap(),
Tag::parse(["h", &ch_id]).unwrap(),
Tag::parse(["expiration", &exp_ts.to_string()]).unwrap(),
])
.sign_with_keys(&owner)
.unwrap();
let expiring_id = expiring.id;
let (ok_e, msg_e) = submit_event_http(&client, &owner, &expiring).await;
assert!(ok_e, "expiring draft must be accepted at ingest: {msg_e}");
let fresh = build_draft(&owner, &d_fresh, "9", &ch_id, &fake_nip44_v2());
let fresh_id = fresh.id;
let (ok_f, msg_f) = submit_event_http(&client, &owner, &fresh).await;
assert!(ok_f, "fresh draft must be accepted: {msg_f}");
// Poll channel-window until the expiring draft drops (timeout 10s).
let deadline = std::time::Instant::now() + Duration::from_secs(10);
loop {
let events = query_channel_window_mixed(&owner, &ch_id, false, false).await;
let ids: Vec<String> = events
.iter()
.filter_map(|e| e["id"].as_str().map(String::from))
.collect();
if !ids.contains(&expiring_id.to_hex()) {
// Suppressed — verify the fresh draft is still present.
assert!(
ids.contains(&fresh_id.to_hex()),
"fresh draft must appear in channel-window; got: {ids:?}"
);
break;
}
assert!(
std::time::Instant::now() < deadline,
"timed out waiting for expiring draft to be suppressed on channel-window"
);
tokio::time::sleep(Duration::from_millis(500)).await;
}
}