From d6418b8386b638de78ea62763ba866893bfc8dc0 Mon Sep 17 00:00:00 2001 From: "@permadeath.com" Date: Fri, 28 Aug 2026 14:20:37 -0400 Subject: [PATCH] test(index): record the firehose wire, and check both halves against it The producer and the consumer declare the frame shape independently, on purpose: an index that compiled against a server's types could only read servers built from this repository. Nothing held the two declarations in step except one test that linked both. `vectors/firehose/` records what a server writes. The server asserts it produces those bodies and the index asserts it reads them, neither reading the other. Two vectors exist only on the consumer side, pinning what an index must tolerate from a server older or newer than itself: absent fields it fills in, and an info name it has never heard of. Co-Authored-By: Claude Opus 5 (1M context) Change-Id: I4dadc374da0ff1b425a924a14cf3eeedf52e525f --- .../vibescrobble-index/tests/firehose_wire.rs | 90 ++++++++++++++++++ .../vibescrobble-serve/tests/firehose_wire.rs | 93 +++++++++++++++++++ vectors/firehose/README.md | 33 +++++++ vectors/firehose/commit-sparse.json | 18 ++++ vectors/firehose/commit.json | 22 +++++ vectors/firehose/info-live.json | 11 +++ vectors/firehose/info-outdated-cursor.json | 11 +++ vectors/firehose/info-unknown-name.json | 11 +++ 8 files changed, 289 insertions(+) create mode 100644 crates/vibescrobble-index/tests/firehose_wire.rs create mode 100644 crates/vibescrobble-serve/tests/firehose_wire.rs create mode 100644 vectors/firehose/README.md create mode 100644 vectors/firehose/commit-sparse.json create mode 100644 vectors/firehose/commit.json create mode 100644 vectors/firehose/info-live.json create mode 100644 vectors/firehose/info-outdated-cursor.json create mode 100644 vectors/firehose/info-unknown-name.json diff --git a/crates/vibescrobble-index/tests/firehose_wire.rs b/crates/vibescrobble-index/tests/firehose_wire.rs new file mode 100644 index 00000000..59ddd73f --- /dev/null +++ b/crates/vibescrobble-index/tests/firehose_wire.rs @@ -0,0 +1,90 @@ +//! What this index reads off a firehose, against the recorded vectors. +//! +//! The consumer's half of `vectors/firehose/`. It asserts that the recorded +//! bodies deserialize into the values this crate acts on, so a server built +//! somewhere else stays readable and a change to the types here fails against +//! the record rather than in production. +//! +//! Two of the vectors exist only on this side. No server built from this +//! repository emits them; they record what an index must tolerate from one +//! that is older or newer than itself. + +use std::path::PathBuf; + +use serde_json::Value; +use vibescrobble_index::firehose::{Commit, StreamInfo}; + +/// Reads one vector, returning its event name and body. +fn vector(name: &str) -> (String, Value) { + let path = PathBuf::from(env!("CARGO_MANIFEST_DIR")) + .join("../../vectors/firehose") + .join(format!("{name}.json")); + let text = std::fs::read_to_string(&path) + .unwrap_or_else(|err| panic!("reading {}: {err}", path.display())); + let doc: Value = serde_json::from_str(&text) + .unwrap_or_else(|err| panic!("parsing {}: {err}", path.display())); + let event = doc["event"].as_str().expect("an event name").to_owned(); + (event, doc["data"].clone()) +} + +#[test] +fn the_opening_frame_of_a_live_connection_reads() { + let (event, data) = vector("info-live"); + let info: StreamInfo = serde_json::from_value(data).expect("a readable info frame"); + + assert_eq!(event, "info"); + assert_eq!(info.name, "Live"); + assert_eq!(info.instance, "01JEXAMPLEINSTANCE"); + assert_eq!((info.oldest, info.newest), (0, 0)); +} + +#[test] +fn an_unreachable_cursor_reads_as_the_name_that_means_a_gap() { + let (_, data) = vector("info-outdated-cursor"); + let info: StreamInfo = serde_json::from_value(data).expect("a readable info frame"); + + // The one name that is not "carry on". The index compares this string, so + // the vector is what stops it drifting from what a server sends. + assert_eq!(info.name, "OutdatedCursor"); + assert_eq!(info.oldest, 41); +} + +#[test] +fn an_info_name_this_index_has_not_heard_of_is_carried_rather_than_refused() { + let (_, data) = vector("info-unknown-name"); + let info: StreamInfo = serde_json::from_value(data).expect("a readable info frame"); + + assert_eq!(info.name, "Rebalancing"); + assert_ne!(info.name, "OutdatedCursor", "this must not read as a gap"); +} + +#[test] +fn a_write_reads_with_every_field_a_server_sent() { + let (event, data) = vector("commit"); + let commit: Commit = serde_json::from_value(data).expect("a readable commit frame"); + + assert_eq!(event, "commit"); + assert_eq!(commit.seq, 42); + assert_eq!(commit.did, "did:web:kestrel.agents.localhost"); + assert_eq!(commit.collection, "zone.quernstone.scrobble"); + assert_eq!(commit.rkey, "3lqm0000abcd2"); + assert_eq!(commit.cursor(), "01JEXAMPLEINSTANCE:42"); + assert_eq!(commit.record["text"], "reading the firehose reconnect path"); +} + +#[test] +fn a_write_from_a_server_that_predates_four_fields_still_reads() { + let (_, data) = vector("commit-sparse"); + let commit: Commit = serde_json::from_value(data).expect("a readable commit frame"); + + // The point of the vector: an index reads servers it did not build, and + // one that does not name a commit must not stop it dead. Empty is the + // honest reading of "this server did not say". + assert_eq!(commit.seq, 43); + assert_eq!(commit.cid, ""); + assert_eq!(commit.commit, ""); + assert_eq!(commit.rev, ""); + assert_eq!(commit.time, ""); + assert_eq!(commit.cursor(), "01JEXAMPLEINSTANCE:43"); + assert_eq!(commit.record["text"], "a frame from an older server"); +} diff --git a/crates/vibescrobble-serve/tests/firehose_wire.rs b/crates/vibescrobble-serve/tests/firehose_wire.rs new file mode 100644 index 00000000..4c7a7d80 --- /dev/null +++ b/crates/vibescrobble-serve/tests/firehose_wire.rs @@ -0,0 +1,93 @@ +//! What this server writes on the firehose, against the recorded vectors. +//! +//! The producer's half of `vectors/firehose/`. It asserts that serializing the +//! types this crate owns yields exactly the recorded bodies, so a rename or a +//! changed field here fails before it reaches a consumer that was built +//! somewhere else. +//! +//! Two of the vectors are consumer-side only and are not read here: they +//! record what an index must tolerate from a server that is not this one, and +//! no server built from this repository emits either. + +use std::path::PathBuf; + +use serde_json::Value; +use vibescrobble_serve::{CommitFrame, Info, InfoName}; + +/// Reads one vector, returning its event name and body. +fn vector(name: &str) -> (String, Value) { + let path = PathBuf::from(env!("CARGO_MANIFEST_DIR")) + .join("../../vectors/firehose") + .join(format!("{name}.json")); + let text = std::fs::read_to_string(&path) + .unwrap_or_else(|err| panic!("reading {}: {err}", path.display())); + let doc: Value = serde_json::from_str(&text) + .unwrap_or_else(|err| panic!("parsing {}: {err}", path.display())); + let event = doc["event"].as_str().expect("an event name").to_owned(); + (event, doc["data"].clone()) +} + +#[test] +fn the_opening_frame_of_a_live_connection_is_what_was_recorded() { + let (event, data) = vector("info-live"); + let info = Info { + name: InfoName::Live, + instance: "01JEXAMPLEINSTANCE".to_owned(), + oldest: 0, + newest: 0, + message: "streaming from now; nothing was replayed".to_owned(), + }; + + assert_eq!(event, "info"); + assert_eq!(serde_json::to_value(&info).expect("serializable"), data); +} + +#[test] +fn an_unreachable_cursor_is_reported_as_what_was_recorded() { + let (event, data) = vector("info-outdated-cursor"); + let info = Info { + name: InfoName::OutdatedCursor, + instance: "01JEXAMPLEINSTANCE".to_owned(), + oldest: 41, + newest: 96, + message: "cursor 01JOLDINSTANCE:7 is older than the buffer; 41 is the oldest still held" + .to_owned(), + }; + + assert_eq!(event, "info"); + assert_eq!(serde_json::to_value(&info).expect("serializable"), data); +} + +#[test] +fn a_write_is_framed_as_what_was_recorded() { + let (event, data) = vector("commit"); + let frame = CommitFrame { + seq: 42, + instance: "01JEXAMPLEINSTANCE".to_owned(), + did: "did:web:kestrel.agents.localhost".to_owned(), + collection: "zone.quernstone.scrobble".to_owned(), + rkey: "3lqm0000abcd2".to_owned(), + uri: "at://did:web:kestrel.agents.localhost/zone.quernstone.scrobble/3lqm0000abcd2" + .to_owned(), + cid: "bafyreigh2akiscaildcqabsyg3dfr6chu3fgpregiymsck7e7aqa4s52zy".to_owned(), + commit: "bafyreidfuxdmqvxvtqcyvzkxpvbdcvwjxlrqvbn5xk3vg6mfumzsrxvz4a".to_owned(), + rev: "3lqm0000abce4".to_owned(), + time: "2026-08-24T12:00:00Z".to_owned(), + record: serde_json::json!({ + "$type": "zone.quernstone.scrobble", + "text": "reading the firehose reconnect path", + "emoji": "🧵", + "createdAt": "2026-08-24T12:00:00Z" + }), + }; + + assert_eq!(event, "commit"); + assert_eq!(serde_json::to_value(&frame).expect("serializable"), data); +} + +#[test] +fn the_cursor_a_frame_hands_back_is_the_instance_and_the_sequence() { + let (_, data) = vector("commit"); + let frame: CommitFrame = serde_json::from_value(data).expect("a readable frame"); + assert_eq!(frame.cursor(), "01JEXAMPLEINSTANCE:42"); +} diff --git a/vectors/firehose/README.md b/vectors/firehose/README.md new file mode 100644 index 00000000..14accf08 --- /dev/null +++ b/vectors/firehose/README.md @@ -0,0 +1,33 @@ +# Firehose wire vectors + +What a server writes on `/firehose`, recorded so that a producer and a consumer +built from different repositories can each check themselves against it without +either compiling against the other. + +The frame shape is not owned by either side. A server declares it as +`vibescrobble_serve::firehose::{CommitFrame, Info}`, an index declares it as +`vibescrobble_index::firehose::{Commit, StreamInfo}`, and the two declarations +are deliberately not the same type: an index that compiled against a server's +types could only read servers built from this repository. These files are what +holds the two in step instead. + +Each file is one frame: + +```json +{ "note": "why this case is here", "event": "commit", "data": { … } } +``` + +`event` is the SSE event name. `data` is the JSON body of that event. + +A producer asserts that serializing its own type yields `data`. A consumer +asserts that deserializing `data` yields the values named beside it. Neither +side reads the other's assertions. + +Two of these exist to pin a tolerance rather than a shape: + +- `info-unknown-name.json` — a consumer carries an `info` name it has never + heard of rather than refusing the stream, so a server that grows a fourth + thing to say does not stop an older index. +- `commit-sparse.json` — a consumer fills in the fields a server did not send, + so an index can read a server that predates them. No server built from this + repository emits this frame; it is a consumer-side vector only. diff --git a/vectors/firehose/commit-sparse.json b/vectors/firehose/commit-sparse.json new file mode 100644 index 00000000..e5cc2a0c --- /dev/null +++ b/vectors/firehose/commit-sparse.json @@ -0,0 +1,18 @@ +{ + "note": "Consumer-side only. A server that predates cid, commit, rev and time sends none of them; an index reads servers it did not build, so it fills them in rather than stopping.", + "event": "commit", + "data": { + "seq": 43, + "instance": "01JEXAMPLEINSTANCE", + "did": "did:web:kestrel.agents.localhost", + "collection": "zone.quernstone.scrobble", + "rkey": "3lqm0000abcf6", + "uri": "at://did:web:kestrel.agents.localhost/zone.quernstone.scrobble/3lqm0000abcf6", + "record": { + "$type": "zone.quernstone.scrobble", + "text": "a frame from an older server", + "emoji": "🕰️", + "createdAt": "2026-08-24T12:00:01Z" + } + } +} diff --git a/vectors/firehose/commit.json b/vectors/firehose/commit.json new file mode 100644 index 00000000..d9208bd6 --- /dev/null +++ b/vectors/firehose/commit.json @@ -0,0 +1,22 @@ +{ + "note": "One write, with every field a server built from this repository sends.", + "event": "commit", + "data": { + "seq": 42, + "instance": "01JEXAMPLEINSTANCE", + "did": "did:web:kestrel.agents.localhost", + "collection": "zone.quernstone.scrobble", + "rkey": "3lqm0000abcd2", + "uri": "at://did:web:kestrel.agents.localhost/zone.quernstone.scrobble/3lqm0000abcd2", + "cid": "bafyreigh2akiscaildcqabsyg3dfr6chu3fgpregiymsck7e7aqa4s52zy", + "commit": "bafyreidfuxdmqvxvtqcyvzkxpvbdcvwjxlrqvbn5xk3vg6mfumzsrxvz4a", + "rev": "3lqm0000abce4", + "time": "2026-08-24T12:00:00Z", + "record": { + "$type": "zone.quernstone.scrobble", + "text": "reading the firehose reconnect path", + "emoji": "🧵", + "createdAt": "2026-08-24T12:00:00Z" + } + } +} diff --git a/vectors/firehose/info-live.json b/vectors/firehose/info-live.json new file mode 100644 index 00000000..a7e836b0 --- /dev/null +++ b/vectors/firehose/info-live.json @@ -0,0 +1,11 @@ +{ + "note": "The opening frame on a connection that gave no cursor. Sent unconditionally, so a consumer learns the instance and the buffer's range without provoking an error.", + "event": "info", + "data": { + "name": "Live", + "instance": "01JEXAMPLEINSTANCE", + "oldest": 0, + "newest": 0, + "message": "streaming from now; nothing was replayed" + } +} diff --git a/vectors/firehose/info-outdated-cursor.json b/vectors/firehose/info-outdated-cursor.json new file mode 100644 index 00000000..feb7969e --- /dev/null +++ b/vectors/firehose/info-outdated-cursor.json @@ -0,0 +1,11 @@ +{ + "note": "The cursor named a frame this server no longer holds, or one from a previous run. Everything still held follows, and the consumer has a hole only a re-read of the repositories can fill.", + "event": "info", + "data": { + "name": "OutdatedCursor", + "instance": "01JEXAMPLEINSTANCE", + "oldest": 41, + "newest": 96, + "message": "cursor 01JOLDINSTANCE:7 is older than the buffer; 41 is the oldest still held" + } +} diff --git a/vectors/firehose/info-unknown-name.json b/vectors/firehose/info-unknown-name.json new file mode 100644 index 00000000..4a8db568 --- /dev/null +++ b/vectors/firehose/info-unknown-name.json @@ -0,0 +1,11 @@ +{ + "note": "Consumer-side only. Every info name except OutdatedCursor means carry on, including one this index has never heard of.", + "event": "info", + "data": { + "name": "Rebalancing", + "instance": "01JEXAMPLEINSTANCE", + "oldest": 41, + "newest": 96, + "message": "moving the buffer; the stream continues" + } +} -- 2.51.2