feat(cli): wire the composite pagination cursors the relay already reads

`buzz messages thread` and `buzz messages get` are both capped (500 and
200) with no way to page past the ceiling, so any thread or channel
larger than its cap is unreadable from the CLI. Both relay-side cursors
already exist and are fully implemented; neither flag was ever added.

## messages thread: --after / --after-id

`get_thread_replies` takes a composite `(created_at, event_id)` keyset
cursor (buzz-db/src/thread.rs:334-343) and the bridge parses it from
`thread_cursor`/`thread_cursor_id` (bridge.rs:1183) — the CLI never sent
the field, so a 265-reply thread ends at the 500-reply wall with no
escape. The wire names are unchanged; the flags are spelled `--after`
because this cursor pages FORWARD (`event_created_at > $n`), so the
house `--before` spelling from `social notes` would name the wrong
direction. Shape and both-or-neither grammar follow that precedent.

The relay reads the cursor only on the depth-limited code path — a
filter without `depth_limit` routes to the generic catch-all, which has
no cursor. Rather than touch that routing guard (removing it would
silently re-anchor every existing caller), a cursor supplied without
`--depth-limit` sends a depth sentinel meaning "no effective bound".
The sentinel is `i32::MAX`, not `u32::MAX`: buzz-db binds depth as i32
(thread.rs:459), so anything above `i32::MAX` wraps negative and matches
zero replies. Measured live: `--depth-limit 2147483648` returns the root
alone, `2147483647` returns the full thread.

## messages get: --before-id

`--before` maps to `until`, which is INCLUSIVE (`created_at <= $u`,
event.rs:505). `before_id` upgrades it to the exclusive keyset branch
(event.rs:493-499) and the bridge has parsed it, with a both-or-neither
grammar and a BAD_REQUEST on malformed input, since bridge.rs:263 — the
doc comment at :424 names the dense-second bug it exists to kill. The
flag did not exist.

## Why the tiebreak is not optional

The two paths have OPPOSITE cursor polarity, so they fail differently
and no timestamp-only rule covers both:

  backward (get)    `created_at <= $u`  INCLUSIVE  event.rs:505
  forward  (thread) `event_created_at > $n` STRICT  thread.rs:443

Measured at this base. Backward, a cap binding inside a second shared by
several messages stalls or drops depending on loop shape. Forward, a page
boundary landing inside a shared second skips the rest of that second
outright: on a 44-reply thread with one 2-multiplicity second, the
timestamp-only walk drops a reply at caps {1, 2, 11, 22} and succeeds at
the other twenty — non-monotonic, because the hazard is where the
boundary lands, not the cap's size. Both failures are silent, rc=0.

With `(ts, id)` the tiebreak travels in the cursor: recovery no longer
depends on loop shape at all. That is what these flags buy — not just a
higher ceiling.

## Tests and help

15 new unit tests over two extracted pure filter builders, covering the
sentinel, its i32 representability, explicit-depth precedence, the
zero-cursor seed, and both-or-neither on each surface. Mutation-checked:
dropping the `before_id` write kills two of them and degrades the live
composite arm to the same drop as timestamp-only paging.

`--help` on both commands now states the default and cap, which end of
the range you get, that the thread cursor is strictly-after and needs
`--after 0` to seed a forward walk, and why to always pass the tiebreak.

Co-authored-by: Sami <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@buzz.block.builderlab.xyz>
Signed-off-by: Sami <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@buzz.block.builderlab.xyz>
This commit is contained in:
Sami
2026-08-15 20:52:31 -04:00
parent 78cbffeb64
commit 8aad722246
2 changed files with 362 additions and 28 deletions
+339 -23
View File
@@ -350,18 +350,38 @@ fn format_events(normalized: &str, format: &crate::OutputFormat) -> String {
}
}
pub async fn cmd_get_messages(
client: &BuzzClient,
/// Composite cursor grammar: the event-id half is meaningless without the
/// timestamp half. The relay rejects it outright (`before_id requires until to
/// be set`, `bridge.rs`), and a half cursor silently demoted to a head request
/// is the dup/loss bug the composite form exists to kill — so refuse locally
/// rather than spend a round trip on a 400.
fn validate_cursor_pair(
ts: Option<i64>,
id: Option<&str>,
id_flag: &str,
ts_flag: &str,
) -> Result<(), CliError> {
if let Some(id) = id {
if ts.is_none() {
return Err(CliError::Usage(format!("{id_flag} requires {ts_flag}")));
}
validate_hex64(id)?;
}
Ok(())
}
/// Build the filter for a channel message query.
///
/// Split out from [`cmd_get_messages`] so the cursor grammar is testable
/// without a live relay.
fn build_messages_filter(
channel_id: &str,
limit: Option<u32>,
limit: u32,
before: Option<i64>,
before_id: Option<&str>,
since: Option<i64>,
kinds: Option<&str>,
format: &crate::OutputFormat,
) -> Result<(), CliError> {
validate_uuid(channel_id)?;
let limit = limit.unwrap_or(50).min(200);
) -> serde_json::Value {
let mut filter = serde_json::json!({
"kinds": [9, 40002, 40008, 45001, 45003],
"#h": [channel_id],
@@ -378,11 +398,35 @@ pub async fn cmd_get_messages(
if let Some(b) = before {
filter["until"] = serde_json::json!(b);
// Both or neither: the id half only rides along with the timestamp.
if let Some(bid) = before_id {
filter["before_id"] = serde_json::json!(bid);
}
}
if let Some(s) = since {
filter["since"] = serde_json::json!(s);
}
filter
}
#[allow(clippy::too_many_arguments)]
pub async fn cmd_get_messages(
client: &BuzzClient,
channel_id: &str,
limit: Option<u32>,
before: Option<i64>,
before_id: Option<&str>,
since: Option<i64>,
kinds: Option<&str>,
format: &crate::OutputFormat,
) -> Result<(), CliError> {
validate_uuid(channel_id)?;
validate_cursor_pair(before, before_id, "--before-id", "--before")?;
let limit = limit.unwrap_or(50).min(200);
let filter = build_messages_filter(channel_id, limit, before, before_id, since, kinds);
let resp = client.query(&filter).await?;
let mut events: Vec<serde_json::Value> = serde_json::from_str(&resp).unwrap_or_default();
events.sort_by_key(|e| e.get("created_at").and_then(|v| v.as_u64()).unwrap_or(0));
@@ -391,30 +435,82 @@ pub async fn cmd_get_messages(
Ok(())
}
pub async fn cmd_get_thread(
client: &BuzzClient,
/// Depth bound sent when `--thread-cursor` is used without `--depth-limit`.
///
/// The relay only reads the thread cursor on the depth-limited code path
/// (`bridge.rs`: filters without `depth_limit` fall through to the generic
/// catch-all query, which has no cursor), so reaching the cursor at all
/// requires sending *some* depth bound. This value expresses "no effective
/// bound": it is `i32::MAX`, and thread nesting cannot approach it.
///
/// It is deliberately `i32::MAX` rather than `u32::MAX` — the relay binds the
/// depth as `i32`, so any value above `i32::MAX` wraps negative and matches
/// zero rows.
const THREAD_CURSOR_DEPTH_SENTINEL: u32 = i32::MAX as u32;
/// Build the reply filter for a thread query.
///
/// Split out from [`cmd_get_thread`] so the cursor/depth interaction is
/// testable without a live relay.
fn build_thread_reply_filter(
channel_id: &str,
event_id: &str,
limit: Option<u32>,
limit: u32,
depth_limit: Option<u32>,
format: &crate::OutputFormat,
) -> Result<(), CliError> {
validate_uuid(channel_id)?;
validate_hex64(event_id)?;
let limit = limit.unwrap_or(100).min(500);
// Two filters ORed in a single HTTP call:
// 1. Replies referencing this event via e-tag (no kind restriction)
// 2. The root event itself by ID
thread_cursor: Option<i64>,
thread_cursor_id: Option<&str>,
) -> serde_json::Value {
let mut reply_filter = serde_json::json!({
"kinds": [9, 40002, 40003, 40008, 45003],
"#h": [channel_id],
"#e": [event_id],
"limit": limit
});
if let Some(d) = depth_limit {
reply_filter["depth_limit"] = serde_json::json!(d);
// A cursor is only honoured on the depth-limited path, so an explicit
// cursor implies a depth bound even when the caller gave none.
match (depth_limit, thread_cursor) {
(Some(d), _) => reply_filter["depth_limit"] = serde_json::json!(d),
(None, Some(_)) => {
reply_filter["depth_limit"] = serde_json::json!(THREAD_CURSOR_DEPTH_SENTINEL)
}
(None, None) => {}
}
if let Some(c) = thread_cursor {
reply_filter["thread_cursor"] = serde_json::json!(c);
if let Some(id) = thread_cursor_id {
reply_filter["thread_cursor_id"] = serde_json::json!(id);
}
}
reply_filter
}
#[allow(clippy::too_many_arguments)]
pub async fn cmd_get_thread(
client: &BuzzClient,
channel_id: &str,
event_id: &str,
limit: Option<u32>,
depth_limit: Option<u32>,
thread_cursor: Option<i64>,
thread_cursor_id: Option<&str>,
format: &crate::OutputFormat,
) -> Result<(), CliError> {
validate_uuid(channel_id)?;
validate_hex64(event_id)?;
validate_cursor_pair(thread_cursor, thread_cursor_id, "--after-id", "--after")?;
let limit = limit.unwrap_or(100).min(500);
// Two filters ORed in a single HTTP call:
// 1. Replies referencing this event via e-tag (no kind restriction)
// 2. The root event itself by ID
let reply_filter = build_thread_reply_filter(
channel_id,
event_id,
limit,
depth_limit,
thread_cursor,
thread_cursor_id,
);
let root_filter = serde_json::json!({
"ids": [event_id],
"limit": 1
@@ -948,6 +1044,7 @@ pub async fn dispatch(
channel,
limit,
before,
before_id,
since,
kinds,
} => {
@@ -956,6 +1053,7 @@ pub async fn dispatch(
&channel,
limit,
before,
before_id.as_deref(),
since,
kinds.as_deref(),
format,
@@ -967,7 +1065,21 @@ pub async fn dispatch(
event,
limit,
depth_limit,
} => cmd_get_thread(client, &channel, &event, limit, depth_limit, format).await,
after,
after_id,
} => {
cmd_get_thread(
client,
&channel,
&event,
limit,
depth_limit,
after,
after_id.as_deref(),
format,
)
.await
}
MessagesCmd::Search {
query,
author,
@@ -1373,3 +1485,207 @@ mod tests {
assert_eq!(match_profiles_by_name(&events, "Aaron").len(), 1);
}
}
#[cfg(test)]
mod thread_cursor_tests {
use super::{build_thread_reply_filter, THREAD_CURSOR_DEPTH_SENTINEL};
const CH: &str = "3928fe05-df61-4b5d-b9c7-d623b9b10ea1";
const ROOT: &str = "f6f7a5212b1a6451f1906406e224c01834dc950826c337046b74b18ecc5785ce";
const CURSOR_ID: &str = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa";
#[test]
fn no_cursor_and_no_depth_sends_neither_field() {
// The default pull must keep its existing shape: absent `depth_limit`
// is what routes the filter to the catch-all (newest-anchored) path,
// so adding the cursor flags must not perturb it.
let f = build_thread_reply_filter(CH, ROOT, 100, None, None, None);
assert!(f.get("depth_limit").is_none());
assert!(f.get("thread_cursor").is_none());
assert!(f.get("thread_cursor_id").is_none());
}
#[test]
fn cursor_without_depth_limit_supplies_the_sentinel() {
// The relay only reads the cursor on the depth-limited path, so a
// cursor with no explicit depth must still carry a depth bound or the
// cursor is silently ignored and the caller re-reads page one forever.
let f = build_thread_reply_filter(CH, ROOT, 500, None, Some(1_786_800_000), None);
assert_eq!(
f["depth_limit"],
serde_json::json!(THREAD_CURSOR_DEPTH_SENTINEL)
);
assert_eq!(f["thread_cursor"], serde_json::json!(1_786_800_000_i64));
}
#[test]
fn sentinel_is_representable_as_i32() {
// buzz-db binds the depth as i32 (`dl as i32`); any value above
// i32::MAX wraps negative and `depth <= -N` matches zero replies.
// Measured live at bff3110a0: --depth-limit 2147483648 returns n=1
// (root only) while 2147483647 returns the full thread.
assert!(i32::try_from(THREAD_CURSOR_DEPTH_SENTINEL).is_ok());
assert_eq!(THREAD_CURSOR_DEPTH_SENTINEL, i32::MAX as u32);
}
#[test]
fn explicit_depth_limit_wins_over_the_sentinel() {
// A caller who asks for depth 2 while paging must get depth 2, not the
// sentinel — otherwise the cursor plumbing would silently widen an
// explicit depth bound.
let f = build_thread_reply_filter(CH, ROOT, 500, Some(2), Some(1_786_800_000), None);
assert_eq!(f["depth_limit"], serde_json::json!(2));
}
#[test]
fn composite_cursor_carries_the_tiebreak_id() {
let f =
build_thread_reply_filter(CH, ROOT, 500, None, Some(1_786_800_000), Some(CURSOR_ID));
assert_eq!(f["thread_cursor"], serde_json::json!(1_786_800_000_i64));
assert_eq!(f["thread_cursor_id"], serde_json::json!(CURSOR_ID));
}
#[test]
fn cursor_id_is_dropped_without_a_cursor_timestamp() {
// Defense in depth: cmd_get_thread rejects this combination up front,
// but the builder must not emit a lone `thread_cursor_id` either —
// the relay's cursor grammar requires both or neither.
let f = build_thread_reply_filter(CH, ROOT, 500, None, None, Some(CURSOR_ID));
assert!(f.get("thread_cursor_id").is_none());
assert!(f.get("thread_cursor").is_none());
}
#[test]
fn limit_and_targeting_fields_are_unchanged_by_paging() {
let f = build_thread_reply_filter(CH, ROOT, 500, None, Some(1), None);
assert_eq!(f["limit"], serde_json::json!(500));
assert_eq!(f["#h"], serde_json::json!([CH]));
assert_eq!(f["#e"], serde_json::json!([ROOT]));
}
#[test]
fn a_zero_cursor_is_a_real_cursor_and_seeds_the_forward_walk() {
// `--after 0` is how a caller opts into the oldest-anchored walk
// without also constraining depth. `Some(0)` must therefore be treated
// as present, not folded into `None` by a falsy check — otherwise the
// filter routes to the newest-anchored catch-all and the walk
// terminates after one page.
let f = build_thread_reply_filter(CH, ROOT, 500, None, Some(0), None);
assert_eq!(f["thread_cursor"], serde_json::json!(0_i64));
assert_eq!(
f["depth_limit"],
serde_json::json!(THREAD_CURSOR_DEPTH_SENTINEL)
);
}
#[test]
fn the_tiebreak_id_is_sent_whenever_it_is_supplied() {
// Polarity asymmetry, measured at bff3110a0: the forward (thread)
// legacy cursor is STRICT `>` (buzz-db/src/thread.rs:443) while the
// backward (messages get) one is INCLUSIVE `<=` (event.rs:505). So on
// this path a timestamp-only cursor whose page boundary lands inside a
// shared second skips the remainder of that second silently — there is
// no inclusive re-return to recover it, and no timestamp-only step rule
// is safe at every boundary alignment. Live: thread 7f2cea28, tie
// second 1785163671 (multiplicity 2 by pinned probe), ts-only drops at
// caps {1,2,11,22} while the composite cursor recovers 2/2 at all of
// 1..24. The tiebreak must therefore ride along on every paged call.
let f =
build_thread_reply_filter(CH, ROOT, 500, None, Some(1_785_163_671), Some(CURSOR_ID));
assert_eq!(f["thread_cursor_id"], serde_json::json!(CURSOR_ID));
}
}
#[cfg(test)]
mod messages_cursor_tests {
use super::{build_messages_filter, validate_cursor_pair};
const CH: &str = "3928fe05-df61-4b5d-b9c7-d623b9b10ea1";
const BEFORE_ID: &str = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb";
#[test]
fn head_request_sends_no_cursor_fields() {
// The default pull must keep its existing wire shape.
let f = build_messages_filter(CH, 50, None, None, None, None);
assert!(f.get("until").is_none());
assert!(f.get("before_id").is_none());
assert_eq!(f["limit"], serde_json::json!(50));
assert_eq!(f["#h"], serde_json::json!([CH]));
}
#[test]
fn timestamp_only_cursor_still_works() {
// Back-compat: `--before` alone is the existing (inclusive) cursor and
// must keep sending a bare `until`, never a half composite.
let f = build_messages_filter(CH, 200, Some(1_786_800_000), None, None, None);
assert_eq!(f["until"], serde_json::json!(1_786_800_000_i64));
assert!(f.get("before_id").is_none());
}
#[test]
fn composite_cursor_sends_both_halves() {
let f = build_messages_filter(CH, 200, Some(1_786_800_000), Some(BEFORE_ID), None, None);
assert_eq!(f["until"], serde_json::json!(1_786_800_000_i64));
assert_eq!(f["before_id"], serde_json::json!(BEFORE_ID));
}
#[test]
fn cursor_id_is_dropped_without_a_timestamp() {
// The relay 400s on `before_id` without `until`; the builder must not
// emit a half cursor even if the caller-level guard is bypassed.
let f = build_messages_filter(CH, 200, None, Some(BEFORE_ID), None, None);
assert!(f.get("before_id").is_none());
assert!(f.get("until").is_none());
}
#[test]
fn since_and_kinds_survive_a_composite_cursor() {
let f = build_messages_filter(
CH,
200,
Some(1_786_800_000),
Some(BEFORE_ID),
Some(1_786_000_000),
Some("9,1984"),
);
assert_eq!(f["since"], serde_json::json!(1_786_000_000_i64));
assert_eq!(f["kinds"], serde_json::json!([9, 1984]));
assert_eq!(f["before_id"], serde_json::json!(BEFORE_ID));
}
#[test]
fn lone_cursor_id_is_a_usage_error() {
let err = validate_cursor_pair(None, Some(BEFORE_ID), "--before-id", "--before")
.expect_err("lone --before-id must be refused");
assert!(
err.to_string().contains("--before-id requires --before"),
"unexpected error: {err}"
);
}
#[test]
fn malformed_cursor_id_is_rejected_locally() {
// The relay rejects a non-64-hex `before_id` with a 400; refusing it
// here means the caller never spends a round trip to learn that.
assert!(validate_cursor_pair(
Some(1_786_800_000),
Some("not-a-hex-id"),
"--before-id",
"--before"
)
.is_err());
}
#[test]
fn a_full_composite_cursor_validates() {
assert!(validate_cursor_pair(
Some(1_786_800_000),
Some(BEFORE_ID),
"--before-id",
"--before"
)
.is_ok());
// And a bare timestamp cursor is legal on both surfaces.
assert!(validate_cursor_pair(Some(1_786_800_000), None, "--before-id", "--before").is_ok());
}
}
+23 -5
View File
@@ -461,18 +461,23 @@ pub enum MessagesCmd {
},
/// Retrieve messages from a channel
#[command(
after_help = "Examples:\n buzz messages get --channel <UUID>\n buzz messages get --channel <UUID> --limit 50 --kinds 1,1984"
after_help = "Pagination:\n Returns up to --limit messages (default 50, max 200), NEWEST-first.\n For channels larger than the cap, page backwards with --before:\n\n buzz messages get --channel <UUID> --limit 200\n buzz messages get --channel <UUID> --limit 200 \\\n --before <created_at of oldest message seen> --before-id <its event id>\n\n --before alone is INCLUSIVE (<=), so a timestamp-only cursor re-returns\n every message sharing that second; if one second holds more messages than\n --limit, paging stalls. Pass --before-id to make the cursor exclusive.\n\nExamples:\n buzz messages get --channel <UUID>\n buzz messages get --channel <UUID> --limit 50 --kinds 1,1984"
)]
Get {
/// Channel UUID
#[arg(long)]
channel: String,
/// Maximum number of results to return
/// Maximum number of results to return (default 50, max 200)
#[arg(long)]
limit: Option<u32>,
/// Unix timestamp — return messages before this time
/// Unix timestamp — return messages before this time (inclusive
/// unless paired with --before-id)
#[arg(long)]
before: Option<i64>,
/// Event ID cursor — tiebreak for messages sharing the cursor second
/// (composite pagination with --before)
#[arg(long)]
before_id: Option<String>,
/// Unix timestamp — return messages after this time
#[arg(long)]
since: Option<i64>,
@@ -481,6 +486,9 @@ pub enum MessagesCmd {
kinds: Option<String>,
},
/// Get a message thread (replies to a root message)
#[command(
after_help = "Pagination:\n Returns up to --limit replies (default 100, max 500) plus the root event.\n Replies without a cursor and without --depth-limit are NEWEST-first; either\n one selects the OLDEST-first walk, so the cursor picks which end you see.\n\n To page a thread larger than the cap, walk FORWARD from the oldest reply.\n Seed with --after 0, then pass the newest reply of each page back in:\n\n buzz messages thread --channel <UUID> --event <ID> --limit 500 --after 0\n buzz messages thread --channel <UUID> --event <ID> --limit 500 \\\n --after <created_at of newest reply seen> --after-id <its event id>\n\n Always pass --after-id. The timestamp-only cursor is STRICTLY greater-than,\n so a page boundary landing inside a second shared by several replies skips\n the rest of that second silently (rc=0). --after-id carries the tiebreak, so\n no reply is skipped regardless of where the boundary lands."
)]
Thread {
/// Channel UUID
#[arg(long)]
@@ -488,12 +496,22 @@ pub enum MessagesCmd {
/// Root message event ID (64-char hex)
#[arg(long)]
event: String,
/// Maximum number of results to return
/// Maximum number of replies to return (default 100, max 500)
#[arg(long)]
limit: Option<u32>,
/// Maximum reply nesting depth to include
/// Maximum reply nesting depth to include. Also selects the
/// oldest-first reply ordering (see Pagination below)
#[arg(long)]
depth_limit: Option<u32>,
/// Unix timestamp cursor — return replies created after this time
/// (STRICTLY after; pass --after 0 to seed a forward walk)
#[arg(long)]
after: Option<i64>,
/// Event ID cursor — tiebreak for replies sharing the cursor second
/// (composite pagination with --after). Always pass this when paging:
/// without it a page boundary inside a shared second skips replies
#[arg(long)]
after_id: Option<String>,
},
/// Full-text search across messages
#[command(