diff --git a/Cargo.lock b/Cargo.lock index ea301a3..75f050b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -36,6 +36,15 @@ version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "250f629c0161ad8107cf89319e990051fae62832fd343083bea452d93e2205fd" +[[package]] +name = "alloca" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e5a7d05ea6aea7e9e64d25b9156ba2fee3fdd659e34e41063cd2fc7cd020d7f4" +dependencies = [ + "cc", +] + [[package]] name = "allocator-api2" version = "0.2.21" @@ -51,6 +60,12 @@ dependencies = [ "libc", ] +[[package]] +name = "anes" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4b46cbb362ab8752921c97e041f5e366ee6297bd428a31275b9fcf1e380f7299" + [[package]] name = "anstream" version = "1.0.0" @@ -331,6 +346,7 @@ dependencies = [ "bobbin-slingshot-client", "bobbin-types", "bytes", + "criterion", "futures", "jacquard-common", "scc", @@ -533,6 +549,12 @@ dependencies = [ "serde", ] +[[package]] +name = "cast" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37b2a672a2cb129a2e41c10b1224bb368f9f37a2b16b612598138befd7b37eb5" + [[package]] name = "cbor4ii" version = "0.2.14" @@ -787,6 +809,39 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "criterion" +version = "0.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "950046b2aa2492f9a536f5f4f9a3de7b9e2476e575e05bd6c333371add4d98f3" +dependencies = [ + "alloca", + "anes", + "cast", + "ciborium", + "clap", + "criterion-plot", + "itertools 0.13.0", + "num-traits", + "oorandom", + "page_size", + "regex", + "serde", + "serde_json", + "tinytemplate", + "walkdir", +] + +[[package]] +name = "criterion-plot" +version = "0.8.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d8d80a2f4f5b554395e47b5d8305bc3d27813bacb73493eb1001e8f76dae29ea" +dependencies = [ + "cast", + "itertools 0.13.0", +] + [[package]] name = "critical-section" version = "1.2.0" @@ -1826,6 +1881,15 @@ version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" +[[package]] +name = "itertools" +version = "0.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "413ee7dfc52ee1a4949ceeb7dbc8a33f2d6c088194d9f922fb8318faf1f01186" +dependencies = [ + "either", +] + [[package]] name = "itertools" version = "0.14.0" @@ -2299,6 +2363,12 @@ version = "0.1.13" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "269bca4c2591a28585d6bf10d9ed0332b7d76900a1b02bec41bdc3a2cdcda107" +[[package]] +name = "oorandom" +version = "11.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6790f58c7ff633d8771f42965289203411a5e5c68388703c06e14f24770b41e" + [[package]] name = "openssl-probe" version = "0.2.1" @@ -2368,6 +2438,16 @@ dependencies = [ "sha2", ] +[[package]] +name = "page_size" +version = "0.6.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "30d5b2194ed13191c1999ae0704b7839fb18384fa22e49b57eeaa97d79ce40da" +dependencies = [ + "libc", + "winapi", +] + [[package]] name = "parking" version = "2.2.1" @@ -3469,7 +3549,7 @@ dependencies = [ "fnv", "fs4", "htmlescape", - "itertools", + "itertools 0.14.0", "levenshtein_automata", "log", "lru", @@ -3518,7 +3598,7 @@ checksum = "c57166f5bcfd478f370ab8445afb4678dce44801fa5ce5c451aaf8595583c5dc" dependencies = [ "downcast-rs", "fastdivide", - "itertools", + "itertools 0.14.0", "serde", "tantivy-bitpacker", "tantivy-common", @@ -3570,7 +3650,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a2cfc3ac5164cbadc28965ffb145a8f47582a60ae5897859ad8d4316596c606" dependencies = [ "futures-util", - "itertools", + "itertools 0.14.0", "tantivy-bitpacker", "tantivy-common", "tantivy-fst", @@ -3699,6 +3779,16 @@ dependencies = [ "zerovec", ] +[[package]] +name = "tinytemplate" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "be4d6b5f19ff7664e8c98d03e2139cb510db9b0a60b55f8e8709b689d939b6bc" +dependencies = [ + "serde", + "serde_json", +] + [[package]] name = "tinyvec" version = "1.11.0" diff --git a/crates/ingest/Cargo.toml b/crates/ingest/Cargo.toml index 2abb708..974afa1 100644 --- a/crates/ingest/Cargo.toml +++ b/crates/ingest/Cargo.toml @@ -30,3 +30,8 @@ tracing-subscriber = { workspace = true } tokio-util = { workspace = true } tokio = { workspace = true, features = ["test-util"] } wiremock = { workspace = true } +criterion = { version = "0.8", default-features = false, features = ["cargo_bench_support"] } + +[[bench]] +name = "json_decode" +harness = false diff --git a/crates/ingest/benches/json_decode.rs b/crates/ingest/benches/json_decode.rs new file mode 100644 index 0000000..40717ed --- /dev/null +++ b/crates/ingest/benches/json_decode.rs @@ -0,0 +1,246 @@ +use std::hint::black_box; + +use bobbin_types::edges::Record; +use criterion::{Criterion, criterion_group, criterion_main}; +use serde::Deserialize; +use serde_json::value::RawValue; + +#[derive(Deserialize)] +struct ValueFrame { + record: ValueRecordFrame, +} + +#[derive(Deserialize)] +struct ValueRecordFrame { + collection: String, + #[serde(default)] + record: Option, +} + +#[derive(Deserialize)] +struct BorrowedRawFrame<'a> { + #[serde(borrow)] + record: BorrowedRawRecordFrame<'a>, +} + +#[derive(Deserialize)] +struct BorrowedRawRecordFrame<'a> { + collection: String, + #[serde(borrow, default)] + record: Option<&'a RawValue>, +} + +#[derive(Deserialize)] +struct OwnedRawFrame { + record: OwnedRawRecordFrame, +} + +#[derive(Deserialize)] +struct OwnedRawRecordFrame { + collection: String, + #[serde(default)] + record: Option>, +} + +fn current_path(text: &str) -> Record { + let frame: ValueFrame = serde_json::from_str(text).expect("frame"); + let value = frame.record.record.expect("upsert body"); + let bytes = serde_json::to_vec(&value).expect("re-serialize Value"); + Record::from_json_bytes(frame.record.collection.as_str(), &bytes).expect("decode record") +} + +fn raw_value_borrowed(text: &str) -> Record { + let frame: BorrowedRawFrame<'_> = serde_json::from_str(text).expect("frame"); + let raw = frame.record.record.expect("upsert body"); + Record::from_json_bytes(frame.record.collection.as_str(), raw.get().as_bytes()) + .expect("decode record") +} + +fn raw_value_owned(text: &str) -> Record { + let frame: OwnedRawFrame = serde_json::from_str(text).expect("frame"); + let raw = frame.record.record.expect("upsert body"); + Record::from_json_bytes(frame.record.collection.as_str(), raw.get().as_bytes()) + .expect("decode record") +} + +fn floor_record_only(collection: &str, record_bytes: &[u8]) -> Record { + Record::from_json_bytes(collection, record_bytes).expect("decode record") +} + +fn star_text() -> &'static str { + r#"{ + "id": 1, + "type": "record", + "record": { + "live": false, + "did": "did:plc:olaren", + "rev": "3lq2zk5wqsh2k", + "collection": "sh.tangled.feed.star", + "rkey": "abcabcabcabcz", + "action": "create", + "record": { + "$type": "sh.tangled.feed.star", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "at://did:plc:abalone/sh.tangled.repo/3lq2zk5wq0000" + } + } + }"# +} + +fn repo_text() -> &'static str { + r#"{ + "id": 2, + "type": "record", + "record": { + "live": false, + "did": "did:plc:nel", + "rev": "3lq2zk5wqsh2l", + "collection": "sh.tangled.repo", + "rkey": "3lq2zk5wq0001", + "action": "create", + "record": { + "$type": "sh.tangled.repo", + "createdAt": "2026-05-01T00:00:00Z", + "knot": "knot.witchcraft.systems", + "name": "limpet", + "description": "demonstration repository for the bench corpus", + "owner": "did:plc:nel", + "repoDid": "did:plc:limpet", + "labels": [ + "at://did:plc:periwinkle/sh.tangled.label.definition/3lq2zk5wq0010", + "at://did:plc:periwinkle/sh.tangled.label.definition/3lq2zk5wq0011" + ] + } + } + }"# +} + +fn issue_text() -> &'static str { + r#"{ + "id": 3, + "type": "record", + "record": { + "live": false, + "did": "did:plc:teq", + "rev": "3lq2zk5wqsh2m", + "collection": "sh.tangled.repo.issue", + "rkey": "3lq2zk5wq0100", + "action": "create", + "record": { + "$type": "sh.tangled.repo.issue", + "createdAt": "2026-05-01T00:00:00Z", + "title": "ingest: single-pass JSON decode follow-up corpus entry", + "body": "Long body to give the bench a realistic decode cost. Repeats: blahhhhh meow meow aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + "repo": "at://did:plc:limpet", + "repoDid": "did:plc:limpet", + "mentions": [ + "did:plc:nel", + "did:plc:olaren", + "did:plc:bailey" + ], + "references": [ + "at://did:plc:limpet/sh.tangled.repo.issue/3lq2zk5wq0099", + "at://did:plc:limpet/sh.tangled.repo.pull/3lq2zk5wq0098" + ] + } + } + }"# +} + +fn pull_comment_text() -> &'static str { + r#"{ + "id": 4, + "type": "record", + "record": { + "live": false, + "did": "did:plc:lyna", + "rev": "3lq2zk5wqsh2n", + "collection": "sh.tangled.repo.pull.comment", + "rkey": "3lq2zk5wq0200", + "action": "create", + "record": { + "$type": "sh.tangled.repo.pull.comment", + "createdAt": "2026-05-01T00:00:00Z", + "body": "lgtm i thinks!!!! but please verify the cursor invariant under buffered(N) before landing. :3", + "pull": "at://did:plc:limpet/sh.tangled.repo.pull/3lq2zk5wq0098", + "owner": "did:plc:lyna" + } + } + }"# +} + +fn follow_text() -> &'static str { + r#"{ + "id": 5, + "type": "record", + "record": { + "live": false, + "did": "did:plc:bailey", + "rev": "3lq2zk5wqsh2o", + "collection": "sh.tangled.graph.follow", + "rkey": "3lq2zk5wq0300", + "action": "create", + "record": { + "$type": "sh.tangled.graph.follow", + "createdAt": "2026-05-01T00:00:00Z", + "subject": "did:plc:nel" + } + } + }"# +} + +fn extract_record_slice(text: &str) -> (String, Vec) { + let frame: ValueFrame = serde_json::from_str(text).expect("frame"); + let value = frame.record.record.expect("upsert body"); + let bytes = serde_json::to_vec(&value).expect("re-serialize Value"); + (frame.record.collection, bytes) +} + +fn bench_decode(c: &mut Criterion) { + let corpus: &[(&str, &str)] = &[ + ("star", star_text()), + ("repo", repo_text()), + ("issue", issue_text()), + ("pull_comment", pull_comment_text()), + ("follow", follow_text()), + ]; + for (name, text) in corpus { + let mut group = c.benchmark_group(format!("decode/{name}")); + let (collection, record_bytes) = extract_record_slice(text); + + group.bench_function("current_path", |b| { + b.iter(|| { + let r = current_path(black_box(text)); + black_box(r); + }); + }); + + group.bench_function("raw_value_borrowed", |b| { + b.iter(|| { + let r = raw_value_borrowed(black_box(text)); + black_box(r); + }); + }); + + group.bench_function("raw_value_owned", |b| { + b.iter(|| { + let r = raw_value_owned(black_box(text)); + black_box(r); + }); + }); + + group.bench_function("floor_record_only", |b| { + let collection = collection.as_str(); + let bytes = record_bytes.as_slice(); + b.iter(|| { + let r = floor_record_only(black_box(collection), black_box(bytes)); + black_box(r); + }); + }); + + group.finish(); + } +} + +criterion_group!(benches, bench_decode); +criterion_main!(benches); diff --git a/crates/ingest/src/frame.rs b/crates/ingest/src/frame.rs index 44a3afd..21f3993 100644 --- a/crates/ingest/src/frame.rs +++ b/crates/ingest/src/frame.rs @@ -5,8 +5,9 @@ use jacquard_common::types::recordkey::Rkey; use jacquard_common::types::string::Cid; use jacquard_common::types::tid::Tid; use serde::Deserialize; +use serde_json::value::RawValue; -#[derive(Clone, Debug, Deserialize)] +#[derive(Debug, Deserialize)] pub struct HydrantFrame { pub id: u64, #[serde(rename = "type")] @@ -45,7 +46,7 @@ pub struct HydrantStreamErrorFrame { pub message: Option, } -#[derive(Clone, Debug, Deserialize)] +#[derive(Debug, Deserialize)] pub struct RecordFrame { pub live: bool, pub did: Did, @@ -54,7 +55,7 @@ pub struct RecordFrame { pub rkey: Rkey, pub action: RecordAction, #[serde(default)] - pub record: Option, + pub record: Option>, #[serde(default)] pub cid: Option>, } diff --git a/crates/ingest/src/lib.rs b/crates/ingest/src/lib.rs index 673979b..78851b3 100644 --- a/crates/ingest/src/lib.rs +++ b/crates/ingest/src/lib.rs @@ -533,7 +533,7 @@ enum PendingOp { source: AtUri, nsid: Nsid, parsed: Record, - bytes: Vec, + bytes: Bytes, cid: Option>, edges: Vec, }, @@ -675,12 +675,11 @@ async fn prepare_record(record: Option, resolver: &RepoIdResolver) }; match record.action { RecordAction::Create | RecordAction::Update => { - let Some(value) = record.record else { + let Some(raw) = record.record else { debug!(collection = %nsid, "create/update missing record body, clearing cache"); return PendingOp::ClearCache { source }; }; - let bytes = serde_json::to_vec(&value) - .expect("serde_json::Value re-serialization is infallible"); + let bytes = Bytes::copy_from_slice(raw.get().as_bytes()); let parsed = match Record::from_json_bytes(record.collection.as_ref(), &bytes) { Ok(r) => r, Err(ExtractError::UnknownCollection(name)) => { @@ -819,7 +818,7 @@ fn cache_body( records: &dyn RecordStore, source: &AtUri, cid: Option>, - bytes: Vec, + bytes: Bytes, ) { match cid { Some(cid) => records.put( @@ -827,7 +826,7 @@ fn cache_body( Arc::new(RecordBody { uri: source.clone(), cid, - value: Bytes::from(bytes), + value: bytes, }), ), None => records.remove(source),