spew_bsky_posts/
spew-bsky-posts.rs1use async_trait::async_trait;
2use rocketman::{
3 connection::JetstreamConnection,
4 handler::{self, Ingestors},
5 ingestion::LexiconIngestor,
6 options::JetstreamOptions,
7 types::event::{Commit, Event},
8};
9use serde_json::Value;
10use std::{sync::Arc, sync::Mutex};
11
12#[tokio::main]
13async fn main() {
14 tracing_subscriber::fmt()
16 .with_max_level(tracing::Level::INFO)
17 .init();
18 let opts = JetstreamOptions::builder()
20 .wanted_collections(vec!["app.bsky.feed.post".to_string()])
22 .build();
23 let jetstream = JetstreamConnection::new(opts);
25
26 let mut ingestors = Ingestors::new();
28
29 ingestors.commits.insert(
31 "app.bsky.feed.post".to_string(),
33 Box::new(PostIngestor),
34 );
35
36 let cursor: Arc<Mutex<Option<u64>>> = Arc::new(Mutex::new(None));
42
43 let msg_rx = jetstream.get_msg_rx();
45 let reconnect_tx = jetstream.get_reconnect_tx();
46
47 let c_cursor = cursor.clone();
50 tokio::spawn(async move {
51 while let Ok(message) = msg_rx.recv().await {
52 if let Err(e) =
53 handler::handle_message(message, &ingestors, reconnect_tx.clone(), c_cursor.clone())
54 .await
55 {
56 eprintln!("Error processing message: {}", e);
57 };
58 }
59 });
60
61 if let Err(e) = jetstream.connect(cursor.clone()).await {
64 eprintln!("Failed to connect to Jetstream: {}", e);
65 std::process::exit(1);
66 }
67}
68
69pub struct PostIngestor;
70
71#[async_trait]
73impl LexiconIngestor for PostIngestor {
74 async fn ingest(&self, message: Event<Value>) -> anyhow::Result<()> {
75 if let Some(Commit {
76 record: Some(record),
77 ..
78 }) = message.commit
79 {
80 if let Some(Value::String(text)) = record.get("text") {
81 println!("{text:?}");
82 }
83 }
84 Ok(())
85 }
86}