From dec1964fc22e43ac1f1dbca238fea3cfb4f7b2dc Mon Sep 17 00:00:00 2001 From: Timothy Quilling Date: Mon, 13 Oct 2025 12:18:42 -0400 Subject: [PATCH] chore: cleanup, reorg processing payloads to the whole event --- consumer/src/indexer/mod.rs | 116 ++++++++--------------------- consumer/src/jetstream/consumer.rs | 26 +++---- consumer/src/jetstream/types.rs | 46 ++++-------- 3 files changed, 53 insertions(+), 135 deletions(-) diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs index 700316d3..8b6ca35a 100644 --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -52,18 +52,6 @@ struct RelayIndexerState { last_seq_time: Option<(u64, chrono::DateTime)>, } -/// Extract time_us sequence number from a JSON preview string -fn extract_sequence_from_preview(preview: &str) -> Option { - // Try to parse the preview as JSON - if let Ok(json) = serde_json::from_str::(preview) { - // Extract time_us field if it exists - if let Some(time_us) = json.get("time_us").and_then(|v| v.as_u64()) { - return Some(time_us); - } - } - None -} - pub struct RelayIndexer { pool: Pool, redis: MultiplexedConnection, @@ -80,35 +68,15 @@ impl RelayIndexer { /// Convert a Jetstream commit event to a Firehose commit event async fn convert_jetstream_commit( &self, - did: &str, - commit: &crate::jetstream::CommitPayload, + commit_event: &crate::jetstream::CommitEvent, ) -> Option { use crate::firehose::{AtpCommitEvent, CommitOp}; use crate::jetstream::CommitOperation; use ipld_core::cid::Cid; use serde_bytes::ByteBuf; - // Jetstream events end up looking something like: - // { - // "did": "did:plc:eygmaihciaxprqvxpfvl6flk", - // "time_us": 1725911162329308, - // "type": "com", - // "commit": { - // "rev": "3l3qo2vutsw2b", - // "type": "c", - // "collection": "app.bsky.feed.like", - // "rkey": "3l3qo2vuowo2b", - // "record": { - // "$type": "app.bsky.feed.like", - // "createdAt": "2024-09-09T19:46:02.102Z", - // "subject": { - // "cid": "bafyreidc6sydkkbchcyg62v77wbhzvb2mvytlmsychqgwf2xojjtirmzj4", - // "uri": "at://did:plc:wa7b35aakoll7hugkrjtf3xf/app.bsky.feed.post/3l3pte3p2e325" - // } - // }, - // "cid": "bafyreidwaivazkwu67xztlmuobx35hs2lnfh3kolmgfmucldvhd3sgzcqi" - // } - // } + let did = &commit_event.did; + let commit = &commit_event.commit; // Handle different operation types with appropriate CID handling let commit_cid = if commit.op == CommitOperation::Delete { @@ -130,9 +98,10 @@ impl RelayIndexer { } }; - // Use current time for timestamp since Jetstream events don't include - // an ISO 8601 timestamp in commit events (only in identity/account events) - let timestamp = chrono::Utc::now(); + // Convert microsecond timestamp to DateTime + let timestamp = + chrono::DateTime::::from_timestamp_micros(commit_event.time_us as i64) + .unwrap_or_else(|| chrono::Utc::now()); // Map Jetstream operation to Firehose operation let op = match commit.op { @@ -271,11 +240,13 @@ impl RelayIndexer { /// Convert a Jetstream identity event to a Firehose identity event fn convert_jetstream_identity( &self, - did: &str, - identity: &crate::jetstream::IdentityPayload, + identity_event: &crate::jetstream::IdentityEvent, ) -> FirehoseEvent { use crate::firehose::AtpIdentityEvent; + let did = &identity_event.did; + let identity = &identity_event.identity; + // Parse the timestamp from ISO 8601 format let timestamp = match chrono::DateTime::parse_from_rfc3339(&identity.time) { Ok(dt) => dt.with_timezone(&chrono::Utc), @@ -298,11 +269,13 @@ impl RelayIndexer { /// Convert a Jetstream account event to a Firehose account event fn convert_jetstream_account( &self, - did: &str, - account: &crate::jetstream::AccountPayload, + account_event: &crate::jetstream::AccountEvent, ) -> FirehoseEvent { use crate::firehose::AtpAccountEvent; + let did = &account_event.did; + let account = &account_event.account; + // Parse the timestamp from ISO 8601 format let timestamp = match chrono::DateTime::parse_from_rfc3339(&account.time) { Ok(dt) => dt.with_timezone(&chrono::Utc), @@ -635,31 +608,22 @@ impl RelayIndexer { tracing::error!("Unrecoverable Jetstream error, exiting"); break; }, - crate::jetstream::JetstreamOutput::Commit(did, commit) => { + crate::jetstream::JetstreamOutput::Commit(commit_event) => { // Map Jetstream commit to Firehose commit event - let start = std::time::Instant::now(); - let firehose_event = self.convert_jetstream_commit(&did, &commit).await; + let firehose_event = self.convert_jetstream_commit(&commit_event).await; if let Some(event) = firehose_event { - tracing::debug!( - "Converted Jetstream commit in {:?}: {} {} {}, op={:?}", - start.elapsed(), - did, - commit.collection, - commit.rkey, - commit.op - ); self.consume_firehose_event(event, threads, &submit).await; } else { tracing::warn!( "Failed to convert Jetstream commit: {} {} {}, op={:?}", - did, - commit.collection, - commit.rkey, - commit.op + commit_event.did, + commit_event.commit.collection, + commit_event.commit.rkey, + commit_event.commit.op ); // Increment failure metric - let collection = commit.collection.to_string(); + let collection = commit_event.commit.collection.to_string(); counter!("jetstream_events.conversion_failed", "collection" => collection).increment(1); } }, @@ -676,44 +640,22 @@ impl RelayIndexer { counter!("jetstream_events.parse_error", "type" => "other").increment(1); } - // Update the sequence number if we can parse it from the JSON - if let Some(time_us) = extract_sequence_from_preview(&text_preview) { - if time_us > self.consumer.current_seq() { - // This is a temporary workaround to ensure we don't get stuck - // processing the same malformed message repeatedly - tracing::debug!("Extracted sequence {} from malformed event", time_us); - - // Since we can't directly update the consumer's sequence, - // we'll just log it for now - a proper fix would update the internal state - } - } - continue; }, - crate::jetstream::JetstreamOutput::Identity(did, identity) => { + crate::jetstream::JetstreamOutput::Identity(identity_event) => { // Map Jetstream identity to Firehose identity event - let firehose_event = self.convert_jetstream_identity(&did, &identity); - tracing::debug!( - "Processing Jetstream identity update for {}: handle={}", - did, - identity.handle - ); + let firehose_event = self.convert_jetstream_identity(&identity_event); // Track identity update metric counter!("jetstream_events.identity_update").increment(1); self.consume_firehose_event(firehose_event, threads, &submit).await; }, - crate::jetstream::JetstreamOutput::Account(did, account) => { + crate::jetstream::JetstreamOutput::Account(account_event) => { // Map Jetstream account to Firehose account event - let firehose_event = self.convert_jetstream_account(&did, &account); - tracing::debug!( - "Processing Jetstream account update for {}: active={}", - did, - account.active - ); + let firehose_event = self.convert_jetstream_account(&account_event); // Track account update metric - counter!("jetstream_events.account_update", "active" => if account.active { "true" } else { "false" }).increment(1); + counter!("jetstream_events.account_update", "active" => if account_event.account.active { "true" } else { "false" }).increment(1); self.consume_firehose_event(firehose_event, threads, &submit).await; } } @@ -1231,7 +1173,7 @@ async fn process_op( } let Some((cid, decoded)) = decode_op(op, blocks) else { - tracing::debug!("failed to decode op for {} {}", op.action, full_path); + tracing::warn!("failed to decode op for {} {}", op.action, full_path); return Ok(()); }; @@ -1258,7 +1200,7 @@ fn decode_op(op: &CommitOp, blocks: &HashMap>) -> Option<(Cid, Reco return Some((cid, record)); } Err(err) => { - tracing::debug!("Failed to deserialize Jetstream record: {}", err); + tracing::warn!("Failed to deserialize Jetstream record: {}", err); return None; } } diff --git a/consumer/src/jetstream/consumer.rs b/consumer/src/jetstream/consumer.rs index 8906951d..d4285c25 100644 --- a/consumer/src/jetstream/consumer.rs +++ b/consumer/src/jetstream/consumer.rs @@ -116,7 +116,7 @@ impl JetstreamConsumer { }; let json = serde_json::to_string(&message)?; - self.sink.send(Message::Text(json.into())).await?; + self.sink.send(Message::Text(json)).await?; debug!("Sent options update to Jetstream"); Ok(()) @@ -163,7 +163,7 @@ impl JetstreamConsumer { self.seq = commit.time_us; } - Ok(JetstreamOutput::Commit(commit.did, commit.commit)) + Ok(JetstreamOutput::Commit(commit)) } JetstreamEvent::Identity(identity) => { // Update cursor @@ -171,7 +171,7 @@ impl JetstreamConsumer { self.seq = identity.time_us; } - Ok(JetstreamOutput::Identity(identity.did, identity.identity)) + Ok(JetstreamOutput::Identity(identity)) } JetstreamEvent::Account(account) => { // Update cursor @@ -179,7 +179,7 @@ impl JetstreamConsumer { self.seq = account.time_us; } - Ok(JetstreamOutput::Account(account.did, account.account)) + Ok(JetstreamOutput::Account(account)) } } } @@ -192,7 +192,6 @@ impl JetstreamConsumer { // Decompress using zstd match Self::decompress_zstd(compressed_data) { Ok(decompressed) => { - trace!("Decompressed message length: {}", decompressed.len()); let text = String::from_utf8(decompressed).map_err(|e| { eyre!("Failed to decode decompressed data as UTF-8: {}", e) })?; @@ -215,19 +214,19 @@ impl JetstreamConsumer { if commit.time_us > self.seq { self.seq = commit.time_us; } - Ok(JetstreamOutput::Commit(commit.did, commit.commit)) + Ok(JetstreamOutput::Commit(commit)) } JetstreamEvent::Identity(identity) => { if identity.time_us > self.seq { self.seq = identity.time_us; } - Ok(JetstreamOutput::Identity(identity.did, identity.identity)) + Ok(JetstreamOutput::Identity(identity)) } JetstreamEvent::Account(account) => { if account.time_us > self.seq { self.seq = account.time_us; } - Ok(JetstreamOutput::Account(account.did, account.account)) + Ok(JetstreamOutput::Account(account)) } } } @@ -238,15 +237,10 @@ impl JetstreamConsumer { } } else { // If it's not a compressed message, try to handle it as a JSON message directly - debug!( - "Received uncompressed binary message of length {}", - data.len() - ); // Try to decode as UTF-8 match String::from_utf8(data.clone()) { Ok(text) => { - trace!("Successfully decoded binary message as UTF-8"); info!( "Binary message decoded as text (first 100 chars): {:.100}...", text.chars().take(100).collect::() @@ -270,19 +264,19 @@ impl JetstreamConsumer { if commit.time_us > self.seq { self.seq = commit.time_us; } - Ok(JetstreamOutput::Commit(commit.did, commit.commit)) + Ok(JetstreamOutput::Commit(commit)) } JetstreamEvent::Identity(identity) => { if identity.time_us > self.seq { self.seq = identity.time_us; } - Ok(JetstreamOutput::Identity(identity.did, identity.identity)) + Ok(JetstreamOutput::Identity(identity)) } JetstreamEvent::Account(account) => { if account.time_us > self.seq { self.seq = account.time_us; } - Ok(JetstreamOutput::Account(account.did, account.account)) + Ok(JetstreamOutput::Account(account)) } } } diff --git a/consumer/src/jetstream/types.rs b/consumer/src/jetstream/types.rs index 9e167537..173959ef 100644 --- a/consumer/src/jetstream/types.rs +++ b/consumer/src/jetstream/types.rs @@ -1,33 +1,21 @@ use serde::{Deserialize, Serialize}; -use std::collections::HashMap; // Jetstream event examples +// Jetstream events have 3 `kinds`s (so far): // -// Jetstream events end up looking something like: +// - `commit`: a Commit to a repo which involves either a create, update, or delete of a record +// - `identity`: an Identity update for a DID which indicates that you may want to purge an identity cache and revalidate the DID doc and handle +// - `account`: an Account event that indicates a change in account status i.e. from `active` to `deactivated`, or to `takendown` if the PDS has taken down the repo. // -// { -// "did": "did:plc:eygmaihciaxprqvxpfvl6flk", -// "time_us": 1725911162329308, -// "kind": "commit", -// "commit": { -// "rev": "3l3qo2vutsw2b", -// "operation": "create", -// "collection": "app.bsky.feed.like", -// "rkey": "3l3qo2vuowo2b", -// "record": { -// "$type": "app.bsky.feed.like", -// "createdAt": "2024-09-09T19:46:02.102Z", -// "subject": { -// "cid": "bafyreidc6sydkkbchcyg62v77wbhzvb2mvytlmsychqgwf2xojjtirmzj4", -// "uri": "at://did:plc:wa7b35aakoll7hugkrjtf3xf/app.bsky.feed.post/3l3pte3p2e325" -// } -// }, -// "cid": "bafyreidwaivazkwu67xztlmuobx35hs2lnfh3kolmgfmucldvhd3sgzcqi" -// } -// } +// Jetstream Commits have 3 `operations`: +// +// - `create`: Create a new record with the contents provided +// - `update`: Update an existing record and replace it with the contents provided +// - `delete`: Delete an existing record with the DID, Collection, and RKey provided +// +// Jetstream events end up looking something like: // // A like committed to a repo - // { // "did": "did:plc:eygmaihciaxprqvxpfvl6flk", // "time_us": 1725911162329308, @@ -48,9 +36,7 @@ use std::collections::HashMap; // "cid": "bafyreidwaivazkwu67xztlmuobx35hs2lnfh3kolmgfmucldvhd3sgzcqi" // } // } - // A deleted follow record - // { // "did": "did:plc:rfov6bpyztcnedeyyzgfq42k", // "time_us": 1725516666833633, @@ -62,9 +48,7 @@ use std::collections::HashMap; // "rkey": "3l3dn7tku762u" // } // } - // An identity update - // { // "did": "did:plc:ufbl4k27gp6kzas5glhz7fim", // "time_us": 1725516665234703, @@ -76,9 +60,7 @@ use std::collections::HashMap; // "time": "2024-09-05T06:11:04.870Z" // } // } - // An account becoming active - // { // "did": "did:plc:ufbl4k27gp6kzas5glhz7fim", // "time_us": 1725516665333808, @@ -218,9 +200,9 @@ pub struct SubscriberSourcedMessage { /// Output variants from the Jetstream consumer #[derive(Debug)] pub enum JetstreamOutput { - Commit(String, CommitPayload), - Identity(String, IdentityPayload), - Account(String, AccountPayload), + Commit(CommitEvent), + Identity(IdentityEvent), + Account(AccountEvent), Error(JetstreamError), /// Represents a parsing error - includes error message and a sample of the text that failed to parse ParseError(String, String), -- 2.51.2