diff --git a/Cargo.toml b/Cargo.toml index d8c5830..a7bba95 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -11,12 +11,15 @@ askama = "0.16" base64 = "0.22" chrono = "0.4" cookie = { version = "0.16", features = ["private"] } -jacquard = "0.12" +jacquard = { version = "0.12", features = ["websocket", "zstd"] } jacquard-identity = "0.12" +futures-util = "0.3" rand = "0.8" +sqlx = { version = "0.8", features = ["sqlite", "runtime-tokio"] } reqwest = { version = "0.12", default-features = false, features = ["rustls-tls", "json"] } serde = { version = "1", features = ["derive"] } serde_json = "1" +tokio = { version = "1", features = ["rt-multi-thread", "time"] } tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter"] } diff --git a/src/db.rs b/src/db.rs new file mode 100644 index 0000000..58ce8ae --- /dev/null +++ b/src/db.rs @@ -0,0 +1,282 @@ +use anyhow::Result; +use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; +use sqlx::{Row, SqlitePool}; + +use jacquard::api::community_lexicon::calendar::event::{Event, EventLocationsItem}; +use jacquard::api::community_lexicon::calendar::rsvp::{Rsvp, RsvpStatus}; + +pub async fn init(path: &str) -> Result { + if let Some(parent) = std::path::Path::new(path).parent() { + std::fs::create_dir_all(parent)?; + } + let options = SqliteConnectOptions::new() + .filename(path) + .create_if_missing(true); + let pool = SqlitePoolOptions::new() + .max_connections(5) + .connect_with(options) + .await?; + sqlx::query( + r#" + CREATE TABLE IF NOT EXISTS events ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + cid TEXT, + name TEXT NOT NULL, + description TEXT, + starts_at TEXT, + ends_at TEXT, + mode TEXT, + status TEXT, + location_text TEXT, + link TEXT, + created_at TEXT + ); + CREATE TABLE IF NOT EXISTS rsvps ( + uri TEXT PRIMARY KEY, + did TEXT NOT NULL, + rkey TEXT NOT NULL, + event_uri TEXT NOT NULL, + status TEXT NOT NULL + ); + CREATE INDEX IF NOT EXISTS idx_rsvps_event ON rsvps(event_uri); + CREATE INDEX IF NOT EXISTS idx_events_starts ON events(starts_at); + CREATE TABLE IF NOT EXISTS indexer_cursor ( + id INTEGER PRIMARY KEY CHECK (id = 1), + time_us INTEGER NOT NULL + ); + "#, + ) + .execute(&pool) + .await?; + Ok(pool) +} + +pub fn normalize_dt(value: Option<&jacquard::common::types::string::Datetime>) -> Option { + let raw = value?.as_str(); + chrono::DateTime::parse_from_rfc3339(raw) + .ok() + .map(|dt| dt.to_utc().format("%Y-%m-%dT%H:%M:%SZ").to_string()) + .or_else(|| Some(raw.to_string())) +} + +fn event_parts(event: &Event) -> (Option, Option) { + let mut location = None; + let mut link = None; + if let Some(items) = &event.locations { + for item in items { + match item { + EventLocationsItem::Address(address) => { + location = address.name.as_ref().map(|n| n.to_string()); + } + EventLocationsItem::Uri(uri) => { + link = Some(uri.uri.as_str().to_string()); + } + _ => {} + } + } + } + (location, link) +} + +pub async fn upsert_event( + pool: &SqlitePool, + uri: &str, + did: &str, + rkey: &str, + cid: Option<&str>, + event: &Event, +) -> Result<()> { + let (location, link) = event_parts(event); + sqlx::query( + r#" + INSERT INTO events (uri, did, rkey, cid, name, description, starts_at, ends_at, mode, status, location_text, link, created_at) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13) + ON CONFLICT(uri) DO UPDATE SET + cid = excluded.cid, + name = excluded.name, + description = excluded.description, + starts_at = excluded.starts_at, + ends_at = excluded.ends_at, + mode = excluded.mode, + status = excluded.status, + location_text = excluded.location_text, + link = excluded.link + "#, + ) + .bind(uri) + .bind(did) + .bind(rkey) + .bind(cid) + .bind(event.name.to_string()) + .bind(event.description.as_ref().map(|d| d.to_string())) + .bind(normalize_dt(event.starts_at.as_ref())) + .bind(normalize_dt(event.ends_at.as_ref())) + .bind(event.mode.as_ref().map(|m| short_token(m.as_str()))) + .bind(event.status.as_ref().map(|s| short_token(s.as_str()))) + .bind(location) + .bind(link) + .bind(normalize_dt(Some(&event.created_at))) + .execute(pool) + .await?; + Ok(()) +} + +pub async fn delete_event(pool: &SqlitePool, uri: &str) -> Result<()> { + sqlx::query("DELETE FROM events WHERE uri = ?1") + .bind(uri) + .execute(pool) + .await?; + Ok(()) +} + +pub fn short_token(value: &str) -> String { + value.rsplit('#').next().unwrap_or(value).to_string() +} + +pub fn rsvp_status_str(status: &RsvpStatus) -> String { + short_token(status.as_str()) +} + +pub async fn upsert_rsvp( + pool: &SqlitePool, + uri: &str, + did: &str, + rkey: &str, + rsvp: &Rsvp, +) -> Result<()> { + sqlx::query( + r#" + INSERT INTO rsvps (uri, did, rkey, event_uri, status) + VALUES (?1, ?2, ?3, ?4, ?5) + ON CONFLICT(uri) DO UPDATE SET status = excluded.status + "#, + ) + .bind(uri) + .bind(did) + .bind(rkey) + .bind(rsvp.subject.uri.as_str()) + .bind(rsvp_status_str(&rsvp.status)) + .execute(pool) + .await?; + Ok(()) +} + +pub async fn delete_rsvp(pool: &SqlitePool, uri: &str) -> Result<()> { + sqlx::query("DELETE FROM rsvps WHERE uri = ?1") + .bind(uri) + .execute(pool) + .await?; + Ok(()) +} + +pub struct GuestRow { + pub did: String, + pub status: String, +} + +pub async fn guests_for_event(pool: &SqlitePool, event_uri: &str) -> Result> { + let rows = sqlx::query("SELECT did, status FROM rsvps WHERE event_uri = ?1 ORDER BY uri") + .bind(event_uri) + .fetch_all(pool) + .await?; + Ok(rows + .into_iter() + .map(|row| GuestRow { + did: row.get("did"), + status: row.get("status"), + }) + .collect()) +} + +pub struct FeedRow { + pub did: String, + pub rkey: String, + pub name: String, + pub starts_at: Option, + pub location_text: Option, + pub going: i64, + pub interested: i64, +} + +pub async fn upcoming_events(pool: &SqlitePool, limit: i64) -> Result> { + let now = chrono::Utc::now().format("%Y-%m-%dT%H:%M:%SZ").to_string(); + let rows = sqlx::query( + r#" + SELECT e.uri, e.did, e.rkey, e.name, e.starts_at, e.location_text, + COALESCE(SUM(CASE WHEN r.status = 'going' THEN 1 ELSE 0 END), 0) AS going, + COALESCE(SUM(CASE WHEN r.status = 'interested' THEN 1 ELSE 0 END), 0) AS interested + FROM events e + LEFT JOIN rsvps r ON r.event_uri = e.uri + WHERE (e.status IS NULL OR e.status != 'cancelled') + AND (e.starts_at IS NULL OR e.starts_at >= ?1) + GROUP BY e.uri + ORDER BY e.starts_at IS NULL, e.starts_at + LIMIT ?2 + "#, + ) + .bind(now) + .bind(limit) + .fetch_all(pool) + .await?; + Ok(rows + .into_iter() + .map(|row| FeedRow { + did: row.get("did"), + rkey: row.get("rkey"), + name: row.get("name"), + starts_at: row.get("starts_at"), + location_text: row.get("location_text"), + going: row.get("going"), + interested: row.get("interested"), + }) + .collect()) +} + +pub async fn counts_for_events( + pool: &SqlitePool, + did: &str, +) -> Result> { + let rows = sqlx::query( + r#" + SELECT r.event_uri, + COALESCE(SUM(CASE WHEN r.status = 'going' THEN 1 ELSE 0 END), 0) AS going, + COALESCE(SUM(CASE WHEN r.status = 'interested' THEN 1 ELSE 0 END), 0) AS interested + FROM rsvps r + JOIN events e ON e.uri = r.event_uri + WHERE e.did = ?1 + GROUP BY r.event_uri + "#, + ) + .bind(did) + .fetch_all(pool) + .await?; + Ok(rows + .into_iter() + .map(|row| { + ( + row.get::("event_uri"), + (row.get("going"), row.get("interested")), + ) + }) + .collect()) +} + +pub async fn load_cursor(pool: &SqlitePool) -> Result> { + let row = sqlx::query("SELECT time_us FROM indexer_cursor WHERE id = 1") + .fetch_optional(pool) + .await?; + Ok(row.map(|r| r.get("time_us"))) +} + +pub async fn store_cursor(pool: &SqlitePool, time_us: i64) -> Result<()> { + sqlx::query( + "INSERT INTO indexer_cursor (id, time_us) VALUES (1, ?1) + ON CONFLICT(id) DO UPDATE SET time_us = excluded.time_us", + ) + .bind(time_us) + .execute(pool) + .await?; + Ok(()) +} diff --git a/src/events.rs b/src/events.rs index 80678ab..223b3c2 100644 --- a/src/events.rs +++ b/src/events.rs @@ -2,7 +2,7 @@ use actix_web::{http::header, web, HttpRequest, HttpResponse}; use anyhow::{anyhow, Context, Result}; use askama::Template; use chrono::TimeZone; -use jacquard::api::app_bsky::actor::get_profile::GetProfile; +use jacquard::api::app_bsky::actor::get_profiles::GetProfiles; use jacquard::api::com_atproto::repo::list_records::ListRecords; use jacquard::api::com_atproto::repo::strong_ref::StrongRef; use jacquard::api::community_lexicon::calendar::event::{self, Event, EventLocationsItem, Mode}; @@ -18,9 +18,10 @@ use jacquard::types::value::from_data; use jacquard::xrpc::{XrpcClient, XrpcExt}; use serde::Deserialize; +use crate::db; use crate::oauth::{session_from_request, AppSession}; use crate::state::AppState; -use crate::templates::{EventFormTemplate, EventTemplate, HomeTemplate}; +use crate::templates::{EventFormTemplate, EventTemplate, GuestView, HomeTemplate}; const STYLE_CSS: &str = include_str!("../static/style.css"); @@ -103,7 +104,21 @@ pub async fn home(req: HttpRequest, state: web::Data) -> HttpResponse Some((session_did, _)) => { let did_string = session_did.to_string(); match list_records::(&agent, &did_string).await { - Ok(records) => { + Ok(mut records) => { + if let Ok(counts) = db::counts_for_events(&state.db, &did_string).await { + for record in &mut records { + let uri = format!( + "at://{}/{}/{}", + did_string, + ::NSID, + record.rkey + ); + if let Some((going, interested)) = counts.get(&uri) { + record.going = *going; + record.interested = *interested; + } + } + } events = records; did = Some(did_string); } @@ -114,10 +129,13 @@ pub async fn home(req: HttpRequest, state: web::Data) -> HttpResponse } } + let feed = db::upcoming_events(&state.db, 25).await.unwrap_or_default(); + render(&HomeTemplate { logged_in, did, events, + feed, }) } @@ -228,17 +246,21 @@ pub async fn event_create( let Some(session) = session_from_request(&req, &state).await else { return redirect("/login"); }; - match create_event_inner(session, form.into_inner()).await { + match create_event_inner(session, form.into_inner(), &state.db).await { Ok(url) => redirect(&url), Err(err) => err_page(&err), } } -async fn create_event_inner(session: AppSession, form: EventForm) -> Result { +async fn create_event_inner( + session: AppSession, + form: EventForm, + pool: &sqlx::SqlitePool, +) -> Result { let event = event_from_form(&form, Datetime::now())?; let agent = Agent::from(session); let output = agent - .create_record(event, None) + .create_record(event.clone(), None) .await .context("failed to create event record")?; let rkey = output @@ -246,7 +268,20 @@ async fn create_event_inner(session: AppSession, form: EventForm) -> Result Result { - let (uri, (event, _cid)) = fetch_event(&authority, &rkey).await?; + let (uri, (event, cid)) = fetch_event(&authority, &rkey).await?; + // read-through indexing: keep the local view of this event fresh + if let Err(err) = db::upsert_event( + &state.db, + uri.as_str(), + &authority, + &rkey, + Some(cid.as_str()), + &event, + ) + .await + { + tracing::warn!("failed to index event {uri}: {err:#}"); + } let session = session_from_request(&req, &state).await; let logged_in = session.is_some(); @@ -287,8 +335,13 @@ async fn event_page_inner( .await .unwrap_or_else(|| authority.clone()); + let guest_rows = db::guests_for_event(&state.db, uri.as_str()) + .await + .unwrap_or_default(); + let guests = fetch_guest_views(&guest_rows).await; + Ok(render(&EventTemplate::from_event( - &event, &authority, &rkey, &host, is_owner, logged_in, my_rsvp, + &event, &authority, &rkey, &host, is_owner, logged_in, my_rsvp, guests, ))) } @@ -346,7 +399,7 @@ pub async fn event_update( let Some(session) = session_from_request(&req, &state).await else { return redirect("/login"); }; - match update_event_inner(session, &authority, &rkey, form.into_inner()).await { + match update_event_inner(session, &authority, &rkey, form.into_inner(), &state.db).await { Ok(()) => redirect(&format!("/e/{authority}/{rkey}")), Err(err) => err_page(&err), } @@ -357,6 +410,7 @@ async fn update_event_inner( authority: &str, rkey: &str, form: EventForm, + pool: &sqlx::SqlitePool, ) -> Result<()> { let agent = Agent::from(session); let owns = matches!(agent.session_info().await, Some((did, _)) if did.to_string() == authority); @@ -366,11 +420,29 @@ async fn update_event_inner( let (_uri, (existing, _cid)) = fetch_event(authority, rkey).await?; let mut event = event_from_form(&form, existing.created_at.clone())?; event.status = existing.status.clone(); - let rkey = Rkey::new(rkey).map_err(|err| anyhow!("invalid rkey: {err}"))?; - agent - .put_record(RecordKey(rkey.into_static()), event) + let rkey_parsed = Rkey::new(rkey).map_err(|err| anyhow!("invalid rkey: {err}"))?; + let output = agent + .put_record(RecordKey(rkey_parsed.into_static()), event.clone()) .await .context("failed to update event record")?; + let uri = format!( + "at://{}/{}/{}", + authority, + ::NSID, + rkey + ); + if let Err(err) = db::upsert_event( + pool, + &uri, + authority, + rkey, + Some(output.cid.as_str()), + &event, + ) + .await + { + tracing::warn!("failed to reindex event {uri}: {err:#}"); + } Ok(()) } @@ -414,13 +486,19 @@ pub async fn rsvp_create( let Some(session) = session_from_request(&req, &state).await else { return redirect("/login"); }; - match rsvp_inner(session, &authority, &rkey, &form.status).await { + match rsvp_inner(session, &authority, &rkey, &form.status, &state.db).await { Ok(()) => redirect(&format!("/e/{authority}/{rkey}")), Err(err) => err_page(&err), } } -async fn rsvp_inner(session: AppSession, authority: &str, rkey: &str, status: &str) -> Result<()> { +async fn rsvp_inner( + session: AppSession, + authority: &str, + rkey: &str, + status: &str, + pool: &sqlx::SqlitePool, +) -> Result<()> { let status = rsvp_status(status)?; let (uri, (_event, cid)) = fetch_event(authority, rkey).await?; let agent = Agent::from(session); @@ -444,21 +522,35 @@ async fn rsvp_inner(session: AppSession, authority: &str, rkey: &str, status: &s .into_iter() .find(|record| record.subject_uri == uri.as_str()); - match existing { + let rsvp_uri = match existing { Some(record) => { let existing_rkey = Rkey::new(record.rkey.as_str()).map_err(|err| anyhow!("invalid rkey: {err}"))?; - agent - .put_record(RecordKey(existing_rkey.into_static()), rsvp) + let output = agent + .put_record(RecordKey(existing_rkey.into_static()), rsvp.clone()) .await .context("failed to update rsvp")?; + output.uri.as_str().to_string() } None => { - agent - .create_record(rsvp, None) + let output = agent + .create_record(rsvp.clone(), None) .await .context("failed to create rsvp")?; + output.uri.as_str().to_string() } + }; + + if let Err(err) = db::upsert_rsvp( + pool, + &rsvp_uri, + &did_string, + rkey_from_uri(&rsvp_uri)?, + &rsvp, + ) + .await + { + tracing::warn!("failed to index rsvp {rsvp_uri}: {err:#}"); } Ok(()) } @@ -472,13 +564,18 @@ pub async fn rsvp_delete( let Some(session) = session_from_request(&req, &state).await else { return redirect("/login"); }; - match rsvp_delete_inner(session, &authority, &rkey).await { + match rsvp_delete_inner(session, &authority, &rkey, &state.db).await { Ok(()) => redirect(&format!("/e/{authority}/{rkey}")), Err(err) => err_page(&err), } } -async fn rsvp_delete_inner(session: AppSession, authority: &str, rkey: &str) -> Result<()> { +async fn rsvp_delete_inner( + session: AppSession, + authority: &str, + rkey: &str, + pool: &sqlx::SqlitePool, +) -> Result<()> { let uri = event_at_uri(authority, rkey)?; let agent = Agent::from(session); let Some((did, _)) = agent.session_info().await else { @@ -495,10 +592,26 @@ async fn rsvp_delete_inner(session: AppSession, authority: &str, rkey: &str) -> .delete_record::(RecordKey(existing_rkey.into_static())) .await .context("failed to delete rsvp")?; + let rsvp_uri = format!( + "at://{}/{}/{}", + did, + ::NSID, + record.rkey + ); + if let Err(err) = db::delete_rsvp(pool, &rsvp_uri).await { + tracing::warn!("failed to unindex rsvp {rsvp_uri}: {err:#}"); + } } Ok(()) } +fn rkey_from_uri(uri: &str) -> Result<&str> { + uri.rsplit('/') + .next() + .filter(|part| !part.is_empty()) + .ok_or_else(|| anyhow!("uri `{uri}` has no rkey")) +} + // -- helpers ------------------------------------------------------------------- pub struct ListedRecord { @@ -507,6 +620,8 @@ pub struct ListedRecord { pub starts_at: String, pub subject_uri: String, pub status: String, + pub going: i64, + pub interested: i64, } async fn list_records(agent: &Agent, did: &str) -> Result> @@ -554,6 +669,8 @@ impl ListedRecordExt for Event { .unwrap_or_default(), subject_uri: String::new(), status: String::new(), + going: 0, + interested: 0, } } } @@ -571,6 +688,8 @@ impl ListedRecordExt for Rsvp { RsvpStatus::Notgoing => "not going".to_string(), RsvpStatus::Other(other) => other.to_string(), }, + going: 0, + interested: 0, } } } @@ -604,20 +723,71 @@ async fn fetch_handle(authority: &str) -> Option { if !authority.starts_with("did:") { return Some(format!("@{authority}")); } + let dids = [authority.to_string()]; + let profiles = fetch_profiles(&dids).await; + profiles + .into_iter() + .next() + .map(|profile| format!("@{}", profile.handle)) +} + +pub struct ProfileInfo { + pub did: String, + pub handle: String, + pub display_name: Option, + pub avatar: Option, +} + +async fn fetch_profiles(dids: &[String]) -> Vec { let http = reqwest::Client::new(); - let base = jacquard::common::deps::fluent_uri::Uri::parse("https://public.api.bsky.app") - .ok()? - .to_owned(); - let did = Did::new(authority).ok()?; - let request = GetProfile::new() - .actor(AtIdentifier::Did(did.into_static())) - .build(); - let output = http - .xrpc(base.borrow()) - .send(&request) - .await - .ok()? - .into_output() - .ok()?; - Some(format!("@{}", output.value.handle.as_str())) + let base = match jacquard::common::deps::fluent_uri::Uri::parse("https://public.api.bsky.app") { + Ok(base) => base.to_owned(), + Err(_) => return Vec::new(), + }; + + let mut out = Vec::new(); + for chunk in dids.chunks(25) { + let actors: Vec = chunk + .iter() + .filter_map(|did| Did::new(did.as_str()).ok()) + .map(|did| AtIdentifier::Did(did.into_static())) + .collect(); + if actors.is_empty() { + continue; + } + let request = GetProfiles::new().actors(actors).build(); + let Ok(response) = http.xrpc(base.borrow()).send(&request).await else { + continue; + }; + let Ok(output) = response.into_output() else { + continue; + }; + for profile in output.profiles { + out.push(ProfileInfo { + did: profile.did.to_string(), + handle: profile.handle.to_string(), + display_name: profile.display_name.as_ref().map(|n| n.to_string()), + avatar: profile.avatar.as_ref().map(|u| u.as_str().to_string()), + }); + } + } + out +} + +async fn fetch_guest_views(rows: &[db::GuestRow]) -> Vec { + let dids: Vec = rows.iter().map(|row| row.did.clone()).collect(); + let profiles = fetch_profiles(&dids).await; + let mut views = Vec::new(); + for row in rows { + let profile = profiles.iter().find(|p| p.did == row.did); + views.push(GuestView { + name: profile + .and_then(|p| p.display_name.clone()) + .or_else(|| profile.map(|p| format!("@{}", p.handle))) + .unwrap_or_else(|| row.did.clone()), + avatar: profile.and_then(|p| p.avatar.clone()), + status: row.status.clone(), + }); + } + views } diff --git a/src/indexer.rs b/src/indexer.rs new file mode 100644 index 0000000..cb4b586 --- /dev/null +++ b/src/indexer.rs @@ -0,0 +1,186 @@ +use std::time::Duration; + +use anyhow::Result; +use futures_util::StreamExt; +use jacquard::api::community_lexicon::calendar::event::Event; +use jacquard::api::community_lexicon::calendar::rsvp::Rsvp; +use jacquard::common::deps::fluent_uri::Uri; +use jacquard::common::jetstream::{CommitOperation, JetstreamMessage, JetstreamParams}; +use jacquard::common::types::collection::Collection; +use jacquard::common::xrpc::{SubscriptionClient, TungsteniteSubscriptionClient}; +use jacquard::common::IntoStatic; +use jacquard::types::value::from_data; +use sqlx::SqlitePool; +use tracing::{info, warn}; + +use crate::db; + +const ENDPOINT: &str = "wss://jetstream2.us-east.bsky.network"; +const UFOS: &str = "https://ufos-api.microcosm.blue/records"; + +pub async fn run(pool: SqlitePool) { + if let Err(err) = backfill(&pool).await { + warn!("ufos backfill failed: {err:#}"); + } + loop { + match run_once(&pool).await { + Ok(()) => { + info!("jetstream connection ended, reconnecting"); + } + Err(err) => { + warn!("jetstream error: {err:#}; reconnecting in 5s"); + } + } + tokio::time::sleep(Duration::from_secs(5)).await; + } +} + +async fn run_once(pool: &SqlitePool) -> Result<()> { + let cursor = db::load_cursor(pool).await?; + let base_url = Uri::parse(ENDPOINT.to_string()).map_err(|(err, _)| anyhow::anyhow!(err))?; + + let collections = vec![ + ::nsid().into_static(), + ::nsid().into_static(), + ]; + let params = match cursor { + Some(cursor) => JetstreamParams::new() + .wanted_collections(collections) + .cursor(cursor) + .compress(true) + .build(), + None => JetstreamParams::new() + .wanted_collections(collections) + .compress(true) + .build(), + }; + + let client = TungsteniteSubscriptionClient::from_base_uri(base_url); + let stream = client + .subscribe(¶ms) + .await + .map_err(|err| anyhow::anyhow!("subscribe failed: {err}"))?; + let (_sink, mut messages) = stream.into_stream(); + + info!(cursor, "jetstream indexer connected"); + + let mut last_flush = std::time::Instant::now(); + let mut latest_cursor: Option = None; + + while let Some(result) = messages.next().await { + let message = result.map_err(|err| anyhow::anyhow!("stream error: {err}"))?; + if let JetstreamMessage::Commit { + did, + time_us, + commit, + } = message + { + latest_cursor = Some(time_us); + let uri = format!("at://{}/{}/{}", did, commit.collection, commit.rkey); + match commit.operation { + CommitOperation::Create | CommitOperation::Update => { + if let Some(record) = &commit.record { + if commit.collection.as_str() == ::NSID { + let parsed: Result = from_data(record); + match parsed { + Ok(event) => { + db::upsert_event( + pool, + &uri, + did.as_str(), + commit.rkey.as_str(), + commit.cid.as_ref().map(|c| c.as_str()), + &event, + ) + .await?; + } + Err(err) => warn!("invalid event record {uri}: {err}"), + } + } else if commit.collection.as_str() == ::NSID { + let parsed: Result = from_data(record); + match parsed { + Ok(rsvp) => { + db::upsert_rsvp( + pool, + &uri, + did.as_str(), + commit.rkey.as_str(), + &rsvp, + ) + .await?; + } + Err(err) => warn!("invalid rsvp record {uri}: {err}"), + } + } + } + } + CommitOperation::Delete => { + if commit.collection.as_str() == ::NSID { + db::delete_event(pool, &uri).await?; + } else if commit.collection.as_str() == ::NSID { + db::delete_rsvp(pool, &uri).await?; + } + } + } + + if last_flush.elapsed() > Duration::from_secs(10) { + if let Some(cursor) = latest_cursor { + db::store_cursor(pool, cursor).await?; + } + last_flush = std::time::Instant::now(); + } + } + } + + if let Some(cursor) = latest_cursor { + db::store_cursor(pool, cursor).await?; + } + Ok(()) +} + +/// Backfill the index with every known record for our collections. +/// ufos (microcosm) tracks the whole network for arbitrary lexicons, which +/// suits these young collections better than a relay crawl. +async fn backfill(pool: &SqlitePool) -> Result<()> { + let http = reqwest::Client::new(); + for collection in [::NSID, ::NSID] { + let rows: Vec = http + .get(UFOS) + .query(&[("collection", collection)]) + .send() + .await? + .error_for_status()? + .json() + .await?; + let mut count = 0u32; + for row in rows { + let uri = format!("at://{}/{}/{}", row.did, collection, row.rkey); + if collection == ::NSID { + match serde_json::from_value::(row.record) { + Ok(event) => { + db::upsert_event(pool, &uri, &row.did, &row.rkey, None, &event).await?; + count += 1; + } + Err(err) => warn!("backfill: invalid event {uri}: {err}"), + } + } else { + match serde_json::from_value::(row.record) { + Ok(rsvp) => { + db::upsert_rsvp(pool, &uri, &row.did, &row.rkey, &rsvp).await?; + count += 1; + } + Err(err) => warn!("backfill: invalid rsvp {uri}: {err}"), + } + } + } + info!(collection, count, "ufos backfill complete"); + } + Ok(()) +} + +#[derive(serde::Deserialize)] +struct UfosRecord { + did: String, + rkey: String, + record: serde_json::Value, +} diff --git a/src/main.rs b/src/main.rs index 9d238d7..c5c0739 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,4 +1,6 @@ +mod db; mod events; +mod indexer; mod oauth; mod state; mod templates; @@ -14,8 +16,12 @@ async fn main() -> std::io::Result<()> { .init(); let config = Config::from_env(); + let db_path = config.data_dir.join("partytool.db"); + let pool = db::init(&db_path.to_string_lossy()) + .await + .map_err(|err| std::io::Error::new(std::io::ErrorKind::Other, err))?; let state = web::Data::new( - AppState::new(&config) + AppState::new(&config, pool.clone()) .map_err(|err| std::io::Error::new(std::io::ErrorKind::Other, err))?, ); let listen = config.listen.clone(); @@ -24,6 +30,8 @@ async fn main() -> std::io::Result<()> { state.base_url ); + tokio::spawn(indexer::run(pool)); + HttpServer::new(move || { App::new() .app_data(state.clone()) diff --git a/src/state.rs b/src/state.rs index d17b37e..1b284fe 100644 --- a/src/state.rs +++ b/src/state.rs @@ -9,6 +9,7 @@ use jacquard::oauth::{ session::ClientData, }; use jacquard_identity::PublicResolver; +use sqlx::SqlitePool; pub type AppOAuthClient = OAuthClient; @@ -16,6 +17,7 @@ pub struct AppState { pub oauth: Arc, pub cookie_key: Key, pub base_url: String, + pub db: SqlitePool, } pub struct Config { @@ -40,7 +42,7 @@ impl Config { } impl AppState { - pub fn new(config: &Config) -> Result { + pub fn new(config: &Config, db: SqlitePool) -> Result { fs::create_dir_all(&config.data_dir) .with_context(|| format!("failed to create {}", config.data_dir.display()))?; @@ -90,6 +92,7 @@ impl AppState { oauth: Arc::new(oauth), cookie_key, base_url: config.base_url.clone(), + db, }) } } diff --git a/src/templates.rs b/src/templates.rs index 12107f5..5178069 100644 --- a/src/templates.rs +++ b/src/templates.rs @@ -1,6 +1,7 @@ use askama::Template; use jacquard::api::community_lexicon::calendar::event::{Event, EventLocationsItem}; +use crate::db::FeedRow; use crate::events::ListedRecord; #[derive(Template)] @@ -9,6 +10,7 @@ pub struct HomeTemplate { pub logged_in: bool, pub did: Option, pub events: Vec, + pub feed: Vec, } #[derive(Template)] @@ -39,6 +41,40 @@ pub struct EventView { pub link: String, } +pub struct GuestView { + pub name: String, + pub avatar: Option, + pub status: String, +} + +pub struct GuestGroups { + pub going: Vec, + pub interested: Vec, + pub notgoing: Vec, +} + +impl GuestGroups { + pub fn from_guests(guests: Vec) -> Self { + let mut groups = Self { + going: Vec::new(), + interested: Vec::new(), + notgoing: Vec::new(), + }; + for guest in guests { + match guest.status.as_str() { + "going" => groups.going.push(guest), + "interested" => groups.interested.push(guest), + _ => groups.notgoing.push(guest), + } + } + groups + } + + pub fn is_empty(&self) -> bool { + self.going.is_empty() && self.interested.is_empty() && self.notgoing.is_empty() + } +} + #[derive(Template)] #[template(path = "event.html")] pub struct EventTemplate { @@ -49,6 +85,7 @@ pub struct EventTemplate { pub is_owner: bool, pub logged_in: bool, pub my_rsvp: Option, + pub guests: GuestGroups, } impl EventTemplate { @@ -60,6 +97,7 @@ impl EventTemplate { is_owner: bool, logged_in: bool, my_rsvp: Option, + guests: Vec, ) -> Self { let mut location = String::new(); let mut link = String::new(); @@ -121,6 +159,7 @@ impl EventTemplate { is_owner, logged_in, my_rsvp, + guests: GuestGroups::from_guests(guests), } } } diff --git a/static/style.css b/static/style.css index 9120bb1..2f3507a 100644 --- a/static/style.css +++ b/static/style.css @@ -96,6 +96,22 @@ h2 { font-size: 1.1rem; color: var(--muted); text-transform: lowercase; letter-s .card.narrow { max-width: 420px; margin: 3rem auto; } .card-title { font-weight: 700; font-size: 1.15rem; display: block; } .card-sub { color: var(--muted); font-size: 0.9rem; } +.card-counts { display: block; margin-top: 0.35rem; color: var(--accent-2); font-size: 0.85rem; font-weight: 600; } + +.guest-heading { font-size: 0.95rem; color: var(--muted); margin: 1rem 0 0.5rem; } +.guest-grid { display: flex; flex-wrap: wrap; gap: 0.9rem; } +.guest { display: flex; flex-direction: column; align-items: center; width: 4.5rem; text-align: center; } +.avatar { + width: 3rem; height: 3rem; + border-radius: 50%; + object-fit: cover; + background: rgba(255, 255, 255, 0.1); +} +.avatar-fallback { + display: inline-block; + background: linear-gradient(135deg, var(--accent), var(--accent-2)); +} +.guest-name { font-size: 0.75rem; margin-top: 0.3rem; color: var(--muted); overflow-wrap: anywhere; } .cta { text-align: center; padding: 2rem 0; } diff --git a/templates/event.html b/templates/event.html index b400b29..573ff9f 100644 --- a/templates/event.html +++ b/templates/event.html @@ -71,6 +71,57 @@

log in to rsvp.

{% endif %} + + {% if !guests.is_empty() %} +
+

guest list

+ {% if !guests.going.is_empty() %} +

going · {{ guests.going.len() }}

+
+ {% for guest in &guests.going %} +
+ {% if let Some(avatar) = &guest.avatar %} + + {% else %} + + {% endif %} + {{ guest.name }} +
+ {% endfor %} +
+ {% endif %} + {% if !guests.interested.is_empty() %} +

interested · {{ guests.interested.len() }}

+
+ {% for guest in &guests.interested %} +
+ {% if let Some(avatar) = &guest.avatar %} + + {% else %} + + {% endif %} + {{ guest.name }} +
+ {% endfor %} +
+ {% endif %} + {% if !guests.notgoing.is_empty() %} +

can't go · {{ guests.notgoing.len() }}

+
+ {% for guest in &guests.notgoing %} +
+ {% if let Some(avatar) = &guest.avatar %} + + {% else %} + + {% endif %} + {{ guest.name }} +
+ {% endfor %} +
+ {% endif %} +
+ {% endif %} diff --git a/templates/home.html b/templates/home.html index 9a9e81a..bb84b25 100644 --- a/templates/home.html +++ b/templates/home.html @@ -38,6 +38,9 @@ {{ event.title }} {{ event.starts_at }} + {% if event.going > 0 || event.interested > 0 %} + {{ event.going }} going{% if event.interested > 0 %} · {{ event.interested }} interested{% endif %} + {% endif %} {% endfor %} @@ -46,6 +49,23 @@ log in with atproto {% endif %} + {% if !feed.is_empty() %} +
+

upcoming on the network

+ {% for event in feed %} + + {{ event.name }} + + {{ event.starts_at.as_deref().unwrap_or("") }} + {% if let Some(location) = &event.location_text %} · {{ location }}{% endif %} + + {% if event.going > 0 || event.interested > 0 %} + {{ event.going }} going{% if event.interested > 0 %} · {{ event.interested }} interested{% endif %} + {% endif %} + + {% endfor %} +
+ {% endif %}