diff --git a/Cargo.lock b/Cargo.lock index 59b4048..f8f9261 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -201,6 +201,27 @@ dependencies = [ "generic-array", ] +[[package]] +name = "bobbin" +version = "0.0.1" +dependencies = [ + "anyhow", + "axum", + "bobbin-edge-index", + "bobbin-ingest", + "bobbin-record-lru", + "bobbin-slingshot-client", + "bobbin-types", + "bobbin-xrpc", + "bytes", + "jacquard-common", + "serde_json", + "tokio", + "tracing", + "tracing-subscriber", + "url", +] + [[package]] name = "bobbin-edge-index" version = "0.0.1" @@ -283,6 +304,7 @@ name = "bobbin-xrpc" version = "0.0.1" dependencies = [ "axum", + "bobbin-edge-index", "bobbin-record-lru", "bobbin-slingshot-client", "bobbin-types", diff --git a/Cargo.toml b/Cargo.toml index 53b8c4c..9a8e91b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,6 +1,7 @@ [workspace] resolver = "2" members = [ + "crates/bobbin", "crates/types", "crates/edge-index", "crates/ingest", diff --git a/containerfiles/hydrant.Containerfile b/containerfiles/hydrant.Containerfile new file mode 100644 index 0000000..8eb64ac --- /dev/null +++ b/containerfiles/hydrant.Containerfile @@ -0,0 +1,14 @@ +FROM docker.io/library/rust:1-alpine3.23 AS builder +RUN apk add --no-cache build-base musl-dev cmake perl pkgconfig +WORKDIR /src +COPY src/ ./ +RUN rm -f .cargo/config.toml +RUN cargo build --release --bin hydrant +RUN strip target/release/hydrant + +FROM docker.io/library/alpine:3.23 +RUN apk add --no-cache ca-certificates +COPY --from=builder /src/target/release/hydrant /usr/local/bin/hydrant +ENV HYDRANT_DATABASE_PATH=/var/lib/hydrant +EXPOSE 3000 +ENTRYPOINT ["/usr/local/bin/hydrant"] diff --git a/containerfiles/slingshot.Containerfile b/containerfiles/slingshot.Containerfile new file mode 100644 index 0000000..130ee66 --- /dev/null +++ b/containerfiles/slingshot.Containerfile @@ -0,0 +1,17 @@ +FROM docker.io/library/rust:1-alpine3.23 AS builder +RUN apk add --no-cache build-base musl-dev cmake perl pkgconfig +WORKDIR /src +COPY src/ ./ +RUN rm -f .cargo/config.toml +RUN cargo build --release --bin slingshot --package slingshot +RUN strip target/release/slingshot + +FROM docker.io/library/alpine:3.23 +RUN apk add --no-cache ca-certificates +WORKDIR /app +COPY --from=builder /src/target/release/slingshot /usr/local/bin/slingshot +COPY --from=builder /src/slingshot/static /app/static +ENV SLINGSHOT_CACHE_DIR=/var/lib/slingshot +ENV SLINGSHOT_BIND=0.0.0.0:8080 +EXPOSE 8080 +ENTRYPOINT ["/usr/local/bin/slingshot"] diff --git a/crates/bobbin/Cargo.toml b/crates/bobbin/Cargo.toml new file mode 100644 index 0000000..8cc6fd4 --- /dev/null +++ b/crates/bobbin/Cargo.toml @@ -0,0 +1,31 @@ +[package] +name = "bobbin" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[[bin]] +name = "bobbin" +path = "src/main.rs" + +[dependencies] +bobbin-edge-index = { workspace = true } +bobbin-ingest = { workspace = true } +bobbin-record-lru = { workspace = true } +bobbin-slingshot-client = { workspace = true } +bobbin-types = { workspace = true } +bobbin-xrpc = { workspace = true } + +jacquard-common = { workspace = true } + +axum = { workspace = true } +serde_json = { workspace = true } +tokio = { workspace = true, features = ["macros", "rt-multi-thread", "signal"] } +url = { workspace = true } +tracing = { workspace = true } +tracing-subscriber = { version = "0.3", features = ["env-filter", "fmt"] } +anyhow = { workspace = true } + +[dev-dependencies] +bytes = { workspace = true } diff --git a/crates/bobbin/src/main.rs b/crates/bobbin/src/main.rs new file mode 100644 index 0000000..b36882d --- /dev/null +++ b/crates/bobbin/src/main.rs @@ -0,0 +1,188 @@ +use std::env; +use std::net::SocketAddr; +use std::sync::Arc; + +use anyhow::{Context, anyhow}; +use bobbin_edge_index::{CoverageWatch, EdgeStore, HydrantCursor}; +use bobbin_ingest::{IngestConfig, RepoDidResolver, ResolveError, run as run_ingest}; +use bobbin_record_lru::{CacheCapacity, LruRecordStore, RecordStore}; +use bobbin_slingshot_client::{SlingshotClient, SlingshotError}; +use bobbin_types::record::RecordBody; +use bobbin_xrpc::{AppState, router}; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use jacquard_common::types::ident::AtIdentifier; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::string::AtUri; +use tracing_subscriber::EnvFilter; +use url::Url; + +const REPO_NSID: &str = "sh.tangled.repo"; + +#[tokio::main] +async fn main() -> anyhow::Result<()> { + tracing_subscriber::fmt() + .with_env_filter(EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"))) + .init(); + + let hydrant_url = env::var("BOBBIN_HYDRANT_URL") + .unwrap_or_else(|_| "http://127.0.0.1:13010".into()); + let slingshot_url = env::var("BOBBIN_SLINGSHOT_URL") + .unwrap_or_else(|_| "http://127.0.0.1:13011".into()); + let bind: SocketAddr = env::var("BOBBIN_BIND") + .unwrap_or_else(|_| "127.0.0.1:8090".into()) + .parse()?; + let cache_bytes: u64 = env::var("BOBBIN_RECORD_LRU_BYTES") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(64 * 1024 * 1024); + let start_cursor: u64 = env::var("BOBBIN_START_CURSOR") + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(0); + + let records: Arc = + Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(cache_bytes))); + let slingshot = SlingshotClient::new(Url::parse(&slingshot_url)?)?; + let edges = Arc::new(EdgeStore::new()); + let coverage = Arc::new(CoverageWatch::new()); + + let ingest_cfg = IngestConfig { + hydrant_base: Url::parse(&hydrant_url)?, + start_cursor: HydrantCursor::new(start_cursor), + }; + let resolver = Arc::new(SlingshotRepoDidResolver { + slingshot: slingshot.clone(), + records: records.clone(), + }); + let ingest_edges = edges.clone(); + let ingest_coverage = coverage.clone(); + let ingest_resolver = resolver.clone(); + let ingest_handle = tokio::spawn(async move { + run_ingest(ingest_cfg, ingest_edges, ingest_coverage, ingest_resolver).await + }); + + let state = AppState::new(records, slingshot, edges, coverage); + let app = router(state); + + tracing::info!(%bind, %hydrant_url, %slingshot_url, "bobbin listening"); + let listener = tokio::net::TcpListener::bind(&bind).await?; + let server_handle = tokio::spawn(async move { axum::serve(listener, app).await }); + + tokio::select! { + res = server_handle => match res { + Ok(Ok(())) => Ok(()), + Ok(Err(e)) => Err(e).context("axum server failed"), + Err(join) => Err(anyhow!("server task panicked: {join}")), + }, + res = ingest_handle => match res { + Ok(Ok(())) => Err(anyhow!("ingest run loop exited; loop is supposed to be infinite")), + Ok(Err(e)) => Err(e).context("ingest exited"), + Err(join) => Err(anyhow!("ingest task panicked: {join}")), + }, + } +} + +struct SlingshotRepoDidResolver { + slingshot: SlingshotClient, + records: Arc, +} + +impl RepoDidResolver for SlingshotRepoDidResolver { + async fn resolve( + &self, + repo_aturi: &AtUri, + ) -> Result>, ResolveError> { + let body = match self.records.get(repo_aturi) { + Some(b) => b, + None => match fetch_repo_record(&self.slingshot, repo_aturi).await? { + Some(b) => { + self.records.put(repo_aturi.clone(), b.clone()); + b + } + None => return Ok(None), + }, + }; + Ok(parse_repo_did(&body)) + } +} + +async fn fetch_repo_record( + slingshot: &SlingshotClient, + repo_aturi: &AtUri, +) -> Result>, ResolveError> { + let owner = match repo_aturi.authority() { + AtIdentifier::Did(d) => d, + AtIdentifier::Handle(_) => return Ok(None), + }; + let Some(rkey) = repo_aturi.rkey() else { + return Ok(None); + }; + let collection = Nsid::<&str>::new(REPO_NSID).expect("REPO_NSID literal validates as NSID"); + match slingshot.get_record(&owner, &collection, &rkey).await { + Ok(body) => Ok(Some(body)), + Err(SlingshotError::NotFound) => Ok(None), + Err(e) => Err(ResolveError::Transport(e.to_string())), + } +} + +fn parse_repo_did(body: &RecordBody) -> Option> { + let v: serde_json::Value = serde_json::from_slice(&body.value).ok()?; + let s = v.get("repoDid")?.as_str()?; + Did::::new(jacquard_common::deps::smol_str::SmolStr::new(s)).ok() +} + +#[cfg(test)] +mod tests { + use super::*; + use bytes::Bytes; + use jacquard_common::types::string::Cid; + + fn body_with_value(value: &str) -> RecordBody { + let cid: Cid = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i" + .parse() + .unwrap(); + RecordBody { + uri: AtUri::new_owned("at://did:plc:abalone/sh.tangled.repo/r1").unwrap(), + cid, + value: Bytes::copy_from_slice(value.as_bytes()), + } + } + + #[test] + fn parses_repo_did_when_field_present() { + let body = body_with_value( + r#"{"$type":"sh.tangled.repo","repoDid":"did:plc:limpet","name":"abalone","knot":"oyster.cafe","createdAt":"2026-05-01T00:00:00Z"}"#, + ); + assert_eq!( + parse_repo_did(&body).map(|d| d.as_ref().to_owned()), + Some("did:plc:limpet".into()), + ); + } + + #[test] + fn returns_none_when_repo_did_missing() { + let body = body_with_value( + r#"{"$type":"sh.tangled.repo","name":"abalone","knot":"oyster.cafe","createdAt":"2026-05-01T00:00:00Z"}"#, + ); + assert!(parse_repo_did(&body).is_none()); + } + + #[test] + fn returns_none_when_repo_did_is_not_a_string() { + let body = body_with_value(r#"{"$type":"sh.tangled.repo","repoDid":42}"#); + assert!(parse_repo_did(&body).is_none()); + } + + #[test] + fn returns_none_when_value_is_not_json() { + let body = body_with_value("not json at all"); + assert!(parse_repo_did(&body).is_none()); + } + + #[test] + fn returns_none_when_repo_did_fails_did_validation() { + let body = body_with_value(r#"{"repoDid":"definitely not a did"}"#); + assert!(parse_repo_did(&body).is_none()); + } +}