diff --git a/jetstream/Cargo.toml b/jetstream/Cargo.toml index 5836560..f05f616 100644 --- a/jetstream/Cargo.toml +++ b/jetstream/Cargo.toml @@ -22,7 +22,7 @@ tokio-tungstenite = { version = "0.26.2", features = [ futures-util = "0.3.31" url = "2.5.4" serde = { version = "1.0.215", features = ["derive"] } -serde_json = "1.0.132" +serde_json = { version = "1.0.140", features = ["raw_value"] } chrono = "0.4.38" zstd = "0.13.2" thiserror = "2.0.3" diff --git a/jetstream/examples/arbitrary_record.rs b/jetstream/examples/arbitrary_record.rs index 794e849..fd162d1 100644 --- a/jetstream/examples/arbitrary_record.rs +++ b/jetstream/examples/arbitrary_record.rs @@ -5,8 +5,9 @@ use atrium_api::types::string; use clap::Parser; use jetstream::{ events::{ - commit::CommitEvent, - JetstreamEvent::Commit, + CommitOp, + EventKind, + JetstreamEvent, }, DefaultJetstreamEndpoints, JetstreamCompression, @@ -30,7 +31,7 @@ async fn main() -> anyhow::Result<()> { let args = Args::parse(); let dids = args.did.unwrap_or_default(); - let config: JetstreamConfig = JetstreamConfig { + let config: JetstreamConfig = JetstreamConfig { endpoint: DefaultJetstreamEndpoints::USEastOne.into(), wanted_collections: vec![args.nsid.clone()], wanted_dids: dids.clone(), @@ -48,8 +49,24 @@ async fn main() -> anyhow::Result<()> { ); while let Some(event) = receiver.recv().await { - if let Commit(CommitEvent::CreateOrUpdate { commit, .. }) = event { - println!("got record: {:?}", commit.record); + if let JetstreamEvent { + kind: EventKind::Commit, + commit: Some(commit), + .. + } = event + { + if commit.collection != args.nsid { + continue; + } + if !(commit.operation == CommitOp::Create || commit.operation == CommitOp::Update) { + continue; + } + let Some(rec) = commit.record else { continue }; + println!( + "New or updated record! ({})\n{:?}\n", + commit.rkey.as_str(), + rec.get() + ); } } diff --git a/jetstream/examples/basic.rs b/jetstream/examples/basic.rs index e043da8..adb7282 100644 --- a/jetstream/examples/basic.rs +++ b/jetstream/examples/basic.rs @@ -7,11 +7,10 @@ use atrium_api::{ use clap::Parser; use jetstream::{ events::{ - commit::{ - CommitEvent, - CommitType, - }, - JetstreamEvent::Commit, + CommitEvent, + CommitOp, + EventKind, + JetstreamEvent, }, DefaultJetstreamEndpoints, JetstreamCompression, @@ -25,9 +24,6 @@ struct Args { /// The DIDs to listen for events on, if not provided we will listen for all DIDs. #[arg(short, long)] did: Option>, - /// The NSID for the collection to listen for (e.g. `app.bsky.feed.post`). - #[arg(short, long)] - nsid: string::Nsid, } #[tokio::main] @@ -37,7 +33,7 @@ async fn main() -> anyhow::Result<()> { let dids = args.did.unwrap_or_default(); let config = JetstreamConfig { endpoint: DefaultJetstreamEndpoints::USEastOne.into(), - wanted_collections: vec![args.nsid.clone()], + wanted_collections: vec![string::Nsid::new("app.bsky.feed.post".to_string()).unwrap()], wanted_dids: dids.clone(), compression: JetstreamCompression::Zstd, ..Default::default() @@ -46,30 +42,23 @@ async fn main() -> anyhow::Result<()> { let jetstream = JetstreamConnector::new(config)?; let mut receiver = jetstream.connect().await?; - println!( - "Listening for '{}' events on DIDs: {:?}", - args.nsid.as_str(), - dids - ); + println!("Listening for 'app.bsky.feed.post' events on DIDs: {dids:?}"); while let Some(event) = receiver.recv().await { - if let Commit(commit) = event { - match commit { - CommitEvent::CreateOrUpdate { info: _, commit } - if commit.info.operation == CommitType::Create => - { - if let AppBskyFeedPost(record) = commit.record { - println!( - "New post created! ({})\n\n'{}'", - commit.info.rkey.as_str(), - record.text - ); - } - } - CommitEvent::Delete { info: _, commit } => { - println!("A post has been deleted. ({})", commit.rkey.as_str()); - } - _ => {} + if let JetstreamEvent { + kind: EventKind::Commit, + commit: + Some(CommitEvent { + operation: CommitOp::Create, + rkey, + record: Some(record), + .. + }), + .. + } = event + { + if let Ok(AppBskyFeedPost(rec)) = serde_json::from_str(record.get()) { + println!("New post created! ({})\n{:?}\n", rkey.as_str(), rec.text); } } } diff --git a/jetstream/src/events/account.rs b/jetstream/src/events/account.rs deleted file mode 100644 index eebea17..0000000 --- a/jetstream/src/events/account.rs +++ /dev/null @@ -1,40 +0,0 @@ -use chrono::Utc; -use serde::Deserialize; - -use crate::{ - events::EventInfo, - exports, -}; - -/// An event representing a change to an account. -#[derive(Deserialize, Debug)] -pub struct AccountEvent { - /// Basic metadata included with every event. - #[serde(flatten)] - pub info: EventInfo, - /// Account specific data bundled with this event. - pub account: AccountData, -} - -/// Account specific data bundled with an account event. -#[derive(Deserialize, Debug)] -pub struct AccountData { - /// Whether the account is currently active. - pub active: bool, - /// The DID of the account. - pub did: exports::Did, - pub seq: u64, - pub time: chrono::DateTime, - /// If `active` is `false` this will be present to explain why the account is inactive. - pub status: Option, -} - -/// The possible reasons an account might be listed as inactive. -#[derive(Deserialize, Debug)] -#[serde(rename_all = "lowercase")] -pub enum AccountStatus { - Deactivated, - Deleted, - Suspended, - TakenDown, -} diff --git a/jetstream/src/events/commit.rs b/jetstream/src/events/commit.rs deleted file mode 100644 index 1014db0..0000000 --- a/jetstream/src/events/commit.rs +++ /dev/null @@ -1,55 +0,0 @@ -use serde::Deserialize; - -use crate::{ - events::EventInfo, - exports, -}; - -/// An event representing a repo commit, which can be a `create`, `update`, or `delete` operation. -#[derive(Deserialize, Debug)] -#[serde(untagged, rename_all = "snake_case")] -pub enum CommitEvent { - CreateOrUpdate { - #[serde(flatten)] - info: EventInfo, - commit: CommitData, - }, - Delete { - #[serde(flatten)] - info: EventInfo, - commit: CommitInfo, - }, -} - -/// The type of commit operation that was performed. -#[derive(Deserialize, Debug, PartialEq)] -#[serde(rename_all = "snake_case")] -pub enum CommitType { - Create, - Update, - Delete, -} - -/// Basic commit specific info bundled with every event, also the only data included with a `delete` -/// operation. -#[derive(Deserialize, Debug)] -pub struct CommitInfo { - /// The type of commit operation that was performed. - pub operation: CommitType, - pub rev: String, - pub rkey: exports::RecordKey, - /// The NSID of the record type that this commit is associated with. - pub collection: exports::Nsid, -} - -/// Detailed data bundled with a commit event. This data is only included when the event is -/// `create` or `update`. -#[derive(Deserialize, Debug)] -pub struct CommitData { - #[serde(flatten)] - pub info: CommitInfo, - /// The CID of the record that was operated on. - pub cid: exports::Cid, - /// The record that was operated on. - pub record: R, -} diff --git a/jetstream/src/events/identity.rs b/jetstream/src/events/identity.rs deleted file mode 100644 index 703ec9f..0000000 --- a/jetstream/src/events/identity.rs +++ /dev/null @@ -1,28 +0,0 @@ -use chrono::Utc; -use serde::Deserialize; - -use crate::{ - events::EventInfo, - exports, -}; - -/// An event representing a change to an identity. -#[derive(Deserialize, Debug)] -pub struct IdentityEvent { - /// Basic metadata included with every event. - #[serde(flatten)] - pub info: EventInfo, - /// Identity specific data bundled with this event. - pub identity: IdentityData, -} - -/// Identity specific data bundled with an identity event. -#[derive(Deserialize, Debug)] -pub struct IdentityData { - /// The DID of the identity. - pub did: exports::Did, - /// The handle associated with the identity. - pub handle: Option, - pub seq: u64, - pub time: chrono::DateTime, -} diff --git a/jetstream/src/events/mod.rs b/jetstream/src/events/mod.rs index 9573270..7a904f5 100644 --- a/jetstream/src/events/mod.rs +++ b/jetstream/src/events/mod.rs @@ -1,7 +1,3 @@ -pub mod account; -pub mod commit; -pub mod identity; - use std::time::{ Duration, SystemTime, @@ -9,33 +5,29 @@ use std::time::{ UNIX_EPOCH, }; +use chrono::Utc; use serde::Deserialize; +use serde_json::value::RawValue; use crate::exports; /// Opaque wrapper for the time_us cursor used by jetstream -/// -/// Generally, you should use a cursor #[derive(Deserialize, Debug, Clone, PartialEq, PartialOrd)] pub struct Cursor(u64); -/// Basic data that is included with every event. -#[derive(Deserialize, Debug)] -pub struct EventInfo { +#[derive(Debug, Deserialize)] +#[serde(rename_all = "snake_case")] +pub struct JetstreamEvent { + #[serde(rename = "time_us")] + pub cursor: Cursor, pub did: exports::Did, - pub time_us: Cursor, pub kind: EventKind, + pub commit: Option, + pub identity: Option, + pub account: Option, } -#[derive(Deserialize, Debug)] -#[serde(untagged)] -pub enum JetstreamEvent { - Commit(commit::CommitEvent), - Identity(identity::IdentityEvent), - Account(account::AccountEvent), -} - -#[derive(Deserialize, Debug)] +#[derive(Debug, Deserialize, PartialEq)] #[serde(rename_all = "snake_case")] pub enum EventKind { Commit, @@ -43,19 +35,40 @@ pub enum EventKind { Account, } -impl JetstreamEvent { - pub fn cursor(&self) -> Cursor { - match self { - JetstreamEvent::Commit(commit::CommitEvent::CreateOrUpdate { info, .. }) => { - info.time_us.clone() - } - JetstreamEvent::Commit(commit::CommitEvent::Delete { info, .. }) => { - info.time_us.clone() - } - JetstreamEvent::Identity(e) => e.info.time_us.clone(), - JetstreamEvent::Account(e) => e.info.time_us.clone(), - } - } +#[derive(Debug, Deserialize)] +#[serde(rename_all = "snake_case")] +pub struct CommitEvent { + pub collection: exports::Nsid, + pub rkey: exports::RecordKey, + pub rev: String, + pub operation: CommitOp, + pub record: Option>, + pub cid: Option, +} + +#[derive(Debug, Deserialize, PartialEq)] +#[serde(rename_all = "snake_case")] +pub enum CommitOp { + Create, + Update, + Delete, +} + +#[derive(Debug, Deserialize, PartialEq)] +pub struct IdentityEvent { + pub did: exports::Did, + pub handle: Option, + pub seq: u64, + pub time: chrono::DateTime, +} + +#[derive(Debug, Deserialize, PartialEq)] +pub struct AccountEvent { + pub active: bool, + pub did: exports::Did, + pub seq: u64, + pub time: chrono::DateTime, + pub status: Option, } impl Cursor { @@ -136,3 +149,48 @@ impl From<&Cursor> for SystemTime { UNIX_EPOCH + Duration::from_micros(c.0) } } + +#[cfg(test)] +mod test { + use super::*; + + #[test] + fn test_parse_commit_event() -> anyhow::Result<()> { + let json = r#"{ + "rev":"3llrdsginou2i", + "operation":"create", + "collection":"app.bsky.feed.post", + "rkey":"3llrdsglqdc2s", + "cid": "bafyreidofvwoqvd2cnzbun6dkzgfucxh57tirf3ohhde7lsvh4fu3jehgy", + "record": {"$type":"app.bsky.feed.post","createdAt":"2025-04-01T16:58:06.154Z","langs":["en"],"text":"I wish apirl 1st would stop existing lol"} + }"#; + let commit: CommitEvent = serde_json::from_str(json)?; + assert_eq!( + commit.cid.unwrap(), + "bafyreidofvwoqvd2cnzbun6dkzgfucxh57tirf3ohhde7lsvh4fu3jehgy".parse()? + ); + assert_eq!( + commit.record.unwrap().get(), + r#"{"$type":"app.bsky.feed.post","createdAt":"2025-04-01T16:58:06.154Z","langs":["en"],"text":"I wish apirl 1st would stop existing lol"}"# + ); + Ok(()) + } + + #[test] + fn test_parse_whole_event() -> anyhow::Result<()> { + let json = r#"{"did":"did:plc:ai3dzf35cth7s3st7n7jsd7r","time_us":1743526687419798,"kind":"commit","commit":{"rev":"3llrdsginou2i","operation":"create","collection":"app.bsky.feed.post","rkey":"3llrdsglqdc2s","record":{"$type":"app.bsky.feed.post","createdAt":"2025-04-01T16:58:06.154Z","langs":["en"],"text":"I wish apirl 1st would stop existing lol"},"cid":"bafyreidofvwoqvd2cnzbun6dkzgfucxh57tirf3ohhde7lsvh4fu3jehgy"}}"#; + let event: JetstreamEvent = serde_json::from_str(json)?; + assert_eq!(event.kind, EventKind::Commit); + assert!(event.commit.is_some()); + let commit = event.commit.unwrap(); + assert_eq!( + commit.cid.unwrap(), + "bafyreidofvwoqvd2cnzbun6dkzgfucxh57tirf3ohhde7lsvh4fu3jehgy".parse()? + ); + assert_eq!( + commit.record.unwrap().get(), + r#"{"$type":"app.bsky.feed.post","createdAt":"2025-04-01T16:58:06.154Z","langs":["en"],"text":"I wish apirl 1st would stop existing lol"}"# + ); + Ok(()) + } +} diff --git a/jetstream/src/lib.rs b/jetstream/src/lib.rs index c17484a..e0011a4 100644 --- a/jetstream/src/lib.rs +++ b/jetstream/src/lib.rs @@ -7,19 +7,16 @@ use std::{ Cursor as IoCursor, Read, }, - marker::PhantomData, time::{ Duration, Instant, }, }; -use atrium_api::record::KnownRecord; use futures_util::{ stream::StreamExt, SinkExt, }; -use serde::de::DeserializeOwned; use tokio::{ net::TcpStream, sync::mpsc::{ @@ -124,16 +121,16 @@ const MAX_WANTED_DIDS: usize = 10_000; const JETSTREAM_ZSTD_DICTIONARY: &[u8] = include_bytes!("../zstd/dictionary"); /// A receiver channel for consuming Jetstream events. -pub type JetstreamReceiver = Receiver>; +pub type JetstreamReceiver = Receiver; /// An internal sender channel for sending Jetstream events to [JetstreamReceiver]'s. -type JetstreamSender = Sender>; +type JetstreamSender = Sender; /// A wrapper connector type for working with a WebSocket connection to a Jetstream instance to /// receive and consume events. See [JetstreamConnector::connect] for more info. -pub struct JetstreamConnector { +pub struct JetstreamConnector { /// The configuration for the Jetstream connection. - config: JetstreamConfig, + config: JetstreamConfig, } pub enum JetstreamCompression { @@ -163,7 +160,7 @@ impl From for JetstreamCompression { } } -pub struct JetstreamConfig { +pub struct JetstreamConfig { /// A Jetstream endpoint to connect to with a WebSocket Scheme i.e. /// `wss://jetstream1.us-east.bsky.network/subscribe`. pub endpoint: String, @@ -200,16 +197,9 @@ pub struct JetstreamConfig { /// can help prevent that if your consumer sometimes pauses, at a cost of higher memory /// usage while events are buffered. pub channel_size: usize, - /// Marker for record deserializable type. - /// - /// See examples/arbitrary_record.rs for an example using serde_json::Value - /// - /// You can omit this if you construct `JetstreamConfig { a: b, ..Default::default() }. - /// If you have to specify it, use `std::marker::PhantomData` with no type parameters. - pub record_type: PhantomData, } -impl Default for JetstreamConfig { +impl Default for JetstreamConfig { fn default() -> Self { JetstreamConfig { endpoint: DefaultJetstreamEndpoints::USEastOne.into(), @@ -220,12 +210,11 @@ impl Default for JetstreamConfig { omit_user_agent_jetstream_info: false, replay_on_reconnect: false, channel_size: 4096, // a few seconds of firehose buffer - record_type: PhantomData, } } } -impl JetstreamConfig { +impl JetstreamConfig { /// Constructs a new endpoint URL with the given [JetstreamConfig] applied. pub fn get_request_builder( &self, @@ -313,11 +302,11 @@ impl JetstreamConfig { } } -impl JetstreamConnector { +impl JetstreamConnector { /// Create a Jetstream connector with a valid [JetstreamConfig]. /// /// After creation, you can call [connect] to connect to the provided Jetstream instance. - pub fn new(config: JetstreamConfig) -> Result { + pub fn new(config: JetstreamConfig) -> Result { // We validate the configuration here so any issues are caught early. config.validate()?; Ok(JetstreamConnector { config }) @@ -327,7 +316,7 @@ impl JetstreamConnector { /// /// A [JetstreamReceiver] is returned which can be used to respond to events. When all instances /// of this receiver are dropped, the connection and task are automatically closed. - pub async fn connect(&self) -> Result, ConnectionError> { + pub async fn connect(&self) -> Result { self.connect_cursor(None).await } @@ -343,7 +332,7 @@ impl JetstreamConnector { pub async fn connect_cursor( &self, cursor: Option, - ) -> Result, ConnectionError> { + ) -> Result { // We validate the config again for good measure. Probably not necessary but it can't hurt. self.config .validate() @@ -424,10 +413,10 @@ impl JetstreamConnector { /// The main task that handles the WebSocket connection and sends [JetstreamEvent]'s to any /// receivers that are listening for them. -async fn websocket_task( +async fn websocket_task( dictionary: DecoderDictionary<'_>, ws: WebSocketStream>, - send_channel: JetstreamSender, + send_channel: JetstreamSender, last_cursor: &mut Option, ) -> Result<(), JetstreamEventError> { // TODO: Use the write half to allow the user to change configuration settings on the fly. @@ -439,9 +428,9 @@ async fn websocket_task( Some(Ok(message)) => { match message { Message::Text(json) => { - let event: JetstreamEvent = serde_json::from_str(&json) + let event: JetstreamEvent = serde_json::from_str(&json) .map_err(JetstreamEventError::ReceivedMalformedJSON)?; - let event_cursor = event.cursor(); + let event_cursor = event.cursor.clone(); if let Some(last) = last_cursor { if event_cursor <= *last { @@ -475,9 +464,11 @@ async fn websocket_task( .read_to_string(&mut json) .map_err(JetstreamEventError::CompressionDecoderError)?; - let event: JetstreamEvent = serde_json::from_str(&json) - .map_err(JetstreamEventError::ReceivedMalformedJSON)?; - let event_cursor = event.cursor(); + let event: JetstreamEvent = serde_json::from_str(&json).map_err(|e| { + eprintln!("lkasjdflkajsd {e:?} {json}"); + JetstreamEventError::ReceivedMalformedJSON(e) + })?; + let event_cursor = event.cursor.clone(); if let Some(last) = last_cursor { if event_cursor <= *last { diff --git a/ufos/src/consumer.rs b/ufos/src/consumer.rs index a0e0699..021d6b9 100644 --- a/ufos/src/consumer.rs +++ b/ufos/src/consumer.rs @@ -1,9 +1,5 @@ use jetstream::{ - events::{ - account::AccountEvent, - commit::{CommitData, CommitEvent, CommitInfo, CommitType}, - Cursor, EventInfo, JetstreamEvent, - }, + events::{CommitEvent, CommitOp, Cursor, EventKind, JetstreamEvent}, exports::Did, DefaultJetstreamEndpoints, JetstreamCompression, JetstreamConfig, JetstreamConnector, JetstreamReceiver, @@ -26,7 +22,7 @@ const BATCH_QUEUE_SIZE: usize = 64; // 4096 got OOM'd. update: 1024 also got OOM #[derive(Debug)] struct Batcher { - jetstream_receiver: JetstreamReceiver, + jetstream_receiver: JetstreamReceiver, batch_sender: Sender, current_batch: EventBatch, } @@ -42,14 +38,14 @@ pub async fn consume( } else { eprintln!("connecting to jetstream at {jetstream_endpoint} => {endpoint}"); } - let config: JetstreamConfig = JetstreamConfig { + let config: JetstreamConfig = JetstreamConfig { endpoint, compression: if no_compress { JetstreamCompression::None } else { JetstreamCompression::Zstd }, - channel_size: 64, // small because we'd rather buffer events into batches + channel_size: 64, // small because we expect to be fast....? ..Default::default() }; let jetstream_receiver = JetstreamConnector::new(config)? @@ -62,10 +58,7 @@ pub async fn consume( } impl Batcher { - fn new( - jetstream_receiver: JetstreamReceiver, - batch_sender: Sender, - ) -> Self { + fn new(jetstream_receiver: JetstreamReceiver, batch_sender: Sender) -> Self { Self { jetstream_receiver, batch_sender, @@ -83,11 +76,8 @@ impl Batcher { } } - async fn handle_event( - &mut self, - event: JetstreamEvent, - ) -> anyhow::Result<()> { - let event_cursor = event.cursor(); + async fn handle_event(&mut self, event: JetstreamEvent) -> anyhow::Result<()> { + let event_cursor = event.cursor; if let Some(earliest) = &self.current_batch.first_jetstream_cursor { if event_cursor.duration_since(earliest)? > Duration::from_secs_f64(MAX_BATCH_SPAN_SECS) @@ -98,28 +88,40 @@ impl Batcher { self.current_batch.first_jetstream_cursor = Some(event_cursor.clone()); } - match event { - JetstreamEvent::Commit(CommitEvent::CreateOrUpdate { commit, info }) => { - match commit.info.operation { - CommitType::Create => self.handle_create_record(commit, info).await?, - CommitType::Update => { - self.handle_modify_record(modify_update(commit, info)) - .await? + match event.kind { + EventKind::Commit if event.commit.is_some() => { + let commit = event.commit.unwrap(); + match commit.operation { + CommitOp::Create => { + self.handle_create_record(event.did, commit, event_cursor.clone()) + .await?; + } + CommitOp::Update => { + self.handle_modify_record(modify_update( + event.did, + commit, + event_cursor.clone(), + )) + .await?; } - CommitType::Delete => { - panic!("jetstream Commit::CreateOrUpdate had Delete operation type") + CommitOp::Delete => { + self.handle_modify_record(modify_delete( + event.did, + commit, + event_cursor.clone(), + )) + .await?; } } } - JetstreamEvent::Commit(CommitEvent::Delete { commit, info }) => { - self.handle_modify_record(modify_delete(commit, info)) - .await? - } - JetstreamEvent::Account(AccountEvent { info, account }) if !account.active => { - self.handle_remove_account(info.did, info.time_us).await? + EventKind::Account if event.account.is_some() => { + let account = event.account.unwrap(); + if !account.active { + self.handle_remove_account(account.did, event_cursor.clone()) + .await?; + } } - JetstreamEvent::Account(_) => {} // ignore account *activations* - JetstreamEvent::Identity(_) => {} // identity events are noops for us + _ => {} }; self.current_batch.last_jetstream_cursor = Some(event_cursor.clone()); @@ -159,27 +161,29 @@ impl Batcher { async fn handle_create_record( &mut self, - commit: CommitData, - info: EventInfo, + did: Did, + commit: CommitEvent, + cursor: Cursor, ) -> anyhow::Result<()> { if !self .current_batch .record_creates - .contains_key(&commit.info.collection) + .contains_key(&commit.collection) && self.current_batch.record_creates.len() >= MAX_BATCHED_COLLECTIONS { self.send_current_batch_now().await?; } + let record = serde_json::from_str(commit.record.unwrap().get())?; let record = CreateRecord { - did: info.did, - rkey: commit.info.rkey, - record: commit.record, - cursor: info.time_us, + did, + rkey: commit.rkey, + record, + cursor, }; let collection = self .current_batch .record_creates - .entry(commit.info.collection) + .entry(commit.collection) .or_default(); collection.total_seen += 1; collection.samples.push_front(record); @@ -206,21 +210,22 @@ impl Batcher { } } -fn modify_update(commit: CommitData, info: EventInfo) -> ModifyRecord { +fn modify_update(did: Did, commit: CommitEvent, cursor: Cursor) -> ModifyRecord { + let record = serde_json::from_str(commit.record.unwrap().get()).unwrap(); ModifyRecord::Update(UpdateRecord { - did: info.did, - collection: commit.info.collection, - rkey: commit.info.rkey, - record: commit.record, - cursor: info.time_us, + did, + collection: commit.collection, + rkey: commit.rkey, + record, + cursor, }) } -fn modify_delete(commit_info: CommitInfo, info: EventInfo) -> ModifyRecord { +fn modify_delete(did: Did, commit: CommitEvent, cursor: Cursor) -> ModifyRecord { ModifyRecord::Delete(DeleteRecord { - did: info.did, - collection: commit_info.collection, - rkey: commit_info.rkey, - cursor: info.time_us, + did, + collection: commit.collection, + rkey: commit.rkey, + cursor, }) } diff --git a/ufos/src/lib.rs b/ufos/src/lib.rs index 7375254..dd7c351 100644 --- a/ufos/src/lib.rs +++ b/ufos/src/lib.rs @@ -1,7 +1,8 @@ pub mod consumer; pub mod db_types; pub mod server; -pub mod store; +// pub mod storage; +pub mod storage_fjall; pub mod store_types; use jetstream::events::Cursor; diff --git a/ufos/src/main.rs b/ufos/src/main.rs index 48ca10b..b91bbe3 100644 --- a/ufos/src/main.rs +++ b/ufos/src/main.rs @@ -1,6 +1,6 @@ use clap::Parser; use std::path::PathBuf; -use ufos::{consumer, server, store}; +use ufos::{consumer, server, storage_fjall}; #[cfg(not(target_env = "msvc"))] use tikv_jemallocator::Jemalloc; @@ -43,7 +43,7 @@ async fn main() -> anyhow::Result<()> { let args = Args::parse(); let (storage, cursor) = - store::Storage::open(args.data, &args.jetstream, args.jetstream_force).await?; + storage_fjall::Storage::open(args.data, &args.jetstream, args.jetstream_force).await?; println!("starting server with storage..."); let serving = server::serve(storage.clone()); diff --git a/ufos/src/server.rs b/ufos/src/server.rs index 4b88841..0e4c847 100644 --- a/ufos/src/server.rs +++ b/ufos/src/server.rs @@ -1,4 +1,4 @@ -use crate::store::{Storage, StorageInfo}; +use crate::storage_fjall::{Storage, StorageInfo}; use crate::{CreateRecord, Nsid}; use dropshot::endpoint; use dropshot::ApiDescription; diff --git a/ufos/src/store.rs b/ufos/src/storage_fjall.rs similarity index 99% rename from ufos/src/store.rs rename to ufos/src/storage_fjall.rs index 3c9e9c2..e43a059 100644 --- a/ufos/src/store.rs +++ b/ufos/src/storage_fjall.rs @@ -117,7 +117,7 @@ impl Storage { // TODO: see rw_loop: enforce single-thread. loop { let t_sleep = Instant::now(); - sleep(Duration::from_secs_f64(0.8)).await; // TODO: minimize during replay + sleep(Duration::from_secs_f64(0.08)).await; // TODO: minimize during replay let slept_for = t_sleep.elapsed(); let queue_size = receiver.len();