diff --git a/src/main.rs b/src/main.rs index e73f1bc..db2f308 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,5 +1,10 @@ use clap::Parser; use url::Url; +use jetstream::{ + JetstreamCompression, JetstreamConfig, JetstreamConnector, + events::{CommitOp, Cursor, EventKind}, + exports::Nsid, +}; use jacquard::{ api::app_bsky::{ feed::post::Post, @@ -19,7 +24,6 @@ use jacquard::{ }, }; - type Result = std::result::Result>; #[derive(Debug, Parser)] @@ -34,6 +38,17 @@ struct Args { /// app password for the bot user #[arg(short, long, env = "BOT_APP_PASSWORD")] app_password: String, + /// lightweight firehose + #[arg(short, long, env = "BOT_JETSTREAM_URL")] + #[clap(default_value = "wss://jetstream1.us-east.fire.hose.cam/subscribe")] + jetstream_url: Url, + /// optional: we can pick up from a past jetstream cursor + /// + /// the default is to just live-tail + /// + /// warning: setting this can lead to rapid bot posting + #[arg(long)] + jetstream_cursor: Option, } async fn post_link(client: &AuthenticatedClient, identifier: &AtIdentifier<'_>, pre_text: &str, link: Url) -> Result<()> { @@ -88,31 +103,104 @@ async fn main() -> Result<()> { env_logger::init(); let args = Args::parse(); - - // Create HTTP client + // Create HTTP client and session let pds_uri = args.pds.as_str().trim_end_matches('/').to_string().into(); let mut client = AuthenticatedClient::new(reqwest::Client::new(), pds_uri); - - // Create session + let bot_id = AtIdentifier::new(&args.identifier)?; let session = Session::from( client .send( CreateSession::new() - .identifier(&args.identifier) + .identifier(&bot_id.to_string()) .password(args.app_password) .build(), ) .await? .into_output()?, ); - println!("logged in as {} ({})", session.handle, session.did); client.set_session(session); - let identifier = AtIdentifier::new(&args.identifier)?; + let jetstream_config: JetstreamConfig = JetstreamConfig { + endpoint: args.jetstream_url.to_string(), + wanted_collections: vec![Nsid::new("sh.tangled.label.op".to_string())?], + user_agent: Some("hacktober_bot".to_string()), + compression: JetstreamCompression::Zstd, + replay_on_reconnect: true, + channel_size: 1024, // buffer up to ~1s of jetstream events + ..Default::default() + }; - let u2: Url = "https://bad-example.com".parse()?; - post_link(&client, &identifier, "link test 2: ", u2).await?; + let mut receiver = JetstreamConnector::new(jetstream_config)? + .connect_cursor(args.jetstream_cursor.map(Cursor::from_raw_u64)) + .await?; + + println!("receiving jetstream messages..."); + loop { + let Some(event) = receiver.recv().await else { + eprintln!("consumer: could not receive event, bailing"); + break; + }; + if event.kind != EventKind::Commit { + continue; + } + let Some(ref commit) = event.commit else { + eprintln!("consumer: commit event missing commit data, ignoring"); + continue; + }; + if commit.operation != CommitOp::Create { + continue; + } + let Some(ref record) = commit.record else { + eprintln!("consumer: commit update/delete missing record, ignoring"); + continue; + }; + let jv: serde_json::Value = match record.get().parse() { + Ok(v) => v, + Err(e) => { + eprintln!("consumer: record failed to parse, ignoring: {e}"); + continue; + } + }; + let serde_json::Value::Object(o) = jv else { + eprintln!("record was not an object, ignoring"); + continue; + }; + let Some(serde_json::Value::Array(adds)) = o.get("add") else { + eprintln!("op did not have label added or was not an array"); + continue; + }; + let mut added_good_first_issue = false; + for added in adds { + let serde_json::Value::Object(a) = added else { + eprintln!("added item was not an obj"); + continue; + }; + let Some(serde_json::Value::String(key)) = a.get("key") else { + eprintln!("added was missing key prop"); + continue; + }; + if key == "at://did:plc:wshs7t2adsemcrrd4snkeqli/sh.tangled.label.definition/good-first-issue" { + println!("found a good first issue label!! {:?}", event.cursor); + added_good_first_issue = true; + break; + } + eprintln!("found a label but it wasn't good-first-issue, ignoring..."); + } + if !added_good_first_issue { + continue; + } + let Some(serde_json::Value::String(subject)) = o.get("subject") else { + eprintln!("could not find `subject` string for the good-first-issue label"); + continue; + }; + + eprintln!("got {subject} {o:?}"); + } + + // let u2: Url = "https://bad-example.com".parse()?; + // let bot_id = AtIdentifier::new(&args.identifier)?; + // post_link(&client, &bot_id, "link test 2: ", u2).await?; Ok(()) }