use anyhow::{anyhow, Result}; use jacquard::{ api::network_bsky::jetstream::subscribe_events::{CommitOperation, SubscribeEventsMessage}, common::{deps::fluent_uri::Uri, websocket::tungstenite_client::TungsteniteClient, DefaultStr}, jetstream::{ archive::JetstreamClient, plan::{CollectionFilter, EventKind, ReplayFilters}, replay::{ReplayItem, ReplayMode, ReplayOptions, ReplayStream}, }, }; use shared::atproto::lexicon::{event, follow, EVENT, FOLLOW}; use std::time::Duration; #[derive(Debug)] pub enum FirehoseAction { Upsert { did: String, rkey: String, event: event::Event, seq: u64, }, FollowUpsert { did: String, rkey: String, follow: follow::Follow, seq: u64, }, Delete { did: String, rkey: String, collection: String, seq: u64, }, } pub struct Firehose { stream: ReplayStream, } impl Firehose { pub fn new(base: &str, api_key: &str, after_seq: Option) -> Result { let base = Uri::parse(base.to_string()) .map_err(|(e, _)| anyhow!("invalid jetstream base url {base}: {e}"))?; let archive = JetstreamClient::new( reqwest::Client::new(), base, Some(jacquard::common::deps::smol_str::SmolStr::from(api_key)), ); let filters = ReplayFilters { kinds: vec![EventKind::Commit], dids: Vec::new(), collections: vec![ CollectionFilter::parse(DefaultStr::from(EVENT)).expect("event filter"), CollectionFilter::parse(DefaultStr::from(FOLLOW)).expect("follow filter"), ], }; let stream = ReplayStream::new( archive, TungsteniteClient::new(), filters, ReplayMode::Replay { after_seq }, ReplayOptions::default(), ); Ok(Self { stream }) } pub async fn next_action(&mut self) -> Result> { loop { match self.stream.next().await { Ok(Some(item)) => { if let Some(action) = to_action(item)? { return Ok(Some(action)); } } Ok(None) => return Ok(None), Err(e) => { tracing::warn!("jetstream stream error (resumable): {e}"); tokio::time::sleep(Duration::from_secs(5)).await; } } } } } fn to_action(item: ReplayItem) -> Result> { let seq = item.last_seq.unwrap_or_default(); let message = item.message; let SubscribeEventsMessage::Commit(commit) = message else { return Ok(None); }; let collection = commit.collection.as_str(); if collection != EVENT && collection != FOLLOW { return Ok(None); } let did = commit.did.as_str().to_string(); let rkey = commit.rkey.0.as_str().to_string(); match commit.operation { CommitOperation::Create | CommitOperation::Update => { let record = commit .record .ok_or_else(|| anyhow!("commit without record: {did} {collection} {rkey}"))?; let value = serde_json::to_value(&record)?; if collection == FOLLOW { let follow: follow::Follow = serde_json::from_value(value)?; Ok(Some(FirehoseAction::FollowUpsert { did, rkey, follow, seq: seq as u64, })) } else { let event: event::Event = serde_json::from_value(value)?; Ok(Some(FirehoseAction::Upsert { did, rkey, event, seq: seq as u64, })) } } CommitOperation::Delete => Ok(Some(FirehoseAction::Delete { did, rkey, collection: collection.to_string(), seq: seq as u64, })), CommitOperation::Other(_) => Ok(None), } } #[cfg(test)] mod tests { use super::*; use jacquard::{ api::network_bsky::jetstream::subscribe_events::{Commit, Info, InfoName}, common::{ deps::smol_str::SmolStr, types::{string::Datetime, value::Data}, }, }; fn commit( operation: CommitOperation, collection: &str, rkey: &str, record: Option, ) -> ReplayItem { let record = record.map(|value| serde_json::from_value::>(value).unwrap()); let message = SubscribeEventsMessage::Commit(Box::new(Commit { cid: None, collection: collection.parse().unwrap(), did: "did:plc:abc".parse().unwrap(), operation, record, rev: "3jzl5bq2x5c2s".parse().unwrap(), rkey: rkey.parse().unwrap(), seq: 123_456_789, time: Datetime::now(), extra_data: None, })); ReplayItem { message, last_seq: Some(123_456_789), } } fn valid_record_json(when: &str) -> serde_json::Value { serde_json::json!({ "$type": "fyi.wed.event", "latitude": 525000000, "longitude": 132000000, "when": when, "name": "Picnic", "before": 1, "after": 3, "location": "Park" }) } fn valid_follow_json() -> serde_json::Value { serde_json::json!({ "$type": "fyi.wed.follow", "subject": "at://did:plc:bob/fyi.wed.event/ev1", "createdAt": "2026-08-13T00:00:00Z" }) } #[test] fn maps_event_create() { let item = commit( CommitOperation::Create, EVENT, "rkey1", Some(valid_record_json("2026-08-10 14:00:00")), ); let action = to_action(item).unwrap().unwrap(); match action { FirehoseAction::Upsert { did, rkey, event, seq, } => { assert_eq!(did, "did:plc:abc"); assert_eq!(rkey, "rkey1"); assert_eq!(event.latitude, 525000000); assert_eq!(event.name.as_deref(), Some("Picnic")); assert_eq!(seq, 123_456_789); } other => panic!("expected Upsert, got {other:?}"), } } #[test] fn maps_event_delete() { let item = commit(CommitOperation::Delete, EVENT, "rkey3", None); let action = to_action(item).unwrap().unwrap(); match action { FirehoseAction::Delete { did, rkey, collection, seq, } => { assert_eq!(did, "did:plc:abc"); assert_eq!(rkey, "rkey3"); assert_eq!(collection, EVENT); assert_eq!(seq, 123_456_789); } other => panic!("expected Delete, got {other:?}"), } } #[test] fn maps_follow_create() { let item = commit( CommitOperation::Create, FOLLOW, "rkey9", Some(valid_follow_json()), ); let action = to_action(item).unwrap().unwrap(); match action { FirehoseAction::FollowUpsert { did, rkey, follow, seq, } => { assert_eq!(did, "did:plc:abc"); assert_eq!(rkey, "rkey9"); assert_eq!( follow.subject.as_str(), "at://did:plc:bob/fyi.wed.event/ev1" ); assert_eq!(seq, 123_456_789); } other => panic!("expected FollowUpsert, got {other:?}"), } } #[test] fn skips_non_wed_collections() { let item = commit( CommitOperation::Create, "app.bsky.feed.post", "rkey4", Some(valid_record_json("2026-08-10 14:00:00")), ); assert!(to_action(item).unwrap().is_none()); } #[test] fn skips_identity_messages() { let item = ReplayItem { message: SubscribeEventsMessage::Info(Box::new(Info { message: Some("hello".to_string().into()), name: InfoName::OutdatedCursor, extra_data: None, })), last_seq: None, }; assert!(to_action(item).unwrap().is_none()); } }