From e8d9e58cd752c974db5789cc3a6652978dc9f971 Mon Sep 17 00:00:00 2001 From: Mia Date: Tue, 28 Jan 2025 19:09:52 +0000 Subject: [PATCH] start on commit processing --- Cargo.lock | 59 +++++++++++++++++++++++++++--- consumer/Cargo.toml | 1 + consumer/src/indexer/mod.rs | 68 ++++++++++++++++++++++++++++++++--- consumer/src/indexer/types.rs | 24 +++++++++++++ 4 files changed, 143 insertions(+), 9 deletions(-) create mode 100644 consumer/src/indexer/types.rs diff --git a/Cargo.lock b/Cargo.lock index d18e0a1f..2073e9a9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -82,6 +82,12 @@ dependencies = [ "windows-sys 0.59.0", ] +[[package]] +name = "anyhow" +version = "1.0.95" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "34ac096ce696dc2fcabef30516bb13c0a68a11d30131d3df6f04711467681b04" + [[package]] name = "async-trait" version = "0.1.85" @@ -251,7 +257,7 @@ dependencies = [ "multihash", "serde", "serde_bytes", - "unsigned-varint", + "unsigned-varint 0.8.0", ] [[package]] @@ -312,6 +318,7 @@ dependencies = [ "figment", "futures", "ipld-core", + "iroh-car", "parakeet-db", "serde", "serde_bytes", @@ -838,6 +845,22 @@ dependencies = [ "serde_bytes", ] +[[package]] +name = "iroh-car" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cb7f8cd4cb9aa083fba8b52e921764252d0b4dcb1cd6d120b809dbfe1106e81a" +dependencies = [ + "anyhow", + "cid", + "futures", + "serde", + "serde_ipld_dagcbor", + "thiserror 1.0.69", + "tokio", + "unsigned-varint 0.7.2", +] + [[package]] name = "is_terminal_polyfill" version = "1.70.1" @@ -949,7 +972,7 @@ checksum = "6b430e7953c29dd6a09afc29ff0bb69c6e306329ee6794700aee27b76a1aea8d" dependencies = [ "core2", "serde", - "unsigned-varint", + "unsigned-varint 0.8.0", ] [[package]] @@ -1095,7 +1118,7 @@ dependencies = [ "eyre", "serde", "serde_json", - "thiserror", + "thiserror 2.0.11", "walkdir", ] @@ -1569,13 +1592,33 @@ dependencies = [ "windows-sys 0.59.0", ] +[[package]] +name = "thiserror" +version = "1.0.69" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b6aaf5339b578ea85b50e080feb250a3e8ae8cfcdff9a461c9ec2904bc923f52" +dependencies = [ + "thiserror-impl 1.0.69", +] + [[package]] name = "thiserror" version = "2.0.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d452f284b73e6d76dd36758a0c8684b1d5be31f92b89d07fd5822175732206fc" dependencies = [ - "thiserror-impl", + "thiserror-impl 2.0.11", +] + +[[package]] +name = "thiserror-impl" +version = "1.0.69" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4fee6c4efc90059e10f81e6d42c60a18f76588c3d74cb83a0b242a2b6c7504c1" +dependencies = [ + "proc-macro2", + "quote", + "syn", ] [[package]] @@ -1812,7 +1855,7 @@ dependencies = [ "native-tls", "rand", "sha1", - "thiserror", + "thiserror 2.0.11", "utf-8", ] @@ -1858,6 +1901,12 @@ version = "0.1.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "e70f2a8b45122e719eb623c01822704c4e0907e7e426a05927e1a1cfff5b75d0" +[[package]] +name = "unsigned-varint" +version = "0.7.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6889a77d49f1f013504cec6bf97a2c730394adedaeb1deb5ea08949a50541105" + [[package]] name = "unsigned-varint" version = "0.8.0" diff --git a/consumer/Cargo.toml b/consumer/Cargo.toml index 2a83b2a3..b35c4423 100644 --- a/consumer/Cargo.toml +++ b/consumer/Cargo.toml @@ -12,6 +12,7 @@ eyre = "0.6.12" figment = { version = "0.10.19", features = ["env", "toml"] } futures = "0.3.31" ipld-core = "0.4.1" +iroh-car = "0.5.1" parakeet-db = { path = "../parakeet-db" } serde = { version = "1.0.217", features = ["derive"] } serde_bytes = "0.11" diff --git a/consumer/src/indexer/mod.rs b/consumer/src/indexer/mod.rs index 764ae824..e14c2a61 100644 --- a/consumer/src/indexer/mod.rs +++ b/consumer/src/indexer/mod.rs @@ -1,16 +1,34 @@ -use diesel_async::AsyncPgConnection; -use diesel_async::pooled_connection::deadpool::Pool; use crate::firehose::{AtpAccountEvent, AtpCommitEvent, AtpIdentityEvent, FirehoseEvent}; +use crate::indexer::types::{CollectionType, RecordTypes}; +use diesel_async::pooled_connection::deadpool::Pool; +use diesel_async::AsyncPgConnection; +use futures::StreamExt; +use ipld_core::cid::Cid; +use std::collections::HashMap; use tokio::sync::mpsc::Receiver; +use tracing::Instrument; -pub async fn relay_indexer(pool: Pool, mut rx: Receiver) -> eyre::Result<()> { +mod types; + +pub async fn relay_indexer( + pool: Pool, + mut rx: Receiver, +) -> eyre::Result<()> { let mut conn = pool.get().await?; while let Some(event) = rx.recv().await { let res = match event { FirehoseEvent::Identity(identity) => index_identity(&mut conn, identity).await, FirehoseEvent::Account(account) => index_account(&mut conn, account).await, - FirehoseEvent::Commit(commit) => index_commit(&mut conn, commit).await, + FirehoseEvent::Commit(commit) => { + let span = tracing::info_span!( + "commit", + seq = commit.seq, + repo = commit.repo, + rev = commit.rev + ); + index_commit(&mut conn, commit).instrument(span).await + } FirehoseEvent::Label(_) => { // We handle all labels through direct connections to labelers tracing::warn!("got #labels from the relay"); @@ -35,5 +53,47 @@ async fn index_account(conn: &mut AsyncPgConnection, account: AtpAccountEvent) - } async fn index_commit(conn: &mut AsyncPgConnection, commit: AtpCommitEvent) -> eyre::Result<()> { + // turn the car slice into a map of cid:block + let car_reader = iroh_car::CarReader::new(commit.blocks.as_slice()).await?; + let blocks = car_reader + .stream() + .filter_map(|car| async move { car.ok() }) + .collect::>>() + .await; + + for op in &commit.ops { + let Some((collection_raw, rkey)) = op.path.split_once("/") else { + tracing::warn!("op contained invalid path {}", op.path); + continue; + }; + + let collection = CollectionType::from_str(collection_raw); + if collection == CollectionType::Unsupported { + tracing::debug!("{} {collection_raw} is unsupported", op.action); + continue; + } + + if op.action == "create" || op.action == "update" { + if let Some(block) = op.cid.and_then(|cid| blocks.get(&cid)) { + let decoded: RecordTypes = match serde_ipld_dagcbor::from_slice(block) { + Ok(decoded) => decoded, + Err(err) => { + tracing::error!("Failed to decode record: {err}"); + continue; + } + }; + + dbg!(decoded); + } else { + tracing::error!("Missing Cid or the block was not found"); + continue; + } + } else if op.action == "delete" { + // + } else { + tracing::warn!("op contained invalid action {}", op.action); + } + } + Ok(()) } diff --git a/consumer/src/indexer/types.rs b/consumer/src/indexer/types.rs new file mode 100644 index 00000000..4e150c86 --- /dev/null +++ b/consumer/src/indexer/types.rs @@ -0,0 +1,24 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Deserialize, Serialize)] +#[serde(tag = "$type")] +pub enum RecordTypes {} + +#[derive(Debug, PartialOrd, PartialEq)] +pub enum CollectionType { + Unsupported, +} + +impl CollectionType { + pub(crate) fn from_str(input: &str) -> CollectionType { + match input { + _ => CollectionType::Unsupported, + } + } + + pub fn can_update(&self) -> bool { + match self { + CollectionType::Unsupported => false, + } + } +} -- 2.51.2