diff --git a/Cargo.lock b/Cargo.lock index 32c0d07..2380546 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -102,6 +102,29 @@ version = "1.1.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0" +[[package]] +name = "atproto-client" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0c1438d67d1fe2d06b5e5f8030c58cce991e074bab03b5e226a949507172f24" +dependencies = [ + "anyhow", + "async-trait", + "atproto-identity", + "atproto-oauth", + "atproto-record", + "bytes", + "reqwest", + "reqwest-chain", + "reqwest-middleware", + "serde", + "serde_json", + "thiserror 2.0.19", + "tokio", + "tracing", + "urlencoding", +] + [[package]] name = "atproto-dasl" version = "0.14.5" @@ -206,6 +229,28 @@ dependencies = [ "thiserror 2.0.19", ] +[[package]] +name = "atproto-tap" +version = "0.14.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "36a17328d4c4478e34f24b97a8112e3ff3dc1aa20fb9169938979f44e678d2fa" +dependencies = [ + "atproto-identity", + "base64", + "compact_str", + "futures", + "http", + "itoa", + "reqwest", + "serde", + "serde_json", + "thiserror 2.0.19", + "tokio", + "tokio-stream", + "tokio-websockets", + "tracing", +] + [[package]] name = "autocfg" version = "1.5.1" @@ -363,6 +408,15 @@ version = "1.12.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fc652a48c352aef3ea3aed32080501cf3ef6ed5da78602a020c991775b0aff04" +[[package]] +name = "castaway" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dec551ab6e7578819132c713a93c022a05d60159dc86e7a7050223577484c55a" +dependencies = [ + "rustversion", +] + [[package]] name = "cc" version = "1.4.0" @@ -476,6 +530,21 @@ version = "0.5.4" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a" +[[package]] +name = "compact_str" +version = "0.9.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9dfdd1c2274d9aa354115b09dc9a901d6c5576818cdf70d14cae2bdb47df00ab" +dependencies = [ + "castaway", + "cfg-if", + "itoa", + "rustversion", + "ryu", + "serde", + "static_assertions", +] + [[package]] name = "const-oid" version = "0.9.6" @@ -520,6 +589,16 @@ dependencies = [ "libc", ] +[[package]] +name = "core-foundation" +version = "0.10.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b2a6cd9ae233e7f62ba4e9353e81a88df7fc8a5987b8d445b4d90c879bd156f6" +dependencies = [ + "core-foundation-sys", + "libc", +] + [[package]] name = "core-foundation-sys" version = "0.8.7" @@ -1944,6 +2023,12 @@ dependencies = [ "portable-atomic", ] +[[package]] +name = "openssl-probe" +version = "0.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" + [[package]] name = "p256" version = "0.13.2" @@ -2503,6 +2588,18 @@ dependencies = [ "zeroize", ] +[[package]] +name = "rustls-native-certs" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dab5152771c58876a2146916e53e35057e1a4dfa2b9df0f0305b07f611fdea4d" +dependencies = [ + "openssl-probe", + "rustls-pki-types", + "schannel", + "security-framework", +] + [[package]] name = "rustls-pki-types" version = "1.15.1" @@ -2545,6 +2642,15 @@ dependencies = [ "winapi-util", ] +[[package]] +name = "schannel" +version = "0.1.29" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "91c1b7e4904c873ef0710c1f407dde2e6287de2bebc1bbbf7d430bb7cbffd939" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "scopeguard" version = "1.2.0" @@ -2566,6 +2672,29 @@ dependencies = [ "zeroize", ] +[[package]] +name = "security-framework" +version = "3.7.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b7f4bc775c73d9a02cde8bf7b2ec4c9d12743edf609006c7facc23998404cd1d" +dependencies = [ + "bitflags", + "core-foundation 0.10.1", + "core-foundation-sys", + "libc", + "security-framework-sys", +] + +[[package]] +name = "security-framework-sys" +version = "2.17.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ce2691df843ecc5d231c0b14ece2acc3efb62c0a398c7e1d875f3983ce020e3" +dependencies = [ + "core-foundation-sys", + "libc", +] + [[package]] name = "semver" version = "1.0.28" @@ -2707,6 +2836,12 @@ dependencies = [ "rand_core 0.6.4", ] +[[package]] +name = "simdutf8" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e3a9fe34e3e7a50316060351f37187a3f546bce95496156754b601a5fa71b76e" + [[package]] name = "slab" version = "0.4.12" @@ -2955,9 +3090,11 @@ version = "0.1.0" dependencies = [ "aes-gcm", "anyhow", + "atproto-client", "atproto-identity", "atproto-oauth", "atproto-record", + "atproto-tap", "axum", "base64", "chrono", @@ -2980,6 +3117,12 @@ dependencies = [ "urlencoding", ] +[[package]] +name = "static_assertions" +version = "1.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a2eb9349b6444b326872e140eb1cf5e7c522154d69e7a0ffb0fb81c06b37543f" + [[package]] name = "stringprep" version = "0.1.5" @@ -3046,7 +3189,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a13f3d0daba03132c0aa9767f98351b3488edc2c100cda2d2ec2b04f3d8d3c8b" dependencies = [ "bitflags", - "core-foundation", + "core-foundation 0.9.4", "system-configuration-sys", ] @@ -3236,6 +3379,28 @@ dependencies = [ "tokio", ] +[[package]] +name = "tokio-websockets" +version = "0.13.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d52efb639344a7c6adb8e62c6f3d2c19c001ff1b79a5041ba1c6ed42e19c6aa5" +dependencies = [ + "base64", + "bytes", + "fastrand", + "futures-core", + "futures-sink", + "http", + "httparse", + "ring", + "rustls-native-certs", + "rustls-pki-types", + "simdutf8", + "tokio", + "tokio-rustls", + "tokio-util", +] + [[package]] name = "tower" version = "0.5.3" diff --git a/Cargo.toml b/Cargo.toml index c7fbd6c..7ff1547 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -44,3 +44,5 @@ tower-livereload = "0.10" rust-embed = { version = "8", features = ["mime-guess"] } notify-debouncer-mini = "0.7" atproto-record = "0.14.5" +atproto-tap = "0.14.5" +atproto-client = "0.14.5" diff --git a/assets/components.css b/assets/components.css index c7f3cb0..d1c47ab 100644 --- a/assets/components.css +++ b/assets/components.css @@ -87,6 +87,7 @@ & .hint { color: var(--text-muted); font-size: var(--text-sm); + text-wrap: pretty; } } @@ -106,11 +107,15 @@ .centered { display: grid; - grid-template-columns: 1fr min(40rem, 100%) 1fr; + grid-template-columns: 1fr min(var(--centered-width, 40rem), 100%) 1fr; & > * { grid-column: 2; } + + &.thin { + --centered-width: 28rem; + } } /* Where a form's buttons sit. */ @@ -164,3 +169,12 @@ transform: none; } } + +.inline-code { + display: inline-block; + font-family: var(--font-mono); + font-size: var(--text-sm); + padding: 0 var(--space-1); + border-radius: var(--space-1); + background: oklch(99% 0 0 / 0.15); +} diff --git a/flake.nix b/flake.nix index 2d7aeae..075b651 100644 --- a/flake.nix +++ b/flake.nix @@ -12,6 +12,12 @@ outputs = inputs@{ self, flake-parts, ... }: + let + # Must match ModListing::NSID. Shared between perSystem and the NixOS + # module below, which are otherwise separate scopes. + lexiconNsids = [ "dev.starhaven.mod.listing" ]; + tapCollections = builtins.concatStringsSep "," lexiconNsids; + in flake-parts.lib.mkFlake { inherit inputs; } { systems = import inputs.systems; perSystem = @@ -57,6 +63,21 @@ ] ); + tap = pkgs.buildGoModule { + pname = "tap"; + version = "0.1.0"; + src = pkgs.fetchFromGitHub { + owner = "bluesky-social"; + repo = "indigo"; + rev = "cbaa83aee9dd4aa015fd0c245e1fb3cfbbe32817"; + sha256 = "sha256-QQvkfNjsfU3vReyd8xB2Dtdqninyv5Zem9SuRVTdnK4="; + }; + subPackages = [ "cmd/tap" ]; + vendorHash = "sha256-s1S+b+QbptqJ2mxqkvsn7M5VWfLrlwpWgRjg6lq2WVE="; + doCheck = false; + meta.mainProgram = "tap"; + }; + devPds = pkgs.buildNpmPackage { pname = "starhaven-dev-pds"; version = "0.1.0"; @@ -77,6 +98,9 @@ formatter = pkgs.nixfmt; packages = { + # Exposed (unlike dev-pds) for the NixOS module's tap service. + inherit tap; + default = let cargoToml = builtins.fromTOML (builtins.readFile ./Cargo.toml); @@ -108,6 +132,7 @@ pkgs.sqlx-cli pkgs.just pkgs.watchexec + tap ] ++ lib.optionals useMold [ pkgs.mold ]; shellHook = '' @@ -115,12 +140,28 @@ export RUST_SRC_PATH="${devToolchain}/lib/rustlib/src/rust/library"; # Local PDS + PLC directory runs in the background + dev_pds_pidfile=/tmp/starhaven-dev-pds.pid if ! kill -0 "$(cat "$dev_pds_pidfile" 2>/dev/null)" 2>/dev/null; then ${devPds}/bin/starhaven-dev-pds &> /tmp/starhaven-dev-pds.log & echo $! > "$dev_pds_pidfile" disown fi + # tap, pointed at the local dev PDS/PLC rather than the real network. + tap_pidfile=/tmp/starhaven-tap.pid + if ! kill -0 "$(cat "$tap_pidfile" 2>/dev/null)" 2>/dev/null; then + TAP_BIND=:2480 \ + TAP_DATABASE_URL=sqlite:///tmp/starhaven-tap.db \ + TAP_PLC_URL=http://localhost:2582 \ + TAP_RELAY_URL=http://localhost:2583 \ + TAP_SIGNAL_COLLECTION=${tapCollections} \ + TAP_COLLECTION_FILTERS=${tapCollections} \ + TAP_NO_REPLAY=true \ + ${tap}/bin/tap run &> /tmp/starhaven-tap.log & + echo $! > "$tap_pidfile" + disown + fi + # parallel frontend stops paying off after 8 threads rustc_threads=$(( $(nproc) < 8 ? $(nproc) : 8 )) @@ -166,27 +207,80 @@ default = "127.0.0.1:3000"; description = "Address:port to listen on."; }; + + tap = { + enable = lib.mkOption { + type = lib.types.bool; + default = cfg.enable; + description = "Whether to run the tap indexer. Defaults on with starhaven, since starhaven depends on it."; + }; + + package = lib.mkOption { + type = lib.types.package; + default = self.packages.${pkgs.stdenv.hostPlatform.system}.tap; + description = "The tap package to run."; + }; + + bind = lib.mkOption { + type = lib.types.str; + default = "127.0.0.1:2480"; + description = "Address:port for tap's HTTP/WebSocket server. starhaven connects to this same address."; + }; + + adminPasswordFile = lib.mkOption { + type = lib.types.nullOr lib.types.path; + default = null; + description = "EnvironmentFile containing `TAP_ADMIN_PASSWORD=...`."; + }; + }; }; config = lib.mkIf cfg.enable { systemd.services.starhaven = { description = "starhaven.dev server"; wantedBy = [ "multi-user.target" ]; + after = [ "network-online.target" ] ++ lib.optionals cfg.tap.enable [ "starhaven-tap.service" ]; + wants = [ "network-online.target" ] ++ lib.optionals cfg.tap.enable [ "starhaven-tap.service" ]; + + serviceConfig = { + ExecStart = lib.escapeShellArgs ( + [ + (lib.getExe cfg.package) + "--external-base" + cfg.externalBase + "--listen" + cfg.listen + "--data-dir" + "/var/lib/${stateDirectory}" + ] + ++ lib.optionals cfg.tap.enable [ + "--tap-hostname" + cfg.tap.bind + ] + ); + DynamicUser = true; + StateDirectory = stateDirectory; + Restart = "on-failure"; + }; + }; + + systemd.services.starhaven-tap = lib.mkIf cfg.tap.enable { + description = "tap indexer for the starhaven.dev server"; + wantedBy = [ "multi-user.target" ]; after = [ "network-online.target" ]; wants = [ "network-online.target" ]; serviceConfig = { - ExecStart = lib.escapeShellArgs [ - (lib.getExe cfg.package) - "--external-base" - cfg.externalBase - "--listen" - cfg.listen - "--data-dir" - "/var/lib/${stateDirectory}" + ExecStart = lib.getExe cfg.tap.package; + Environment = [ + "TAP_BIND=${cfg.tap.bind}" + "TAP_DATABASE_URL=sqlite:///var/lib/starhaven-tap/tap.db" + "TAP_SIGNAL_COLLECTION=${tapCollections}" + "TAP_COLLECTION_FILTERS=${tapCollections}" ]; + EnvironmentFile = lib.mkIf (cfg.tap.adminPasswordFile != null) cfg.tap.adminPasswordFile; DynamicUser = true; - StateDirectory = stateDirectory; + StateDirectory = "starhaven-tap"; Restart = "on-failure"; }; }; diff --git a/src/atproto/actor.rs b/src/atproto/actor.rs index 20c5a67..66ac145 100644 --- a/src/atproto/actor.rs +++ b/src/atproto/actor.rs @@ -16,9 +16,7 @@ use sqlx::{Row, SqlitePool}; use crate::atproto::id::{Did, Handle}; use crate::state::AppState; -/// `.test` handles live only on the local dev PDS (see `dev-pds/`), which -/// isn't on real DNS or the public PLC directory, so they need their own -/// resolution path rather than the DNS/HTTPS one real handles use. +/// `.test` handles only exist on the local dev PDS, not real DNS/PLC. fn is_dev_handle(handle: &str) -> bool { handle.ends_with(".test") } @@ -52,9 +50,7 @@ async fn dev_resolve_document(http_client: &reqwest::Client, handle: &str) -> Op } async fn dev_resolve_document_by_did(http_client: &reqwest::Client, did: &str) -> Option { - // Not atproto_identity::plc::query: the pinned crate version always - // prepends "https://" to the hostname, which breaks against the dev - // PLC's plain http. + // Not atproto_identity::plc::query: the pinned version always forces https. http_client .get(format!("{DEV_PLC_URL}/{did}")) .send() @@ -211,9 +207,7 @@ pub async fn verified_handle( resolver: &SharedIdentityResolver, did: &Did, ) -> Result> { - // Public PLC directory first, then the dev PLC: a real DID never lives - // on the dev PLC, and a dev DID never lives on the public one, so trying - // both costs nothing but one extra local round-trip on a dev miss. + // Public PLC first, then dev PLC as a fallback on a dev DID. let document = match resolver.resolve(did.as_str()).await { Ok(document) => document, Err(_) => match dev_resolve_document_by_did(&resolver.http_client, did.as_str()).await { diff --git a/src/atproto/id.rs b/src/atproto/id.rs index 178dd9d..c0256ae 100644 --- a/src/atproto/id.rs +++ b/src/atproto/id.rs @@ -95,6 +95,12 @@ impl Rkey { Self(atproto_record::tid::Tid::new().encode()) } + // Not validated: for an rkey out of tap's event stream, which already + // verified the record's MST inclusion and signature before delivering it. + pub fn from_verified(value: impl Into) -> Self { + Self(value.into()) + } + pub fn as_str(&self) -> &str { &self.0 } diff --git a/src/atproto/mod.rs b/src/atproto/mod.rs index 4e9c10b..6d69a76 100644 --- a/src/atproto/mod.rs +++ b/src/atproto/mod.rs @@ -8,3 +8,4 @@ pub mod actor; pub mod id; pub mod lexicon; +pub mod tap; diff --git a/src/atproto/tap.rs b/src/atproto/tap.rs new file mode 100644 index 0000000..ae6c320 --- /dev/null +++ b/src/atproto/tap.rs @@ -0,0 +1,63 @@ +//! Background indexer. Reads records from the tap service (i.e. the firehose) and indexes them into the local database. + +use atproto_tap::{connect, RecordAction, RecordEvent, TapConfig, TapEvent}; +use tokio_stream::StreamExt; + +use crate::atproto::actor; +use crate::atproto::id::{Did, Handle, Rkey}; +use crate::atproto::lexicon::ModListing; +use crate::state::AppState; + +/// Spawn the indexer. Runs until the process exits; reconnects on its own. +pub fn spawn(state: AppState) { + tokio::spawn(async move { + let config = TapConfig::builder() + .hostname(state.config.tap_hostname.clone()) + .build(); + + let mut stream = connect(config); + while let Some(result) = stream.next().await { + match result { + Ok(event) => handle_event(&state, &event).await, + Err(err) => eprintln!("tap: stream error: {err}"), + } + } + }); +} + +async fn handle_event(state: &AppState, event: &TapEvent) { + match event { + TapEvent::Record { record, .. } if &*record.collection == ModListing::NSID => { + if let Err(err) = handle_record(state, record).await { + eprintln!("tap: failed to index {}: {err}", record.at_uri()); + } + } + TapEvent::Record { .. } => {} + TapEvent::Identity { identity, .. } => { + let Ok(handle) = Handle::new(&identity.handle) else { + return; + }; + let did = Did::new(identity.did.to_string()); + if let Err(err) = actor::upsert(&state.db, &did, &handle).await { + eprintln!("tap: failed to index identity for {}: {err}", identity.did); + } + } + } +} + +async fn handle_record(state: &AppState, record: &RecordEvent) -> anyhow::Result<()> { + let did = Did::new(record.did.to_string()); + let rkey = Rkey::from_verified(record.rkey.to_string()); + + match record.action { + RecordAction::Delete => ModListing::delete(&state.db, &did, &rkey).await, + RecordAction::Create | RecordAction::Update => { + let value = record + .record_value() + .ok_or_else(|| anyhow::anyhow!("{} action with no record", record.action))?; + let listing = ModListing::from_value(value)?; + listing.validate()?; + listing.save(&state.db, &did, &rkey).await + } + } +} diff --git a/src/config.rs b/src/config.rs index 15067f6..5e62564 100644 --- a/src/config.rs +++ b/src/config.rs @@ -27,6 +27,11 @@ pub struct Config { /// State directory. #[arg(long, default_value = ".")] pub data_dir: PathBuf, + + /// `tap` service hostname, for indexing repo/identity events into the + /// local cache independent of who wrote them. + #[arg(long, default_value = "localhost:2480")] + pub tap_hostname: String, } impl Config { diff --git a/src/main.rs b/src/main.rs index f30826c..257feb2 100644 --- a/src/main.rs +++ b/src/main.rs @@ -26,6 +26,8 @@ async fn main() { .expect("failed to initialize application state"); let listen = state.config.listen; + atproto::tap::spawn(state.clone()); + let app = Router::new() .merge(oauth::router()) .merge(assets::router()) diff --git a/src/oauth/login.rs b/src/oauth/login.rs index 14748f4..be64a37 100644 --- a/src/oauth/login.rs +++ b/src/oauth/login.rs @@ -55,11 +55,29 @@ pub struct LoginForm { /// `GET /login` -- a minimal handle-entry form. pub async fn login_form() -> Html { layout(html! { - h1 { "log in" } - form method="post" action="/login" { - label for="handle" { "handle, DID, or PDS URL" } - input type="text" id="handle" name="handle" placeholder="alice.bsky.social" required; - button type="submit" { "continue" } + div.centered.thin { + form method="post" action="/login" { + div.form-item { + label for="handle" { "Handle" } + input type="text" id="handle" name="handle" placeholder="alex.starhaven.dev" + autocapitalize="none" autocorrect="off" autocomplete="username" inputmode="url" + tabindex="1" required; + p.hint { + "Use your "; + a href="https://atproto.com" { "AT Protocol" } + " handle to log in. If you're unsure, this is likely your Star Haven (" + code.inline-code { ".starhaven.dev" } + ") or "; + a href="https://bsky.app" { "Bluesky" } + " ("; + code.inline-code { ".bsky.app" } + ") account."; + } + } + div.form-actions { + button.button type="submit" tabindex="2" { "Login" } + } + } } }) }