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" + } +}