diff --git a/README.md b/README.md index fabfa7e..f6cdc5b 100644 --- a/README.md +++ b/README.md @@ -72,7 +72,6 @@ See the `config.example.yml` file for additional examples. # TODO * use i64, it's fine -* look up keys on startup * possible scoring function for queries * add likes * support deletes diff --git a/migrations/20241103180245_init.down.sql b/migrations/20241103180245_init.down.sql index 8b00173..4b833b6 100644 --- a/migrations/20241103180245_init.down.sql +++ b/migrations/20241103180245_init.down.sql @@ -2,4 +2,4 @@ DROP TABLE feed_content; DROP TABLE consumer_control; - +DROP TABLE verification_method_cache; diff --git a/migrations/20241103180245_init.up.sql b/migrations/20241103180245_init.up.sql index 0034f36..a735894 100644 --- a/migrations/20241103180245_init.up.sql +++ b/migrations/20241103180245_init.up.sql @@ -4,17 +4,16 @@ CREATE TABLE feed_content ( feed_id TEXT NOT NULL, uri TEXT NOT NULL, indexed_at INTEGER NOT NULL, - indexed_at_more INTEGER NOT NULL, cid TEXT NOT NULL, updated_at DATETIME NOT NULL DEFAULT (datetime('now')), PRIMARY KEY (feed_id, uri) ); -CREATE INDEX feed_content_idx_feed ON feed_content(feed_id, indexed_at DESC, indexed_at_more DESC, cid DESC); +CREATE INDEX feed_content_idx_feed ON feed_content(feed_id, indexed_at DESC, cid DESC); CREATE TABLE consumer_control ( source TEXT NOT NULL, - time_us VARCHAR NOT NULL, + time_us INTEGER NOT NULL, updated_at DATETIME NOT NULL DEFAULT (datetime('now')), PRIMARY KEY (source) ); diff --git a/src/consumer.rs b/src/consumer.rs index 3edee85..a35e5f0 100644 --- a/src/consumer.rs +++ b/src/consumer.rs @@ -96,7 +96,7 @@ impl ConsumerTask { let sleeper = sleep(interval); tokio::pin!(sleeper); - let mut time_usec = 0u64; + let mut time_usec = 0i64; loop { tokio::select! { @@ -104,7 +104,7 @@ impl ConsumerTask { break; }, () = &mut sleeper => { - consumer_control_insert(&self.pool, &self.config.jetstream_hostname, &time_usec.to_string()).await?; + consumer_control_insert(&self.pool, &self.config.jetstream_hostname, time_usec).await?; sleeper.as_mut().reset(Instant::now() + interval); }, item = client.next() => { @@ -164,7 +164,12 @@ impl ConsumerTask { if feed_matcher.matches(&event_value) { tracing::debug!(feed_id = ?feed_matcher.feed, "matched event"); if let Some((uri, cid)) = model::to_post_strong_ref(&event) { - let feed_content = storage::model::FeedContent::new(feed_matcher.feed.clone(), uri, event.clone().time_us, cid); + let feed_content = storage::model::FeedContent{ + feed_id: feed_matcher.feed.clone(), + uri, + indexed_at: event.clone().time_us, + cid, + }; feed_content_insert(&self.pool, &feed_content).await?; } } @@ -200,7 +205,7 @@ pub(crate) mod model { max_message_size_bytes: u64, #[serde(skip_serializing_if = "Option::is_none")] - cursor: Option, + cursor: Option, }, } @@ -273,7 +278,7 @@ pub(crate) mod model { pub(crate) struct Event { pub(crate) did: String, pub(crate) kind: String, - pub(crate) time_us: u64, + pub(crate) time_us: i64, pub(crate) commit: Option, } diff --git a/src/http/handle_get_feed_skeleton.rs b/src/http/handle_get_feed_skeleton.rs index 686ea10..f57da13 100644 --- a/src/http/handle_get_feed_skeleton.rs +++ b/src/http/handle_get_feed_skeleton.rs @@ -58,7 +58,7 @@ pub async fn handle_get_feed_skeleton( let feed_control = feed_control.unwrap(); - if feed_control.allowed.len() > 0 { + if !feed_control.allowed.is_empty() { let authorization = headers.get("Authorization").and_then(|value| { value .to_str() @@ -104,7 +104,7 @@ pub async fn handle_get_feed_skeleton( let cursor = feed_items .iter() .last() - .map(|last_feed_item| format!("{},{}", last_feed_item.time_us(), last_feed_item.cid)); + .map(|last_feed_item| format!("{},{}", last_feed_item.indexed_at, last_feed_item.cid)); let feed_item_views = feed_items .iter() @@ -197,7 +197,7 @@ async fn did_from_jwt( Ok(claims.iss) } -fn parse_cursor(value: Option) -> Option<(u64, u32, u32, String)> { +fn parse_cursor(value: Option) -> Option<(i64, String)> { let value = value.as_ref()?; let parts = value.split(",").collect::>(); @@ -205,24 +205,11 @@ fn parse_cursor(value: Option) -> Option<(u64, u32, u32, String)> { return None; } - let time_us = parts[0].parse::(); + let time_us = parts[0].parse::(); if time_us.is_err() { return None; } let time_us = time_us.unwrap(); - let time_us_bytes = time_us.to_be_bytes(); - let indexed_at = u32::from_be_bytes([ - time_us_bytes[0], - time_us_bytes[1], - time_us_bytes[2], - time_us_bytes[3], - ]); - let indexed_at_more = u32::from_be_bytes([ - time_us_bytes[4], - time_us_bytes[5], - time_us_bytes[6], - time_us_bytes[7], - ]); - Some((time_us, indexed_at, indexed_at_more, parts[1].to_string())) + Some((time_us, parts[1].to_string())) } diff --git a/src/storage.rs b/src/storage.rs index 6b4a8e6..b590f12 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -14,51 +14,9 @@ pub mod model { pub struct FeedContent { pub feed_id: String, pub uri: String, - pub indexed_at: u32, - pub indexed_at_more: u32, + pub indexed_at: i64, pub cid: String, } - - impl FeedContent { - pub fn new(feed_id: String, uri: String, time_us: u64, cid: String) -> Self { - // Are their better ways to do this? Probably. - let time_us_bytes = time_us.to_be_bytes(); - let indexed_at = u32::from_be_bytes([ - time_us_bytes[0], - time_us_bytes[1], - time_us_bytes[2], - time_us_bytes[3], - ]); - let indexed_at_more = u32::from_be_bytes([ - time_us_bytes[4], - time_us_bytes[5], - time_us_bytes[6], - time_us_bytes[7], - ]); - - Self { - feed_id, - uri, - indexed_at, - indexed_at_more, - cid, - } - } - pub fn time_us(&self) -> u64 { - let indexed_at_bytes = self.indexed_at.to_be_bytes(); - let indexed_at_more_bytes = self.indexed_at_more.to_be_bytes(); - u64::from_be_bytes([ - indexed_at_bytes[0], - indexed_at_bytes[1], - indexed_at_bytes[2], - indexed_at_bytes[3], - indexed_at_more_bytes[0], - indexed_at_more_bytes[1], - indexed_at_more_bytes[2], - indexed_at_more_bytes[3], - ]) - } - } } pub async fn feed_content_insert( @@ -68,11 +26,10 @@ pub async fn feed_content_insert( let mut tx = pool.begin().await.context("failed to begin transaction")?; let now = Utc::now(); - sqlx::query("INSERT OR REPLACE INTO feed_content (feed_id, uri, indexed_at, indexed_at_more, cid, updated_at) VALUES (?, ?, ?, ?, ?, ?)") + sqlx::query("INSERT OR REPLACE INTO feed_content (feed_id, uri, indexed_at, cid, updated_at) VALUES (?, ?, ?, ?, ?)") .bind(&feed_content.feed_id) .bind(&feed_content.uri) .bind(feed_content.indexed_at) - .bind(feed_content.indexed_at_more) .bind(&feed_content.cid) .bind(now) .execute(tx.as_mut()) @@ -85,25 +42,24 @@ pub async fn feed_content_paginate( pool: &StoragePool, feed_uri: &str, limit: Option, - cursor: Option<(u64, u32, u32, String)>, + cursor: Option<(i64, String)>, ) -> Result> { let mut tx = pool.begin().await.context("failed to begin transaction")?; let limit = limit.unwrap_or(20).clamp(1, 100); - let results = if let Some((_time_us, indexed_at, indexed_at_more, cid)) = cursor { - let query = "SELECT * FROM feed_content WHERE feed_id = ? AND (indexed_at, indexed_at_more, cid) < (?, ?, ?) ORDER BY indexed_at DESC, indexed_at_more DESC, cid DESC LIMIT ?"; + let results = if let Some((indexed_at, cid)) = cursor { + let query = "SELECT * FROM feed_content WHERE feed_id = ? AND (indexed_at, cid) < (?, ?) ORDER BY indexed_at DESC, cid DESC LIMIT ?"; sqlx::query_as::<_, FeedContent>(query) .bind(feed_uri) .bind(indexed_at) - .bind(indexed_at_more) .bind(cid) .bind(limit) .fetch_all(tx.as_mut()) .await? } else { - let query = "SELECT * FROM feed_content WHERE feed_id = ? ORDER BY indexed_at DESC, indexed_at_more DESC, cid DESC LIMIT ?"; + let query = "SELECT * FROM feed_content WHERE feed_id = ? ORDER BY indexed_at DESC, cid DESC LIMIT ?"; sqlx::query_as::<_, FeedContent>(query) .bind(feed_uri) @@ -117,11 +73,7 @@ pub async fn feed_content_paginate( Ok(results) } -pub async fn consumer_control_insert( - pool: &StoragePool, - source: &str, - time_us: &str, -) -> Result<()> { +pub async fn consumer_control_insert(pool: &StoragePool, source: &str, time_us: i64) -> Result<()> { let mut tx = pool.begin().await.context("failed to begin transaction")?; let now = Utc::now(); @@ -137,11 +89,11 @@ pub async fn consumer_control_insert( tx.commit().await.context("failed to commit transaction") } -pub async fn consumer_control_get(pool: &StoragePool, source: &str) -> Result> { +pub async fn consumer_control_get(pool: &StoragePool, source: &str) -> Result> { let mut tx = pool.begin().await.context("failed to begin transaction")?; let result = - sqlx::query_scalar::<_, String>("SELECT time_us FROM consumer_control WHERE source = ?") + sqlx::query_scalar::<_, i64>("SELECT time_us FROM consumer_control WHERE source = ?") .bind(source) .fetch_optional(tx.as_mut()) .await @@ -149,7 +101,7 @@ pub async fn consumer_control_get(pool: &StoragePool, source: &str) -> Result().ok())) + Ok(result) } pub async fn verifcation_method_insert( @@ -226,12 +178,13 @@ mod tests { #[sqlx::test] async fn record_feed_content(pool: SqlitePool) -> sqlx::Result<()> { - let record = super::model::FeedContent::new( - "feed".to_string(), - "at://did:plc:qadlgs4xioohnhi2jg54mqds/app.bsky.feed.post/3la3bqjg4hx2n".to_string(), - 1730673934229172_u64, - "bafyreih74qdc6zskq7yarqi3xm634vnubf4g3ac5ieegbvakprxpjnsj74".to_string(), - ); + let record = super::model::FeedContent { + feed_id: "feed".to_string(), + uri: "at://did:plc:qadlgs4xioohnhi2jg54mqds/app.bsky.feed.post/3la3bqjg4hx2n" + .to_string(), + indexed_at: 1730673934229172_i64, + cid: "bafyreih74qdc6zskq7yarqi3xm634vnubf4g3ac5ieegbvakprxpjnsj74".to_string(), + }; super::feed_content_insert(&pool, &record) .await .expect("failed to insert record"); @@ -246,14 +199,14 @@ mod tests { records[0].uri, "at://did:plc:qadlgs4xioohnhi2jg54mqds/app.bsky.feed.post/3la3bqjg4hx2n" ); - assert_eq!(records[0].time_us(), 1730673934229172_u64); + assert_eq!(records[0].indexed_at, 1730673934229172_i64); Ok(()) } #[sqlx::test] async fn consumer_control(pool: SqlitePool) -> sqlx::Result<()> { - super::consumer_control_insert(&pool, "foo", "1730673934229172") + super::consumer_control_insert(&pool, "foo", 1730673934229172_i64) .await .expect("failed to insert record"); @@ -261,10 +214,10 @@ mod tests { super::consumer_control_get(&pool, "foo") .await .expect("failed to get record"), - Some(1730673934229172_u64) + Some(1730673934229172_i64) ); - super::consumer_control_insert(&pool, "foo", "1730673934229173") + super::consumer_control_insert(&pool, "foo", 1730673934229173_i64) .await .expect("failed to insert record"); @@ -272,7 +225,7 @@ mod tests { super::consumer_control_get(&pool, "foo") .await .expect("failed to get record"), - Some(1730673934229173_u64) + Some(1730673934229173_i64) ); Ok(())