Something went wrong. Try again.
Bluesky -> traQ
Something went wrong. Try again.
12 kB · 326 lines
Rust
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327use std::sync::Arc;use std::time::{Duration, Instant};
use atproto_jetstream::{Consumer, ConsumerTaskConfig, EventHandler, JetstreamEvent};use tokio::sync::{Mutex, Notify};use tokio::time::sleep;use tracing::{error, info, warn};
use crate::app_config::config::Config;use crate::database::DbPool;use crate::model::bluesky_types::build_at_proto_uri;use crate::model::post_event::{PostCreateEvent, QueuedEventType};use crate::repository::queued_event;use crate::repository::system_state;use crate::repository::user;use crate::service::bluesky_client::BlueskyClient;
/// Debounce interval for persisting the Jetstream cursor./// Replays within this window are masked by the duplicate-post check.const CURSOR_SAVE_DEBOUNCE: Duration = Duration::from_secs(30);const NO_DIDS_RETRY: Duration = Duration::from_secs(30);const MAX_BACKOFF: Duration = Duration::from_secs(300);/// Connections lasting at least this long are considered healthy, so the/// reconnect backoff counter resets.const STABLE_CONNECTION: Duration = Duration::from_secs(60);
pub async fn start( endpoints: &[&str], initial_dids: &[String], initial_cursor: Option<i64>, pool: &DbPool, config: &Config, cancel_token: atproto_jetstream::CancellationToken, notify: Arc<Notify>,) -> anyhow::Result<()> { let mut backoff_attempt: u32 = 0; // Fall back to the startup snapshot until the first DB refresh succeeds. let mut dids = initial_dids.to_vec(); let mut cursor = initial_cursor; let mut dids_loaded = false; let cursor_saver = Arc::new(CursorSaver::new(pool.clone()));
loop { if cancel_token.is_cancelled() { info!("Jetstream shutting down..."); if let Err(e) = cursor_saver.flush().await { error!("Failed to flush Jetstream cursor on shutdown: {}", e); } return Ok(()); }
match user::get_all_dids(pool).await { Ok(fresh) => { dids = fresh; dids_loaded = true; } Err(e) => { error!("Failed to load DIDs for Jetstream: {}", e); if !dids_loaded { // Keep the startup snapshot on first failure; otherwise reuse last known. } } }
match system_state::get_jetstream_cursor(pool).await { Ok(fresh) => cursor = fresh, Err(e) => { error!("Failed to load Jetstream cursor: {}", e); } }
if dids.is_empty() { warn!("No DIDs to subscribe. Retrying..."); tokio::select! { _ = cancel_token.cancelled() => { if let Err(e) = cursor_saver.flush().await { error!("Failed to flush Jetstream cursor on shutdown: {}", e); } return Ok(()); } _ = sleep(NO_DIDS_RETRY) => {} } continue; }
let bluesky_client = match BlueskyClient::new(config).await { Ok(c) => c, Err(e) => { backoff_attempt += 1; let backoff = calculate_backoff(backoff_attempt); error!( "Failed to create Bluesky client for Jetstream: {}. Retrying in {:?}", e, backoff ); tokio::select! { _ = cancel_token.cancelled() => { if let Err(e) = cursor_saver.flush().await { error!("Failed to flush Jetstream cursor on shutdown: {}", e); } return Ok(()); } _ = sleep(backoff) => {} } continue; } };
for endpoint in endpoints { if cancel_token.is_cancelled() { if let Err(e) = cursor_saver.flush().await { error!("Failed to flush Jetstream cursor on shutdown: {}", e); } return Ok(()); }
info!("Connecting to Jetstream endpoint: {}", endpoint);
let task_config = ConsumerTaskConfig { user_agent: format!("qonstellation/{}", env!("CARGO_PKG_VERSION")), compression: false, zstd_dictionary_location: String::new(), jetstream_hostname: endpoint.to_string(), collections: vec!["app.bsky.feed.post".to_string()], dids: dids.clone(), max_message_size_bytes: None, cursor, require_hello: false, };
let consumer = Consumer::new(task_config); let handler = Arc::new(JetstreamEventHandler { pool: pool.clone(), bluesky_client: bluesky_client.clone(), notify: notify.clone(), cursor_saver: cursor_saver.clone(), });
if let Err(e) = consumer.register_handler(handler).await { error!("Failed to register handler: {}", e); continue; }
let token = cancel_token.clone(); let connected_at = Instant::now(); match consumer.run_background(token).await { Ok(_) if cancel_token.is_cancelled() => { info!("Jetstream consumer exited gracefully"); if let Err(e) = cursor_saver.flush().await { error!("Failed to flush Jetstream cursor on shutdown: {}", e); } return Ok(()); } Ok(_) => { warn!( "Jetstream endpoint {} closed the connection, reconnecting...", endpoint ); } Err(e) => { error!("Jetstream endpoint {} failed: {}", endpoint, e); } }
// Save the debounced cursor so a reconnect replays as little as // possible, and reuse the flushed value for the next endpoint. match cursor_saver.flush().await { Ok(Some(fresh)) => cursor = Some(fresh), Ok(None) => {} Err(e) => { error!("Failed to flush Jetstream cursor before reconnect: {}", e) } } if connected_at.elapsed() >= STABLE_CONNECTION { backoff_attempt = 0; } }
backoff_attempt += 1; let backoff = calculate_backoff(backoff_attempt); error!("All Jetstream endpoints failed. Retrying in {:?}", backoff); tokio::select! { _ = cancel_token.cancelled() => { if let Err(e) = cursor_saver.flush().await { error!("Failed to flush Jetstream cursor on shutdown: {}", e); } return Ok(()); } _ = sleep(backoff) => {} } }}
fn calculate_backoff(attempt: u32) -> Duration { let secs = 2u64 .saturating_pow(attempt.min(8)) .min(MAX_BACKOFF.as_secs()); Duration::from_secs(secs.max(1))}
struct CursorSaver { pool: DbPool, last_write: Mutex<Option<Instant>>, last_cursor: Mutex<Option<i64>>,}
impl CursorSaver { fn new(pool: DbPool) -> Self { Self { pool, last_write: Mutex::new(None), last_cursor: Mutex::new(None), } }
async fn save(&self, cursor: i64) -> anyhow::Result<()> { *self.last_cursor.lock().await = Some(cursor);
let should_write = match *self.last_write.lock().await { None => true, Some(t) => t.elapsed() >= CURSOR_SAVE_DEBOUNCE, };
if should_write { system_state::save_jetstream_cursor(&self.pool, cursor).await?; *self.last_write.lock().await = Some(Instant::now()); }
Ok(()) }
/// Persists the debounced cursor and returns it, so reconnects can reuse /// the freshest value instead of a stale local one. async fn flush(&self) -> anyhow::Result<Option<i64>> { let cursor = *self.last_cursor.lock().await; if let Some(cursor) = cursor { system_state::save_jetstream_cursor(&self.pool, cursor).await?; *self.last_write.lock().await = Some(Instant::now()); }
Ok(cursor) }}
struct JetstreamEventHandler { pool: DbPool, bluesky_client: BlueskyClient, notify: Arc<Notify>, cursor_saver: Arc<CursorSaver>,}
#[async_trait::async_trait]impl EventHandler for JetstreamEventHandler { async fn handle_event(&self, event: Arc<JetstreamEvent>) -> anyhow::Result<()> { match &*event { JetstreamEvent::Commit { did, time_us, commit, .. } => { if commit.operation != "create" || commit.collection != "app.bsky.feed.post" { self.cursor_saver.save(*time_us as i64).await?; return Ok(()); }
let post_event = match PostCreateEvent::from_commit(did, *time_us, commit) { Ok(post_event) => post_event, Err(e) => { warn!("Invalid record: {}", e); self.cursor_saver.save(*time_us as i64).await?; return Ok(()); } }; let at_proto_uri = build_at_proto_uri(did, &commit.rkey);
match self .bluesky_client .is_self_thread(&post_event.record.reply, did) .await { Ok(true) => {} Ok(false) => { warn!( "Skipping post {} because it is not a self thread", at_proto_uri ); self.cursor_saver.save(*time_us as i64).await?; return Ok(()); } Err(e) => { // Transient Bluesky failure: don't advance the cursor so // this event is replayed, and return Err without killing // the connection (the outer loop reconnects on fatal errors). error!("Error checking self thread: {}", e); return Err(e); } }
let event_type = QueuedEventType::Post(post_event); let json = serde_json::to_value(&event_type).map_err(|e| { error!("Error handling Jetstream event: {}", e); anyhow::anyhow!("Failed to serialize event: {}", e) })?; queued_event::add_queued_event(&self.pool, &json) .await .map_err(|e| { error!("Error handling Jetstream event: {}", e); e })?; self.notify.notify_one(); self.cursor_saver.save(*time_us as i64).await?; Ok(()) } JetstreamEvent::Delete { time_us, .. } | JetstreamEvent::Identity { time_us, .. } | JetstreamEvent::Account { time_us, .. } => { self.cursor_saver.save(*time_us as i64).await?; Ok(()) } } }
fn handler_id(&self) -> &str { "qonstellation_jetstream_handler" }}