···11+[ 156ms] [ERROR] Failed to load resource: the server responded with a status of 403 () @ https://www.printables.com/:0
22+[ 328ms] [ERROR] Failed to load resource: the server responded with a status of 403 () @ https://www.printables.com/favicon.ico:0
···11+- generic [ref=e3]:
22+ - generic [ref=e4]:
33+ - heading "Sorry, you have been blocked" [level=1] [ref=e5]
44+ - heading "You are unable to access printables.com" [level=2] [ref=e6]
55+ - generic [ref=e12]:
66+ - generic [ref=e13]:
77+ - heading "Why have I been blocked?" [level=2] [ref=e14]
88+ - paragraph [ref=e15]: This website is using a security service to protect itself from online attacks. The action you just performed triggered the security solution. There are several actions that could trigger this block including submitting a certain word or phrase, a SQL command or malformed data.
99+ - generic [ref=e16]:
1010+ - heading "What can I do to resolve this?" [level=2] [ref=e17]
1111+ - paragraph [ref=e18]: You can email the site owner to let them know you were blocked. Please include what you were doing when this page came up and the Cloudflare Ray ID found at the bottom of this page.
1212+ - paragraph [ref=e20]:
1313+ - generic [ref=e21]:
1414+ - text: "Cloudflare Ray ID:"
1515+ - strong [ref=e22]: a0f533fa896139cc
1616+ - generic [ref=e23]: •
1717+ - generic [ref=e24]:
1818+ - text: "Your IP:"
1919+ - button "Click to reveal" [ref=e25] [cursor=pointer]
2020+ - generic [ref=e26]: •
2121+ - generic [ref=e27]:
2222+ - text: Performance & security by
2323+ - link "Cloudflare" [ref=e28] [cursor=pointer]:
2424+ - /url: https://www.cloudflare.com/5xx-error-landing
+20
.playwright-mcp/page-2026-06-21T18-51-41-927Z.yml
···11+- generic [active] [ref=e1]:
22+ - main [ref=e2]:
33+ - generic [ref=e3]:
44+ - generic [ref=e4]:
55+ - img "Icon for makerworld.com" [ref=e5]
66+ - heading "makerworld.com" [level=1] [ref=e6]
77+ - heading "Performing security verification" [level=2] [ref=e7]
88+ - paragraph [ref=e8]: This website uses a security service to protect against malicious bots. This page is displayed while the website verifies you are not a bot.
99+ - contentinfo [ref=e12]:
1010+ - generic [ref=e14]:
1111+ - generic [ref=e16]:
1212+ - text: "Ray ID:"
1313+ - code [ref=e17]: a0f534abacbea281
1414+ - generic [ref=e18]:
1515+ - generic [ref=e19]:
1616+ - text: Performance and Security by
1717+ - link "Cloudflare" [ref=e20] [cursor=pointer]:
1818+ - /url: https://www.cloudflare.com?utm_source=challenge&utm_campaign=m
1919+ - link "Privacy" [ref=e22] [cursor=pointer]:
2020+ - /url: https://www.cloudflare.com/privacypolicy/
+18
.playwright-mcp/page-2026-06-21T18-51-42-273Z.yml
···11+- generic [active] [ref=e1]:
22+ - main [ref=e2]:
33+ - generic [ref=e3]:
44+ - heading "www.thingiverse.com" [level=1] [ref=e5]
55+ - heading "Performing security verification" [level=2] [ref=e6]
66+ - paragraph [ref=e7]: This website uses a security service to protect against malicious bots. This page is displayed while the website verifies you are not a bot.
77+ - contentinfo [ref=e14]:
88+ - generic [ref=e16]:
99+ - generic [ref=e18]:
1010+ - text: "Ray ID:"
1111+ - code [ref=e19]: a0f534b1bf60a2c9
1212+ - generic [ref=e20]:
1313+ - generic [ref=e21]:
1414+ - text: Performance and Security by
1515+ - link "Cloudflare" [ref=e22] [cursor=pointer]:
1616+ - /url: https://www.cloudflare.com?utm_source=challenge&utm_campaign=m
1717+ - link "Privacy" [ref=e24] [cursor=pointer]:
1818+ - /url: https://www.cloudflare.com/privacypolicy/
+49
.playwright-mcp/page-2026-06-21T18-53-57-432Z.yml
···11+- generic [active] [ref=e1]:
22+ - generic [ref=e5] [cursor=pointer]:
33+ - generic [ref=e6]:
44+ - img [ref=e7]
55+ - heading [level=3] [ref=e11]: Your app is being rebuilt.
66+ - paragraph [ref=e12]: A non-hot-reloadable change occurred and we must rebuild.
77+ - main [ref=e14]:
88+ - generic [ref=e15]:
99+ - heading "Polymodel" [level=1] [ref=e16]
1010+ - generic [ref=e17]:
1111+ - strong [ref=e18]: Thingiverse on atproto.
1212+ - text: Project/profile/like lexicons, social feeds, and a browser STL viewer.
1313+ - strong [ref=e19]: "Workflow:"
1414+ - text: Jira + Confluence + jj + Polytoken.
1515+ - link "Open viewer" [ref=e20] [cursor=pointer]:
1616+ - /url: /viewer/body-f-chest-v4
1717+ - generic [ref=e21]:
1818+ - article [ref=e22]:
1919+ - generic [ref=e24]: NSID
2020+ - generic [ref=e25]:
2121+ - heading "Project lexicon" [level=2] [ref=e26]
2222+ - paragraph [ref=e27]: "Define the model/project record: title, description, files, preview assets, tags, license, and attribution."
2323+ - article [ref=e28]:
2424+ - generic [ref=e30]: DID
2525+ - generic [ref=e31]:
2626+ - heading "Profile publishing" [level=2] [ref=e32]
2727+ - paragraph [ref=e33]: Use ATProto identity and profile-oriented feeds so makers can publish and browse from their own repos.
2828+ - article [ref=e34]:
2929+ - generic [ref=e36]: ★
3030+ - generic [ref=e37]:
3131+ - heading "Likes and saves" [level=2] [ref=e38]
3232+ - paragraph [ref=e39]: Start with lightweight feedback records that can drive hot/recent discovery without inventing a heavyweight backend.
3333+ - article [ref=e40]:
3434+ - generic [ref=e42]: STL
3535+ - generic [ref=e43]:
3636+ - heading "Model viewer" [level=2] [ref=e44]
3737+ - paragraph [ref=e45]: Spike three-d and kiss3d as WASM viewer candidates; choose based on STL loading, orbit controls, rendering quality, and Dioxus integration friction.
3838+ - link "Open viewer" [ref=e46] [cursor=pointer]:
3939+ - /url: /viewer/body-f-chest-v4
4040+ - article [ref=e47]:
4141+ - generic [ref=e49]: PM
4242+ - generic [ref=e50]:
4343+ - heading "Jira and Confluence" [level=2] [ref=e51]
4444+ - paragraph [ref=e52]: Use Jira and Confluence as both planning/demo narrative and agent-readable workflow state.
4545+ - article [ref=e53]:
4646+ - generic [ref=e55]: pt
4747+ - generic [ref=e56]:
4848+ - heading "Polytoken workflow" [level=2] [ref=e57]
4949+ - paragraph [ref=e58]: Use facets, skills, hooks, project vars, jj workspaces, and review gates to make the demo workflow real.
+1-1
.polytoken/skills/jira-solo-tasker-plan/SKILL.md
···104104Use jj terminology consistently:
105105106106- workspace: isolated checkout created with `jj workspace add` or the project’s configured wrapper;
107107-- change: a fresh jj change created for the ticket work with `jj new -m "<ticket-key>: <summary>"` before any edits;
107107+- change: rename the change with the right context for the ticket work with `jj desc -m "<ticket-key>: <summary>"` before any edits;
108108- bookmark: optional named ref for pushing/sharing, normally the same slug as the workspace.
109109110110Do not mention Git worktrees, `.worktrees/`, `git merge --squash`, or `.merge-lock` unless the current project explicitly uses those. This workflow defaults to jj workspaces.
···2323web = ["dioxus/web"]
2424desktop = ["dioxus/desktop"]
2525mobile = ["dioxus/mobile"]
2626-server = ["dioxus/server", "dep:axum", "dep:tower", "dep:tower-http"]
2626+server = [
2727+ "dioxus/server",
2828+ "dep:axum",
2929+ "dep:tower",
3030+ "dep:tower-http",
3131+ # Indexing infrastructure (Hydrant firehose + SQLite projection)
3232+ "dep:hydrant",
3333+ "dep:sqlx",
3434+ "dep:libsqlite3-sys",
3535+ "dep:tokio",
3636+ "dep:futures",
3737+ "dep:anyhow",
3838+ "dep:dotenvy",
3939+ "dep:chrono",
4040+ "dep:tracing-subscriber",
4141+]
27422843[dependencies]
2944dioxus = { version = "0.7", features = ["router", "fullstack"] }
···4055axum = { version = "0.8", optional = true }
4156tower = { version = "0.5", optional = true }
4257tower-http = { version = "0.6", optional = true, features = ["trace"] }
5858+5959+# Server-only indexing infrastructure.
6060+# hydrant is pinned to a specific commit to protect against breaking API changes.
6161+hydrant = { git = "https://tangled.org/ptr.pet/hydrant", rev = "1145fa1bb194372ceea81ab48b5d68d304da34ad", optional = true }
6262+# Compile-time-checked queries. sqlx version matches the `sqlx-cli` provided by
6363+# the nix dev shell so `cargo sqlx prepare` and the macros agree on the .sqlx
6464+# cache format. `cargo sqlx prepare --features server` generates the committed
6565+# .sqlx cache for offline/CI builds.
6666+sqlx = { version = "0.9", default-features = false, features = [
6767+ "runtime-tokio",
6868+ "sqlite",
6969+ "macros",
7070+ "migrate",
7171+], optional = true }
7272+# Bundle SQLite (libsqlite3-sys bundled enables FTS5) so no system libsqlite3 is
7373+# required. Version follows sqlx 0.9's libsqlite3-sys requirement.
7474+libsqlite3-sys = { version = "0.37", features = ["bundled"], optional = true }
7575+tokio = { version = "1", features = ["rt-multi-thread", "macros", "signal"], optional = true }
7676+futures = { version = "0.3", optional = true }
7777+anyhow = { version = "1", optional = true }
7878+dotenvy = { version = "0.15", optional = true }
7979+chrono = { version = "0.4", optional = true }
8080+tracing-subscriber = { version = "0.3", optional = true, default-features = false, features = [
8181+ "std",
8282+ "fmt",
8383+ "env-filter",
8484+] }
43854486[target.'cfg(all(target_family = "wasm", target_os = "unknown"))'.dependencies]
4587console_error_panic_hook = "0.1"
+39-1
build.rs
···11+//! Build script.
22+//!
33+//! Reads `.env` and bakes any `POLYMODEL_*` variables into `src/env.rs` as
44+//! compile-time constants. This is how configuration reaches the WASM bundle,
55+//! which has no filesystem or process environment to read at runtime.
66+//!
77+//! Server-side runtime configuration (DATABASE_URL, HYDRANT_*) is NOT handled
88+//! here: the server binary loads `.env` directly via `dotenvy` and reads it into
99+//! a typed `ServerConfig` struct (see `src/indexing/config.rs`).
1010+1111+use std::env;
1212+use std::fs::File;
1313+use std::io::Write;
1414+115fn main() {
22- dotenvy::dotenv().ok();
316 println!("cargo:rerun-if-changed=.env");
1717+1818+ dotenvy::dotenv().ok();
1919+2020+ let mut f = File::create("./src/env.rs").expect("create src/env.rs");
2121+ f.write_all(
2222+ b"// This file is automatically generated by build.rs.\n\
2323+ // Frontend (WASM) compile-time configuration.\n\
2424+ // Do not edit; regenerate by rebuilding.\n\n",
2525+ )
2626+ .expect("write env.rs header");
2727+2828+ // Bake every POLYMODEL_* env var into the bundle as a `&'static str` const.
2929+ for (key, value) in env::vars() {
3030+ if let Some(rest) = key.strip_prefix("POLYMODEL_") {
3131+ if rest.is_empty() {
3232+ continue;
3333+ }
3434+ let line = format!(
3535+ "#[allow(unused)]\npub const {}: &str = \"{}\";\n",
3636+ key,
3737+ value.replace('\\', "\\\\").replace('"', "\\\"")
3838+ );
3939+ f.write_all(line.as_bytes()).expect("write env.rs const");
4040+ }
4141+ }
442}
+2
flake.nix
···4848 cargo-nextest
4949 cargo-insta
5050 cargo-bloat
5151+ # Provides the `cargo sqlx` subcommand (package name in nixpkgs is sqlx-cli).
5252+ sqlx-cli
51535254 dioxus-cli
5355 wasm-bindgen-cli
+15
justfile
···99 cargo clippy -p polymodel --all-targets --features server --fix --allow-dirty --allow-staged --allow-no-vcs -- -D warnings
10101111# Compile checks across the whole workspace plus the app's server/wasm targets.
1212+# Note: the server check uses compile-time `query!` macros that validate against a
1313+# migrated SQLite database. Run `just migrate` first (or run `just sqlx-prepare`)
1414+# so the macros can resolve.
1215check:
1316 cargo check --workspace
1417 cargo check -p polymodel --features server
1518 cargo check -p polymodel --target wasm32-unknown-unknown --features web
1919+2020+# Create and migrate the SQLite projection database.
2121+# Required by the compile-time `query!` macros: run before the first server
2222+# check, and re-run after editing migrations.
2323+migrate:
2424+ mkdir -p ./data
2525+ -cargo sqlx database create --database-url 'sqlite:./data/polymodel.db'
2626+ cargo sqlx migrate run --source migrations --database-url 'sqlite:./data/polymodel.db'
2727+2828+# Regenerate the committed `.sqlx` offline query cache after changing SQL.
2929+sqlx-prepare:
3030+ cargo sqlx prepare -- -p polymodel --features server
16311732# Run unit tests across the workspace.
1833test:
+167
migrations/001_projection_schema.sql
···11+-- Polymodel SQLite projection schema.
22+-- Derived views projected from Hydrant ATProto events for space.polymodel.* collections.
33+--
44+-- NOTE: WAL mode is intentionally NOT set here. sqlx runs migrations inside a
55+-- transaction and SQLite silently ignores journal_mode changes within one.
66+-- WAL is enabled programmatically on the connection pool (see src/indexing/db.rs).
77+88+-- Projection state (durable cursor for crash recovery)
99+CREATE TABLE projection_state (
1010+ key TEXT PRIMARY KEY,
1111+ value TEXT NOT NULL,
1212+ updated_at INTEGER NOT NULL
1313+);
1414+1515+-- Derived content views
1616+CREATE TABLE things (
1717+ did TEXT NOT NULL,
1818+ rkey TEXT NOT NULL,
1919+ uri TEXT NOT NULL UNIQUE,
2020+ cid TEXT NOT NULL,
2121+ name TEXT NOT NULL,
2222+ summary TEXT,
2323+ license TEXT,
2424+ tags_json TEXT,
2525+ tags_text TEXT,
2626+ instructions_text TEXT,
2727+ cover_json TEXT,
2828+ derived_from_uri TEXT,
2929+ created_at INTEGER NOT NULL,
3030+ indexed_at INTEGER NOT NULL,
3131+ PRIMARY KEY (did, rkey)
3232+);
3333+CREATE INDEX idx_things_did ON things(did, created_at DESC);
3434+CREATE INDEX idx_things_created ON things(created_at DESC);
3535+3636+CREATE TABLE models (
3737+ did TEXT NOT NULL,
3838+ rkey TEXT NOT NULL,
3939+ uri TEXT NOT NULL UNIQUE,
4040+ cid TEXT NOT NULL,
4141+ name TEXT NOT NULL,
4242+ summary TEXT,
4343+ created_at INTEGER NOT NULL,
4444+ indexed_at INTEGER NOT NULL,
4545+ PRIMARY KEY (did, rkey)
4646+);
4747+4848+CREATE TABLE parts (
4949+ did TEXT NOT NULL,
5050+ rkey TEXT NOT NULL,
5151+ uri TEXT NOT NULL UNIQUE,
5252+ cid TEXT NOT NULL,
5353+ name TEXT NOT NULL,
5454+ format TEXT,
5555+ file_json TEXT,
5656+ created_at INTEGER NOT NULL,
5757+ indexed_at INTEGER NOT NULL,
5858+ PRIMARY KEY (did, rkey)
5959+);
6060+6161+-- Junction tables (relationship lookups + ordering)
6262+CREATE TABLE thing_models (
6363+ thing_uri TEXT NOT NULL,
6464+ model_uri TEXT NOT NULL,
6565+ position INTEGER NOT NULL,
6666+ PRIMARY KEY (thing_uri, model_uri)
6767+);
6868+CREATE INDEX idx_thing_models_model ON thing_models(model_uri);
6969+CREATE INDEX idx_thing_models_thing_pos ON thing_models(thing_uri, position);
7070+7171+CREATE TABLE model_parts (
7272+ model_uri TEXT NOT NULL,
7373+ part_uri TEXT NOT NULL,
7474+ position INTEGER NOT NULL,
7575+ PRIMARY KEY (model_uri, part_uri)
7676+);
7777+CREATE INDEX idx_model_parts_part ON model_parts(part_uri);
7878+CREATE INDEX idx_model_parts_model_pos ON model_parts(model_uri, position);
7979+8080+-- Social record projections
8181+CREATE TABLE likes (
8282+ did TEXT NOT NULL,
8383+ rkey TEXT NOT NULL,
8484+ subject_uri TEXT NOT NULL,
8585+ created_at INTEGER NOT NULL,
8686+ PRIMARY KEY (did, rkey),
8787+ UNIQUE (did, subject_uri)
8888+);
8989+CREATE INDEX idx_likes_subject ON likes(subject_uri);
9090+9191+CREATE TABLE saves (
9292+ did TEXT NOT NULL,
9393+ rkey TEXT NOT NULL,
9494+ subject_uri TEXT NOT NULL,
9595+ note TEXT,
9696+ created_at INTEGER NOT NULL,
9797+ PRIMARY KEY (did, rkey),
9898+ UNIQUE (did, subject_uri)
9999+);
100100+CREATE INDEX idx_saves_subject ON saves(subject_uri);
101101+102102+CREATE TABLE tags (
103103+ did TEXT NOT NULL,
104104+ rkey TEXT NOT NULL,
105105+ subject_uri TEXT NOT NULL,
106106+ tag TEXT NOT NULL,
107107+ created_at INTEGER NOT NULL,
108108+ PRIMARY KEY (did, rkey)
109109+);
110110+CREATE INDEX idx_tags_subject ON tags(subject_uri);
111111+CREATE INDEX idx_tags_tag ON tags(tag);
112112+113113+CREATE TABLE listitems (
114114+ did TEXT NOT NULL,
115115+ rkey TEXT NOT NULL,
116116+ list_uri TEXT NOT NULL,
117117+ subject_uri TEXT NOT NULL,
118118+ created_at INTEGER NOT NULL,
119119+ PRIMARY KEY (did, rkey)
120120+);
121121+CREATE INDEX idx_listitems_list ON listitems(list_uri, created_at);
122122+123123+-- Denormalized counts
124124+CREATE TABLE content_stats (
125125+ uri TEXT PRIMARY KEY,
126126+ like_count INTEGER NOT NULL DEFAULT 0,
127127+ save_count INTEGER NOT NULL DEFAULT 0,
128128+ tag_count INTEGER NOT NULL DEFAULT 0
129129+);
130130+131131+-- Full-text search (FTS5, external content)
132132+CREATE VIRTUAL TABLE things_fts USING fts5(
133133+ name, summary, tags_text, instructions_text,
134134+ content=things, content_rowid=rowid,
135135+ tokenize='porter unicode61'
136136+);
137137+138138+CREATE TRIGGER things_ai AFTER INSERT ON things BEGIN
139139+ INSERT INTO things_fts(rowid, name, summary, tags_text, instructions_text)
140140+ VALUES (new.rowid, new.name, new.summary, new.tags_text, new.instructions_text);
141141+END;
142142+CREATE TRIGGER things_ad AFTER DELETE ON things BEGIN
143143+ INSERT INTO things_fts(things_fts, rowid, name, summary, tags_text, instructions_text)
144144+ VALUES ('delete', old.rowid, old.name, old.summary, old.tags_text, old.instructions_text);
145145+END;
146146+CREATE TRIGGER things_au AFTER UPDATE ON things BEGIN
147147+ INSERT INTO things_fts(things_fts, rowid, name, summary, tags_text, instructions_text)
148148+ VALUES ('delete', old.rowid, old.name, old.summary, old.tags_text, old.instructions_text);
149149+ INSERT INTO things_fts(rowid, name, summary, tags_text, instructions_text)
150150+ VALUES (new.rowid, new.name, new.summary, new.tags_text, new.instructions_text);
151151+END;
152152+153153+-- OAuth session persistence (DB-backed ClientAuthStore; used in Task 3)
154154+CREATE TABLE oauth_sessions (
155155+ did TEXT NOT NULL,
156156+ session_id TEXT NOT NULL,
157157+ session_data BLOB NOT NULL,
158158+ created_at INTEGER NOT NULL,
159159+ updated_at INTEGER NOT NULL,
160160+ PRIMARY KEY (did, session_id)
161161+);
162162+163163+CREATE TABLE oauth_auth_requests (
164164+ state TEXT PRIMARY KEY,
165165+ request_data BLOB NOT NULL,
166166+ created_at INTEGER NOT NULL
167167+);
+31
src/indexing/config.rs
···11+//! Server runtime configuration.
22+//!
33+//! Configuration that varies per deployment (e.g. `DATABASE_URL`) and is only
44+//! needed by the server binary. These values are read from the process
55+//! environment at startup, not baked into the WASM bundle.
66+//!
77+//! The server binary calls `dotenvy::dotenv()` in `main` before
88+//! [`ServerConfig::load`]. Hydrant reads its own `HYDRANT_*` variables via its
99+//! own `Config::from_env`; this struct holds only app-level server config and is
1010+//! the single place server-side environment variables are read.
1111+1212+use std::env;
1313+1414+/// App-level server configuration loaded from the environment.
1515+#[derive(Debug, Clone)]
1616+pub struct ServerConfig {
1717+ /// sqlx connect URL, e.g. `sqlite:./data/polymodel.db`.
1818+ pub database_url: String,
1919+}
2020+2121+impl ServerConfig {
2222+ /// Read the server config from the process environment.
2323+ ///
2424+ /// Call `dotenvy::dotenv()` first so variables in `.env` are visible. This
2525+ /// is the only place server-side env reads happen.
2626+ pub fn load() -> Self {
2727+ let database_url =
2828+ env::var("DATABASE_URL").unwrap_or_else(|_| "sqlite:./data/polymodel.db".to_string());
2929+ Self { database_url }
3030+ }
3131+}
+74
src/indexing/db.rs
···11+//! SQLite pool initialization and projection-cursor persistence.
22+//!
33+//! Two durable stores back the server:
44+//! 1. Hydrant's fjall store (`HYDRANT_DATABASE_PATH`) for raw ATProto records.
55+//! 2. This SQLite projection database (`DATABASE_URL`) for app-specific derived
66+//! views. The pool here owns the second.
77+88+use std::str::FromStr;
99+1010+use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions};
1111+use sqlx::{Sqlite, SqlitePool, Transaction};
1212+1313+use super::config::ServerConfig;
1414+1515+/// Open the SQLite pool, enable WAL mode, run migrations, and return the pool.
1616+///
1717+/// WAL mode is set on the connection options (not in a migration) because SQLite
1818+/// silently ignores `journal_mode` changes inside the transaction sqlx uses to
1919+/// run migrations.
2020+pub async fn init_db(cfg: &ServerConfig) -> anyhow::Result<SqlitePool> {
2121+ ensure_data_dir(&cfg.database_url)?;
2222+2323+ let options = SqliteConnectOptions::from_str(&cfg.database_url)?
2424+ .create_if_missing(true)
2525+ .journal_mode(SqliteJournalMode::Wal);
2626+2727+ let pool = SqlitePoolOptions::new()
2828+ .max_connections(5)
2929+ .connect_with(options)
3030+ .await?;
3131+3232+ sqlx::migrate!("./migrations").run(&pool).await?;
3333+ tracing::info!("sqlite migrations applied");
3434+ Ok(pool)
3535+}
3636+3737+fn ensure_data_dir(db_url: &str) -> anyhow::Result<()> {
3838+ if let Some(path) = db_url.strip_prefix("sqlite:")
3939+ && let Some(parent) = std::path::Path::new(path).parent()
4040+ {
4141+ std::fs::create_dir_all(parent)?;
4242+ }
4343+ Ok(())
4444+}
4545+4646+struct CursorRow {
4747+ value: String,
4848+}
4949+5050+/// Load the last persisted projection cursor (the Hydrant event id), if any.
5151+pub async fn load_cursor(db: &SqlitePool) -> anyhow::Result<Option<u64>> {
5252+ let row: Option<CursorRow> = sqlx::query_as!(
5353+ CursorRow,
5454+ "SELECT value FROM projection_state WHERE key = ?",
5555+ "cursor"
5656+ )
5757+ .fetch_optional(db)
5858+ .await?;
5959+ Ok(row.and_then(|r| r.value.parse().ok()))
6060+}
6161+6262+/// Persist the projection cursor inside an in-flight projection transaction.
6363+pub async fn save_cursor(tx: &mut Transaction<'_, Sqlite>, cursor: u64) -> anyhow::Result<()> {
6464+ let now = chrono::Utc::now().timestamp_millis();
6565+ sqlx::query!(
6666+ "INSERT INTO projection_state (key, value, updated_at) VALUES ('cursor', ?, ?) \
6767+ ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = excluded.updated_at",
6868+ cursor.to_string(),
6969+ now
7070+ )
7171+ .execute(&mut **tx)
7272+ .await?;
7373+ Ok(())
7474+}
+135
src/indexing/mod.rs
···11+//! Server-side indexing infrastructure.
22+//!
33+//! [Hydrant](hydrant) consumes the ATProto firehose for `space.polymodel.*`
44+//! collections. A Tokio task projects each event into the SQLite derived tables
55+//! (see [`projection`]) and persists a durable cursor so the pipeline resumes
66+//! after a crash without skipping or re-processing beyond the last success.
77+//!
88+//! Two durable stores back the server:
99+//! 1. Hydrant's fjall store (`HYDRANT_DATABASE_PATH`) for raw ATProto records.
1010+//! 2. The SQLite projection database (`DATABASE_URL`) for app-specific views.
1111+1212+pub mod config;
1313+pub mod db;
1414+pub mod projection;
1515+pub mod setup;
1616+1717+use std::time::Duration;
1818+1919+use futures::StreamExt;
2020+use hydrant::control::Hydrant;
2121+use sqlx::SqlitePool;
2222+2323+/// Start the indexing pipeline: initialize Hydrant, resume from the persisted
2424+/// cursor, and spawn the firehose driver plus the projection consumer.
2525+///
2626+/// Both background tasks run for the lifetime of the server's Tokio runtime, so
2727+/// this must be called from within that runtime (the `dioxus::serve` closure in
2828+/// `main` runs there).
2929+pub async fn start_indexing(db: SqlitePool) -> anyhow::Result<()> {
3030+ let hydrant = setup::init_hydrant().await?;
3131+3232+ // Drive the firehose + backfill for the process lifetime. run() takes &self
3333+ // and returns a future that resolves when a fatal component exits; it is
3434+ // called inside the spawned task because the returned future borrows the
3535+ // handle (Rust 2024 lifetime capture). Hydrant's fatal paths (crawler,
3636+ // firehose worker, stats ticker) call process::abort() on unexpected exit, so
3737+ // a non-aborting resolution here is rare; in any case the supervised consumer
3838+ // keeps indexing alive by resubscribing from the persisted cursor.
3939+ let runner = hydrant.clone();
4040+ tokio::spawn(async move {
4141+ let fut = match runner.run() {
4242+ Ok(fut) => fut,
4343+ Err(e) => {
4444+ tracing::error!(error = %e, "hydrant::run() failed to start");
4545+ return;
4646+ }
4747+ };
4848+ if let Err(e) = fut.await {
4949+ tracing::error!(error = %e, "hydrant driver exited with error");
5050+ }
5151+ });
5252+5353+ // Supervised projection consumer. The stream can terminate (a slow consumer
5454+ // gets a StreamError and the channel closes); when it does, resubscribe from
5555+ // the persisted cursor after backoff instead of stopping. A process restart
5656+ // also resumes correctly, but supervision keeps indexing live without one.
5757+ tokio::spawn(consume_loop(hydrant.clone(), db.clone()));
5858+5959+ Ok(())
6060+}
6161+6262+/// Maximum backoff for resubscribe / projection retry (exponential, capped).
6363+const MAX_BACKOFF: Duration = Duration::from_secs(60);
6464+6565+/// Supervised consumer loop: subscribe, project each event until the stream
6666+/// ends, then resubscribe from the persisted cursor with exponential backoff.
6767+///
6868+/// The cursor never advances past an event whose projection failed: that event
6969+/// is retried in place before any later event is read. Projection is idempotent,
7070+/// so replaying across a resubscribe is safe.
7171+async fn consume_loop(hydrant: Hydrant, db: SqlitePool) {
7272+ let mut reconnect = Duration::from_secs(2);
7373+ loop {
7474+ // Resume from the last successfully projected event.
7575+ let cursor = match db::load_cursor(&db).await {
7676+ Ok(c) => c.unwrap_or(0),
7777+ Err(e) => {
7878+ tracing::error!(error = %e, "cursor load failed; will retry");
7979+ tokio::time::sleep(reconnect).await;
8080+ reconnect = (reconnect * 2).min(MAX_BACKOFF);
8181+ continue;
8282+ }
8383+ };
8484+8585+ let mut events = hydrant.subscribe(Some(cursor));
8686+ tracing::info!(cursor, "(re)subscribed to hydrant event stream");
8787+ reconnect = Duration::from_secs(2); // reset after a successful subscribe
8888+8989+ // Consume until the stream terminates.
9090+ while let Some(result) = events.next().await {
9191+ let event = match result {
9292+ Ok(event) => event,
9393+ Err(e) => {
9494+ // Transient stream error: keep draining. If the channel has
9595+ // closed, the next poll returns None and we resubscribe.
9696+ tracing::warn!(error = %e, "hydrant stream error");
9797+ continue;
9898+ }
9999+ };
100100+101101+ let pe = match projection::ProjectionEvent::from_hydrant(&event) {
102102+ Ok(Some(pe)) => pe,
103103+ Ok(None) => continue, // non-record event: not projected
104104+ Err(e) => {
105105+ tracing::warn!(error = %e, "dropping unconvertible event");
106106+ continue;
107107+ }
108108+ };
109109+110110+ // Retry the same event until it projects. A persistently-failing
111111+ // ("poison") event wedges here and eventually stalls the stream,
112112+ // triggering a resubscribe; it never advances the cursor, so no data
113113+ // is lost. Loud per-retry logging keeps it observable.
114114+ let mut retry = Duration::from_secs(1);
115115+ loop {
116116+ match projection::project_event(&db, &pe).await {
117117+ Ok(()) => break,
118118+ Err(e) => {
119119+ tracing::error!(error = %e, seq = pe.seq, "projection failed; retrying");
120120+ tokio::time::sleep(retry).await;
121121+ retry = (retry * 2).min(MAX_BACKOFF);
122122+ }
123123+ }
124124+ }
125125+ }
126126+127127+ // Stream ended (next() returned None): back off and resubscribe.
128128+ tracing::warn!(
129129+ reconnect_secs = reconnect.as_secs(),
130130+ "hydrant event stream ended; resubscribing"
131131+ );
132132+ tokio::time::sleep(reconnect).await;
133133+ reconnect = (reconnect * 2).min(MAX_BACKOFF);
134134+ }
135135+}
+1071
src/indexing/projection.rs
···11+//! Projection of Hydrant ATProto events into the SQLite derived tables.
22+//!
33+//! [`project_event`] is the single entry point. Each event is projected inside
44+//! one SQLite transaction, and the projection cursor is persisted in that same
55+//! transaction so crash recovery never advances past an unprojected event.
66+//!
77+//! All upserts are idempotent (`INSERT OR REPLACE` for content, `INSERT OR
88+//! IGNORE` for social records) so replaying events from an earlier cursor is
99+//! safe. Social counts are maintained with a two-step delete: look up the
1010+//! subject of the existing row, delete it, and only decrement if a row was
1111+//! actually removed.
1212+1313+use serde_json::Value;
1414+use sqlx::{Sqlite, SqlitePool, Transaction};
1515+1616+/// A domain-neutral projection event, decoupled from Hydrant's internal types so
1717+/// the projection can be unit-tested with synthetic data.
1818+#[derive(Debug, Clone)]
1919+pub struct ProjectionEvent {
2020+ /// Monotonic Hydrant event id; used as the durable projection cursor.
2121+ pub seq: u64,
2222+ pub did: String,
2323+ pub collection: String,
2424+ pub rkey: String,
2525+ /// `"create"`, `"update"`, or `"delete"`.
2626+ pub action: String,
2727+ /// The record body for create/update events; `None` for delete events.
2828+ pub record: Option<Value>,
2929+ pub cid: Option<String>,
3030+}
3131+3232+impl ProjectionEvent {
3333+ /// Convert a Hydrant [`hydrant::control::Event`] into a projection event.
3434+ ///
3535+ /// Returns `Ok(None)` for non-record events (identity/account) which are not
3636+ /// projected into SQLite.
3737+ pub fn from_hydrant(event: &hydrant::control::Event) -> anyhow::Result<Option<Self>> {
3838+ let Some(rec) = event.record.as_ref() else {
3939+ return Ok(None);
4040+ };
4141+ Ok(Some(ProjectionEvent {
4242+ seq: event.id,
4343+ did: rec.did.to_string(),
4444+ collection: rec.collection.to_string(),
4545+ rkey: rec.rkey.to_string(),
4646+ action: rec.action.to_string(),
4747+ record: rec.record.clone(),
4848+ cid: rec.cid.as_ref().map(|c| c.to_string()),
4949+ }))
5050+ }
5151+}
5252+5353+/// Project a single event into SQLite, advancing the cursor on success.
5454+pub async fn project_event(db: &SqlitePool, event: &ProjectionEvent) -> anyhow::Result<()> {
5555+ let mut tx = db.begin().await?;
5656+ let action = event.action.as_str();
5757+ let collection = event.collection.as_str();
5858+5959+ match (action, collection) {
6060+ ("create" | "update", "space.polymodel.library.thing") => {
6161+ project_thing_upsert(&mut tx, event).await?;
6262+ }
6363+ ("delete", "space.polymodel.library.thing") => {
6464+ project_thing_delete(&mut tx, event).await?;
6565+ }
6666+ ("create" | "update", "space.polymodel.library.model") => {
6767+ project_model_upsert(&mut tx, event).await?;
6868+ }
6969+ ("delete", "space.polymodel.library.model") => {
7070+ project_model_delete(&mut tx, event).await?;
7171+ }
7272+ ("create" | "update", "space.polymodel.library.part") => {
7373+ project_part_upsert(&mut tx, event).await?;
7474+ }
7575+ ("delete", "space.polymodel.library.part") => {
7676+ project_part_delete(&mut tx, event).await?;
7777+ }
7878+ ("create", "space.polymodel.graph.like") => {
7979+ project_like_create(&mut tx, event).await?;
8080+ }
8181+ ("delete", "space.polymodel.graph.like") => {
8282+ project_like_delete(&mut tx, event).await?;
8383+ }
8484+ ("create", "space.polymodel.graph.save") => {
8585+ project_save_create(&mut tx, event).await?;
8686+ }
8787+ ("delete", "space.polymodel.graph.save") => {
8888+ project_save_delete(&mut tx, event).await?;
8989+ }
9090+ ("create", "space.polymodel.graph.tag") => {
9191+ project_tag_create(&mut tx, event).await?;
9292+ }
9393+ ("delete", "space.polymodel.graph.tag") => {
9494+ project_tag_delete(&mut tx, event).await?;
9595+ }
9696+ ("create", "space.polymodel.graph.listitem") => {
9797+ project_listitem_create(&mut tx, event).await?;
9898+ }
9999+ ("delete", "space.polymodel.graph.listitem") => {
100100+ project_listitem_delete(&mut tx, event).await?;
101101+ }
102102+ // graph.list lives in Hydrant's fjall store; no SQLite projection.
103103+ // account/identity events are filtered out by from_hydrant.
104104+ _ => {
105105+ tracing::debug!(action, collection, "no projection rule for event");
106106+ }
107107+ }
108108+109109+ super::db::save_cursor(&mut tx, event.seq).await?;
110110+ tx.commit().await?;
111111+ Ok(())
112112+}
113113+114114+// ---------------------------------------------------------------------------
115115+// Field extraction helpers (records arrive as serde_json::Value).
116116+// ---------------------------------------------------------------------------
117117+118118+fn s_str(v: &Value, key: &str) -> Option<String> {
119119+ v.get(key).and_then(|x| x.as_str()).map(String::from)
120120+}
121121+122122+/// URI from a single strongRef field, e.g. `derivedFrom: { uri, cid }`.
123123+fn ref_uri(v: &Value, key: &str) -> Option<String> {
124124+ v.get(key)
125125+ .and_then(|r| r.get("uri"))
126126+ .and_then(|u| u.as_str())
127127+ .map(String::from)
128128+}
129129+130130+/// URIs from an array of strongRefs, e.g. `models: [{ uri, cid }, ...]`.
131131+fn ref_uris(v: &Value, key: &str) -> Vec<String> {
132132+ v.get(key)
133133+ .and_then(|x| x.as_array())
134134+ .map(|arr| {
135135+ arr.iter()
136136+ .filter_map(|el| el.get("uri").and_then(|u| u.as_str()).map(String::from))
137137+ .collect()
138138+ })
139139+ .unwrap_or_default()
140140+}
141141+142142+fn str_array(v: &Value, key: &str) -> Vec<String> {
143143+ v.get(key)
144144+ .and_then(|x| x.as_array())
145145+ .map(|arr| {
146146+ arr.iter()
147147+ .filter_map(|el| el.as_str().map(String::from))
148148+ .collect()
149149+ })
150150+ .unwrap_or_default()
151151+}
152152+153153+fn parse_millis(v: &Value, key: &str) -> Option<i64> {
154154+ chrono::DateTime::parse_from_rfc3339(v.get(key)?.as_str()?)
155155+ .ok()
156156+ .map(|dt| dt.timestamp_millis())
157157+}
158158+159159+fn now_millis() -> i64 {
160160+ chrono::Utc::now().timestamp_millis()
161161+}
162162+163163+fn at_uri(did: &str, collection: &str, rkey: &str) -> String {
164164+ format!("at://{did}/{collection}/{rkey}")
165165+}
166166+167167+// ---------------------------------------------------------------------------
168168+// Library content
169169+// ---------------------------------------------------------------------------
170170+171171+async fn project_thing_upsert(
172172+ tx: &mut Transaction<'_, Sqlite>,
173173+ e: &ProjectionEvent,
174174+) -> anyhow::Result<()> {
175175+ let Some(rec) = &e.record else {
176176+ return Ok(());
177177+ };
178178+ let uri = at_uri(&e.did, &e.collection, &e.rkey);
179179+ // Required content fields: skip the event gracefully if absent (malformed
180180+ // record) rather than storing a blank row.
181181+ let Some(cid) = e.cid.clone() else {
182182+ return Ok(());
183183+ };
184184+ let Some(name) = s_str(rec, "name") else {
185185+ return Ok(());
186186+ };
187187+ let summary = s_str(rec, "summary");
188188+ let license = s_str(rec, "license");
189189+ let instructions = s_str(rec, "instructions");
190190+ let cover_json = rec.get("cover").map(|v| v.to_string());
191191+ let derived_from_uri = ref_uri(rec, "derivedFrom");
192192+193193+ let tags = str_array(rec, "tags");
194194+ let tags_json = if tags.is_empty() {
195195+ None
196196+ } else {
197197+ serde_json::to_string(&tags).ok()
198198+ };
199199+ let tags_text = if tags.is_empty() {
200200+ None
201201+ } else {
202202+ Some(tags.join(" "))
203203+ };
204204+205205+ let created_at = parse_millis(rec, "createdAt").unwrap_or_else(now_millis);
206206+ let indexed_at = now_millis();
207207+208208+ // UPSERT (not INSERT OR REPLACE): REPLACE performs DELETE+INSERT with a new
209209+ // rowid, which skips the AFTER UPDATE FTS trigger and (with recursive_triggers
210210+ // off, the sqlx default) the AFTER DELETE trigger too — corrupting things_fts
211211+ // on every update/replay. ON CONFLICT DO UPDATE preserves the rowid and fires
212212+ // things_au so the external-content FTS index stays in sync.
213213+ sqlx::query!(
214214+ "INSERT INTO things \
215215+ (did, rkey, uri, cid, name, summary, license, tags_json, tags_text, \
216216+ instructions_text, cover_json, derived_from_uri, created_at, indexed_at) \
217217+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) \
218218+ ON CONFLICT(did, rkey) DO UPDATE SET \
219219+ uri = excluded.uri, cid = excluded.cid, name = excluded.name, \
220220+ summary = excluded.summary, license = excluded.license, \
221221+ tags_json = excluded.tags_json, tags_text = excluded.tags_text, \
222222+ instructions_text = excluded.instructions_text, cover_json = excluded.cover_json, \
223223+ derived_from_uri = excluded.derived_from_uri, created_at = excluded.created_at, \
224224+ indexed_at = excluded.indexed_at",
225225+ e.did,
226226+ e.rkey,
227227+ uri,
228228+ cid,
229229+ name,
230230+ summary,
231231+ license,
232232+ tags_json,
233233+ tags_text,
234234+ instructions,
235235+ cover_json,
236236+ derived_from_uri,
237237+ created_at,
238238+ indexed_at,
239239+ )
240240+ .execute(&mut **tx)
241241+ .await?;
242242+243243+ // Rebuild thing_models junction preserving array order.
244244+ sqlx::query!("DELETE FROM thing_models WHERE thing_uri = ?", uri)
245245+ .execute(&mut **tx)
246246+ .await?;
247247+ for (pos, model_uri) in ref_uris(rec, "models").into_iter().enumerate() {
248248+ sqlx::query!(
249249+ "INSERT OR REPLACE INTO thing_models (thing_uri, model_uri, position) \
250250+ VALUES (?, ?, ?)",
251251+ uri,
252252+ model_uri,
253253+ pos as i64,
254254+ )
255255+ .execute(&mut **tx)
256256+ .await?;
257257+ }
258258+259259+ // Ensure a stats row exists so feeds can inner-join.
260260+ sqlx::query!(
261261+ "INSERT OR IGNORE INTO content_stats (uri, like_count, save_count, tag_count) \
262262+ VALUES (?, 0, 0, 0)",
263263+ uri,
264264+ )
265265+ .execute(&mut **tx)
266266+ .await?;
267267+268268+ Ok(())
269269+}
270270+271271+async fn project_thing_delete(
272272+ tx: &mut Transaction<'_, Sqlite>,
273273+ e: &ProjectionEvent,
274274+) -> anyhow::Result<()> {
275275+ let uri = at_uri(&e.did, &e.collection, &e.rkey);
276276+ sqlx::query!(
277277+ "DELETE FROM things WHERE did = ? AND rkey = ?",
278278+ e.did,
279279+ e.rkey
280280+ )
281281+ .execute(&mut **tx)
282282+ .await?;
283283+ sqlx::query!("DELETE FROM thing_models WHERE thing_uri = ?", uri)
284284+ .execute(&mut **tx)
285285+ .await?;
286286+ sqlx::query!("DELETE FROM content_stats WHERE uri = ?", uri)
287287+ .execute(&mut **tx)
288288+ .await?;
289289+ Ok(())
290290+}
291291+292292+async fn project_model_upsert(
293293+ tx: &mut Transaction<'_, Sqlite>,
294294+ e: &ProjectionEvent,
295295+) -> anyhow::Result<()> {
296296+ let Some(rec) = &e.record else {
297297+ return Ok(());
298298+ };
299299+ let uri = at_uri(&e.did, &e.collection, &e.rkey);
300300+ let Some(cid) = e.cid.clone() else {
301301+ return Ok(());
302302+ };
303303+ let Some(name) = s_str(rec, "name") else {
304304+ return Ok(());
305305+ };
306306+ let summary = s_str(rec, "summary");
307307+ let created_at = parse_millis(rec, "createdAt").unwrap_or_else(now_millis);
308308+ let indexed_at = now_millis();
309309+310310+ sqlx::query!(
311311+ "INSERT INTO models \
312312+ (did, rkey, uri, cid, name, summary, created_at, indexed_at) \
313313+ VALUES (?, ?, ?, ?, ?, ?, ?, ?) \
314314+ ON CONFLICT(did, rkey) DO UPDATE SET \
315315+ uri = excluded.uri, cid = excluded.cid, name = excluded.name, \
316316+ summary = excluded.summary, created_at = excluded.created_at, \
317317+ indexed_at = excluded.indexed_at",
318318+ e.did,
319319+ e.rkey,
320320+ uri,
321321+ cid,
322322+ name,
323323+ summary,
324324+ created_at,
325325+ indexed_at,
326326+ )
327327+ .execute(&mut **tx)
328328+ .await?;
329329+330330+ sqlx::query!("DELETE FROM model_parts WHERE model_uri = ?", uri)
331331+ .execute(&mut **tx)
332332+ .await?;
333333+ for (pos, part_uri) in ref_uris(rec, "parts").into_iter().enumerate() {
334334+ sqlx::query!(
335335+ "INSERT OR REPLACE INTO model_parts (model_uri, part_uri, position) \
336336+ VALUES (?, ?, ?)",
337337+ uri,
338338+ part_uri,
339339+ pos as i64,
340340+ )
341341+ .execute(&mut **tx)
342342+ .await?;
343343+ }
344344+ Ok(())
345345+}
346346+347347+async fn project_model_delete(
348348+ tx: &mut Transaction<'_, Sqlite>,
349349+ e: &ProjectionEvent,
350350+) -> anyhow::Result<()> {
351351+ let uri = at_uri(&e.did, &e.collection, &e.rkey);
352352+ sqlx::query!(
353353+ "DELETE FROM models WHERE did = ? AND rkey = ?",
354354+ e.did,
355355+ e.rkey
356356+ )
357357+ .execute(&mut **tx)
358358+ .await?;
359359+ sqlx::query!("DELETE FROM model_parts WHERE model_uri = ?", uri)
360360+ .execute(&mut **tx)
361361+ .await?;
362362+ Ok(())
363363+}
364364+365365+async fn project_part_upsert(
366366+ tx: &mut Transaction<'_, Sqlite>,
367367+ e: &ProjectionEvent,
368368+) -> anyhow::Result<()> {
369369+ let Some(rec) = &e.record else {
370370+ return Ok(());
371371+ };
372372+ let uri = at_uri(&e.did, &e.collection, &e.rkey);
373373+ let Some(cid) = e.cid.clone() else {
374374+ return Ok(());
375375+ };
376376+ let Some(name) = s_str(rec, "name") else {
377377+ return Ok(());
378378+ };
379379+ let format = s_str(rec, "format");
380380+ let file_json = rec.get("file").map(|v| v.to_string());
381381+ let created_at = parse_millis(rec, "createdAt").unwrap_or_else(now_millis);
382382+ let indexed_at = now_millis();
383383+384384+ sqlx::query!(
385385+ "INSERT INTO parts \
386386+ (did, rkey, uri, cid, name, format, file_json, created_at, indexed_at) \
387387+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) \
388388+ ON CONFLICT(did, rkey) DO UPDATE SET \
389389+ uri = excluded.uri, cid = excluded.cid, name = excluded.name, \
390390+ format = excluded.format, file_json = excluded.file_json, \
391391+ created_at = excluded.created_at, indexed_at = excluded.indexed_at",
392392+ e.did,
393393+ e.rkey,
394394+ uri,
395395+ cid,
396396+ name,
397397+ format,
398398+ file_json,
399399+ created_at,
400400+ indexed_at,
401401+ )
402402+ .execute(&mut **tx)
403403+ .await?;
404404+ Ok(())
405405+}
406406+407407+async fn project_part_delete(
408408+ tx: &mut Transaction<'_, Sqlite>,
409409+ e: &ProjectionEvent,
410410+) -> anyhow::Result<()> {
411411+ sqlx::query!(
412412+ "DELETE FROM parts WHERE did = ? AND rkey = ?",
413413+ e.did,
414414+ e.rkey
415415+ )
416416+ .execute(&mut **tx)
417417+ .await?;
418418+ Ok(())
419419+}
420420+421421+// ---------------------------------------------------------------------------
422422+// Social records (two-step delete + conditional count maintenance)
423423+// ---------------------------------------------------------------------------
424424+425425+#[derive(sqlx::FromRow)]
426426+struct SubjectRow {
427427+ subject_uri: String,
428428+}
429429+430430+async fn project_like_create(
431431+ tx: &mut Transaction<'_, Sqlite>,
432432+ e: &ProjectionEvent,
433433+) -> anyhow::Result<()> {
434434+ let Some(rec) = &e.record else {
435435+ return Ok(());
436436+ };
437437+ let Some(subject) = ref_uri(rec, "subject") else {
438438+ return Ok(());
439439+ };
440440+ let created_at = parse_millis(rec, "createdAt").unwrap_or_else(now_millis);
441441+442442+ let res = sqlx::query!(
443443+ "INSERT OR IGNORE INTO likes (did, rkey, subject_uri, created_at) VALUES (?, ?, ?, ?)",
444444+ e.did,
445445+ e.rkey,
446446+ subject,
447447+ created_at,
448448+ )
449449+ .execute(&mut **tx)
450450+ .await?;
451451+452452+ if res.rows_affected() == 1 {
453453+ sqlx::query!(
454454+ "INSERT INTO content_stats (uri, like_count, save_count, tag_count) VALUES (?, 1, 0, 0) \
455455+ ON CONFLICT(uri) DO UPDATE SET like_count = like_count + 1",
456456+ subject,
457457+ )
458458+ .execute(&mut **tx)
459459+ .await?;
460460+ }
461461+ Ok(())
462462+}
463463+464464+async fn project_like_delete(
465465+ tx: &mut Transaction<'_, Sqlite>,
466466+ e: &ProjectionEvent,
467467+) -> anyhow::Result<()> {
468468+ let old: Option<SubjectRow> = sqlx::query_as!(
469469+ SubjectRow,
470470+ "SELECT subject_uri FROM likes WHERE did = ? AND rkey = ?",
471471+ e.did,
472472+ e.rkey
473473+ )
474474+ .fetch_optional(&mut **tx)
475475+ .await?;
476476+477477+ let res = sqlx::query!(
478478+ "DELETE FROM likes WHERE did = ? AND rkey = ?",
479479+ e.did,
480480+ e.rkey
481481+ )
482482+ .execute(&mut **tx)
483483+ .await?;
484484+485485+ if res.rows_affected() == 1
486486+ && let Some(row) = old
487487+ {
488488+ sqlx::query!(
489489+ "UPDATE content_stats SET like_count = MAX(0, like_count - 1) WHERE uri = ?",
490490+ row.subject_uri,
491491+ )
492492+ .execute(&mut **tx)
493493+ .await?;
494494+ }
495495+ Ok(())
496496+}
497497+498498+async fn project_save_create(
499499+ tx: &mut Transaction<'_, Sqlite>,
500500+ e: &ProjectionEvent,
501501+) -> anyhow::Result<()> {
502502+ let Some(rec) = &e.record else {
503503+ return Ok(());
504504+ };
505505+ let Some(subject) = ref_uri(rec, "subject") else {
506506+ return Ok(());
507507+ };
508508+ let note = s_str(rec, "note");
509509+ let created_at = parse_millis(rec, "createdAt").unwrap_or_else(now_millis);
510510+511511+ let res = sqlx::query!(
512512+ "INSERT OR IGNORE INTO saves (did, rkey, subject_uri, note, created_at) \
513513+ VALUES (?, ?, ?, ?, ?)",
514514+ e.did,
515515+ e.rkey,
516516+ subject,
517517+ note,
518518+ created_at,
519519+ )
520520+ .execute(&mut **tx)
521521+ .await?;
522522+523523+ if res.rows_affected() == 1 {
524524+ sqlx::query!(
525525+ "INSERT INTO content_stats (uri, like_count, save_count, tag_count) VALUES (?, 0, 1, 0) \
526526+ ON CONFLICT(uri) DO UPDATE SET save_count = save_count + 1",
527527+ subject,
528528+ )
529529+ .execute(&mut **tx)
530530+ .await?;
531531+ }
532532+ Ok(())
533533+}
534534+535535+async fn project_save_delete(
536536+ tx: &mut Transaction<'_, Sqlite>,
537537+ e: &ProjectionEvent,
538538+) -> anyhow::Result<()> {
539539+ let old: Option<SubjectRow> = sqlx::query_as!(
540540+ SubjectRow,
541541+ "SELECT subject_uri FROM saves WHERE did = ? AND rkey = ?",
542542+ e.did,
543543+ e.rkey
544544+ )
545545+ .fetch_optional(&mut **tx)
546546+ .await?;
547547+548548+ let res = sqlx::query!(
549549+ "DELETE FROM saves WHERE did = ? AND rkey = ?",
550550+ e.did,
551551+ e.rkey
552552+ )
553553+ .execute(&mut **tx)
554554+ .await?;
555555+556556+ if res.rows_affected() == 1
557557+ && let Some(row) = old
558558+ {
559559+ sqlx::query!(
560560+ "UPDATE content_stats SET save_count = MAX(0, save_count - 1) WHERE uri = ?",
561561+ row.subject_uri,
562562+ )
563563+ .execute(&mut **tx)
564564+ .await?;
565565+ }
566566+ Ok(())
567567+}
568568+569569+async fn project_tag_create(
570570+ tx: &mut Transaction<'_, Sqlite>,
571571+ e: &ProjectionEvent,
572572+) -> anyhow::Result<()> {
573573+ let Some(rec) = &e.record else {
574574+ return Ok(());
575575+ };
576576+ let Some(subject) = ref_uri(rec, "subject") else {
577577+ return Ok(());
578578+ };
579579+ let Some(tag) = s_str(rec, "tag") else {
580580+ return Ok(());
581581+ };
582582+ let created_at = parse_millis(rec, "createdAt").unwrap_or_else(now_millis);
583583+584584+ let res = sqlx::query!(
585585+ "INSERT OR IGNORE INTO tags (did, rkey, subject_uri, tag, created_at) \
586586+ VALUES (?, ?, ?, ?, ?)",
587587+ e.did,
588588+ e.rkey,
589589+ subject,
590590+ tag,
591591+ created_at,
592592+ )
593593+ .execute(&mut **tx)
594594+ .await?;
595595+596596+ if res.rows_affected() == 1 {
597597+ sqlx::query!(
598598+ "INSERT INTO content_stats (uri, like_count, save_count, tag_count) VALUES (?, 0, 0, 1) \
599599+ ON CONFLICT(uri) DO UPDATE SET tag_count = tag_count + 1",
600600+ subject,
601601+ )
602602+ .execute(&mut **tx)
603603+ .await?;
604604+ }
605605+ Ok(())
606606+}
607607+608608+async fn project_tag_delete(
609609+ tx: &mut Transaction<'_, Sqlite>,
610610+ e: &ProjectionEvent,
611611+) -> anyhow::Result<()> {
612612+ let old: Option<SubjectRow> = sqlx::query_as!(
613613+ SubjectRow,
614614+ "SELECT subject_uri FROM tags WHERE did = ? AND rkey = ?",
615615+ e.did,
616616+ e.rkey
617617+ )
618618+ .fetch_optional(&mut **tx)
619619+ .await?;
620620+621621+ let res = sqlx::query!("DELETE FROM tags WHERE did = ? AND rkey = ?", e.did, e.rkey)
622622+ .execute(&mut **tx)
623623+ .await?;
624624+625625+ if res.rows_affected() == 1
626626+ && let Some(row) = old
627627+ {
628628+ sqlx::query!(
629629+ "UPDATE content_stats SET tag_count = MAX(0, tag_count - 1) WHERE uri = ?",
630630+ row.subject_uri,
631631+ )
632632+ .execute(&mut **tx)
633633+ .await?;
634634+ }
635635+ Ok(())
636636+}
637637+638638+async fn project_listitem_create(
639639+ tx: &mut Transaction<'_, Sqlite>,
640640+ e: &ProjectionEvent,
641641+) -> anyhow::Result<()> {
642642+ let Some(rec) = &e.record else {
643643+ return Ok(());
644644+ };
645645+ let Some(list_uri) = ref_uri(rec, "list") else {
646646+ return Ok(());
647647+ };
648648+ let Some(subject_uri) = ref_uri(rec, "subject") else {
649649+ return Ok(());
650650+ };
651651+ let created_at = parse_millis(rec, "createdAt").unwrap_or_else(now_millis);
652652+653653+ // No denormalized count for listitems; list membership is queried directly.
654654+ sqlx::query!(
655655+ "INSERT OR IGNORE INTO listitems (did, rkey, list_uri, subject_uri, created_at) \
656656+ VALUES (?, ?, ?, ?, ?)",
657657+ e.did,
658658+ e.rkey,
659659+ list_uri,
660660+ subject_uri,
661661+ created_at,
662662+ )
663663+ .execute(&mut **tx)
664664+ .await?;
665665+ Ok(())
666666+}
667667+668668+async fn project_listitem_delete(
669669+ tx: &mut Transaction<'_, Sqlite>,
670670+ e: &ProjectionEvent,
671671+) -> anyhow::Result<()> {
672672+ sqlx::query!(
673673+ "DELETE FROM listitems WHERE did = ? AND rkey = ?",
674674+ e.did,
675675+ e.rkey
676676+ )
677677+ .execute(&mut **tx)
678678+ .await?;
679679+ Ok(())
680680+}
681681+682682+#[cfg(test)]
683683+mod tests {
684684+ //! These tests exercise the projection against a real (in-memory) SQLite
685685+ //! database with migrations applied. Production SQL is compile-time checked
686686+ //! via `query!`; these assertions use runtime queries so the committed
687687+ //! `.sqlx` cache stays focused on production SQL.
688688+689689+ use super::*;
690690+ use serde_json::json;
691691+ use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions};
692692+ use std::str::FromStr;
693693+694694+ /// A fresh in-memory SQLite pool with migrations applied.
695695+ ///
696696+ /// `max_connections(1)` so every query hits the same `:memory:` database.
697697+ async fn test_db() -> SqlitePool {
698698+ let options =
699699+ SqliteConnectOptions::from_str("sqlite::memory:").expect("parse in-memory url");
700700+ let pool = SqlitePoolOptions::new()
701701+ .max_connections(1)
702702+ .connect_with(options)
703703+ .await
704704+ .expect("connect in-memory pool");
705705+ sqlx::migrate!("./migrations")
706706+ .run(&pool)
707707+ .await
708708+ .expect("apply migrations");
709709+ pool
710710+ }
711711+712712+ fn thing_event(
713713+ seq: u64,
714714+ did: &str,
715715+ rkey: &str,
716716+ name: &str,
717717+ models: &[&str],
718718+ ) -> ProjectionEvent {
719719+ ProjectionEvent {
720720+ seq,
721721+ did: did.to_string(),
722722+ collection: "space.polymodel.library.thing".to_string(),
723723+ rkey: rkey.to_string(),
724724+ action: "create".to_string(),
725725+ record: Some(json!({
726726+ "name": name,
727727+ "summary": "a thing summary",
728728+ "tags": ["cube", "calibration"],
729729+ "createdAt": "2024-01-15T12:00:00.000Z",
730730+ "models": models.iter().map(|m| json!({"uri": m, "cid": "c"})).collect::<Vec<_>>(),
731731+ })),
732732+ cid: Some("thingcid".to_string()),
733733+ }
734734+ }
735735+736736+ fn like_event(seq: u64, did: &str, rkey: &str, subject: &str) -> ProjectionEvent {
737737+ ProjectionEvent {
738738+ seq,
739739+ did: did.to_string(),
740740+ collection: "space.polymodel.graph.like".to_string(),
741741+ rkey: rkey.to_string(),
742742+ action: "create".to_string(),
743743+ record: Some(json!({
744744+ "subject": {"uri": subject, "cid": "c"},
745745+ "createdAt": "2024-02-01T00:00:00.000Z",
746746+ })),
747747+ cid: Some("likecid".to_string()),
748748+ }
749749+ }
750750+751751+ fn delete_event(seq: u64, did: &str, collection: &str, rkey: &str) -> ProjectionEvent {
752752+ ProjectionEvent {
753753+ seq,
754754+ did: did.to_string(),
755755+ collection: collection.to_string(),
756756+ rkey: rkey.to_string(),
757757+ action: "delete".to_string(),
758758+ record: None,
759759+ cid: None,
760760+ }
761761+ }
762762+763763+ #[derive(sqlx::FromRow)]
764764+ struct CountRow {
765765+ count: i64,
766766+ }
767767+768768+ #[derive(sqlx::FromRow)]
769769+ struct OptValue {
770770+ value: Option<String>,
771771+ }
772772+773773+ async fn stat(pool: &SqlitePool, uri: &str, column: &str) -> i64 {
774774+ // `column` is a hardcoded test literal (like_count / save_count), never
775775+ // user input, so AssertSqlSafe is appropriate here. Production queries
776776+ // use static `query!` literals.
777777+ let sql = format!("SELECT COALESCE({column}, 0) AS count FROM content_stats WHERE uri = ?");
778778+ sqlx::query_as::<_, CountRow>(sqlx::AssertSqlSafe(sql.as_str()))
779779+ .bind(uri)
780780+ .fetch_optional(pool)
781781+ .await
782782+ .expect("stat query")
783783+ .map(|r| r.count)
784784+ .unwrap_or(0)
785785+ }
786786+787787+ async fn likes_table_count(pool: &SqlitePool, subject: &str) -> i64 {
788788+ sqlx::query_as::<_, CountRow>("SELECT COUNT(*) AS count FROM likes WHERE subject_uri = ?")
789789+ .bind(subject)
790790+ .fetch_one(pool)
791791+ .await
792792+ .expect("likes count")
793793+ .count
794794+ }
795795+796796+ async fn junction_count(pool: &SqlitePool, thing_uri: &str) -> i64 {
797797+ sqlx::query_as::<_, CountRow>(
798798+ "SELECT COUNT(*) AS count FROM thing_models WHERE thing_uri = ?",
799799+ )
800800+ .bind(thing_uri)
801801+ .fetch_one(pool)
802802+ .await
803803+ .expect("junction count")
804804+ .count
805805+ }
806806+807807+ async fn fts_uris(pool: &SqlitePool, term: &str) -> Vec<String> {
808808+ #[derive(sqlx::FromRow)]
809809+ struct UriRow {
810810+ uri: String,
811811+ }
812812+ sqlx::query_as::<_, UriRow>(
813813+ "SELECT t.uri FROM things_fts \
814814+ JOIN things t ON t.rowid = things_fts.rowid \
815815+ WHERE things_fts MATCH ?",
816816+ )
817817+ .bind(term)
818818+ .fetch_all(pool)
819819+ .await
820820+ .expect("fts query")
821821+ .into_iter()
822822+ .map(|r| r.uri)
823823+ .collect()
824824+ }
825825+826826+ async fn cursor(pool: &SqlitePool) -> Option<u64> {
827827+ sqlx::query_as::<_, OptValue>("SELECT value FROM projection_state WHERE key = 'cursor'")
828828+ .fetch_one(pool)
829829+ .await
830830+ .expect("cursor query")
831831+ .value
832832+ .and_then(|v| v.parse().ok())
833833+ }
834834+835835+ #[tokio::test]
836836+ async fn thing_create_projects_row_junction_stats_and_fts() {
837837+ let db = test_db().await;
838838+ let did = "did:plc:abc";
839839+ let uri = "at://did:plc:abc/space.polymodel.library.thing/r1";
840840+ let e = thing_event(10, did, "r1", "Calibration Cube", &["at://m1", "at://m2"]);
841841+ project_event(&db, &e).await.expect("project thing");
842842+843843+ assert_eq!(junction_count(&db, uri).await, 2);
844844+ assert_eq!(stat(&db, uri, "like_count").await, 0); // stats row exists, zeroed
845845+ assert!(fts_uris(&db, "cube").await.contains(&uri.to_string()));
846846+ assert!(
847847+ fts_uris(&db, "calibration")
848848+ .await
849849+ .contains(&uri.to_string())
850850+ );
851851+ assert_eq!(cursor(&db).await, Some(10));
852852+ }
853853+854854+ #[tokio::test]
855855+ async fn thing_update_rebuilds_junction() {
856856+ let db = test_db().await;
857857+ let did = "did:plc:abc";
858858+ let uri = "at://did:plc:abc/space.polymodel.library.thing/r1";
859859+860860+ // Create with two models.
861861+ project_event(
862862+ &db,
863863+ &thing_event(1, did, "r1", "Cube", &["at://m1", "at://m2"]),
864864+ )
865865+ .await
866866+ .unwrap();
867867+ assert_eq!(junction_count(&db, uri).await, 2);
868868+869869+ // Update: same rkey, different models. Old junction rows must be replaced.
870870+ let mut upd = thing_event(2, did, "r1", "Cube v2", &["at://m3"]);
871871+ upd.action = "update".to_string();
872872+ project_event(&db, &upd).await.unwrap();
873873+874874+ assert_eq!(junction_count(&db, uri).await, 1);
875875+ // m1/m2 gone, m3 present.
876876+ let has_m3 = sqlx::query_as::<_, CountRow>(
877877+ "SELECT COUNT(*) AS count FROM thing_models WHERE thing_uri = ? AND model_uri = 'at://m3'",
878878+ )
879879+ .bind(uri)
880880+ .fetch_one(&db)
881881+ .await
882882+ .unwrap()
883883+ .count;
884884+ assert_eq!(has_m3, 1);
885885+ }
886886+887887+ #[tokio::test]
888888+ async fn thing_delete_cascades() {
889889+ let db = test_db().await;
890890+ let did = "did:plc:abc";
891891+ let uri = "at://did:plc:abc/space.polymodel.library.thing/r1";
892892+ project_event(&db, &thing_event(1, did, "r1", "Cube", &["at://m1"]))
893893+ .await
894894+ .unwrap();
895895+ assert!(!fts_uris(&db, "cube").await.is_empty());
896896+897897+ project_event(
898898+ &db,
899899+ &delete_event(2, did, "space.polymodel.library.thing", "r1"),
900900+ )
901901+ .await
902902+ .unwrap();
903903+904904+ assert_eq!(junction_count(&db, uri).await, 0);
905905+ assert_eq!(stat(&db, uri, "like_count").await, 0); // stats row removed
906906+ assert!(fts_uris(&db, "cube").await.is_empty()); // FTS delete trigger fired
907907+ }
908908+909909+ #[tokio::test]
910910+ async fn like_create_increments_count() {
911911+ let db = test_db().await;
912912+ let subject = "at://did:plc:author/.../thing";
913913+ project_event(&db, &like_event(1, "did:plc:a", "k1", subject))
914914+ .await
915915+ .unwrap();
916916+ project_event(&db, &like_event(2, "did:plc:b", "k2", subject))
917917+ .await
918918+ .unwrap();
919919+920920+ assert_eq!(likes_table_count(&db, subject).await, 2);
921921+ assert_eq!(stat(&db, subject, "like_count").await, 2);
922922+ }
923923+924924+ #[tokio::test]
925925+ async fn like_delete_decrements_count_two_step() {
926926+ let db = test_db().await;
927927+ let subject = "at://did:plc:author/.../thing";
928928+ project_event(&db, &like_event(1, "did:plc:a", "k1", subject))
929929+ .await
930930+ .unwrap();
931931+ project_event(
932932+ &db,
933933+ &delete_event(2, "did:plc:a", "space.polymodel.graph.like", "k1"),
934934+ )
935935+ .await
936936+ .unwrap();
937937+938938+ assert_eq!(likes_table_count(&db, subject).await, 0);
939939+ assert_eq!(stat(&db, subject, "like_count").await, 0);
940940+ }
941941+942942+ #[tokio::test]
943943+ async fn like_delete_when_absent_is_idempotent() {
944944+ // Deleting a like that was never there must not panic nor go negative.
945945+ let db = test_db().await;
946946+ let subject = "at://did:plc:author/.../thing";
947947+ project_event(
948948+ &db,
949949+ &delete_event(1, "did:plc:a", "space.polymodel.graph.like", "nope"),
950950+ )
951951+ .await
952952+ .expect("delete absent like must succeed");
953953+ assert_eq!(stat(&db, subject, "like_count").await, 0);
954954+ }
955955+956956+ #[tokio::test]
957957+ async fn duplicate_like_does_not_double_count() {
958958+ // Replay: same like event twice (e.g. eager write + firehose delivery).
959959+ let db = test_db().await;
960960+ let subject = "at://did:plc:author/.../thing";
961961+ project_event(&db, &like_event(1, "did:plc:a", "k1", subject))
962962+ .await
963963+ .unwrap();
964964+ project_event(&db, &like_event(2, "did:plc:a", "k1", subject))
965965+ .await
966966+ .unwrap();
967967+968968+ assert_eq!(likes_table_count(&db, subject).await, 1);
969969+ assert_eq!(stat(&db, subject, "like_count").await, 1);
970970+ }
971971+972972+ #[tokio::test]
973973+ async fn cursor_advances_with_each_event() {
974974+ let db = test_db().await;
975975+ project_event(&db, &thing_event(7, "did:plc:a", "r1", "Cube", &[]))
976976+ .await
977977+ .unwrap();
978978+ assert_eq!(cursor(&db).await, Some(7));
979979+ project_event(&db, &like_event(99, "did:plc:b", "k1", "at://x"))
980980+ .await
981981+ .unwrap();
982982+ assert_eq!(cursor(&db).await, Some(99));
983983+ }
984984+985985+ #[tokio::test]
986986+ async fn save_create_and_delete_maintain_count() {
987987+ let db = test_db().await;
988988+ let subject = "at://did:plc:author/.../thing";
989989+ let save = ProjectionEvent {
990990+ seq: 1,
991991+ did: "did:plc:a".to_string(),
992992+ collection: "space.polymodel.graph.save".to_string(),
993993+ rkey: "s1".to_string(),
994994+ action: "create".to_string(),
995995+ record: Some(json!({
996996+ "subject": {"uri": subject, "cid": "c"},
997997+ "note": "useful",
998998+ "createdAt": "2024-03-01T00:00:00.000Z",
999999+ })),
10001000+ cid: Some("savecid".to_string()),
10011001+ };
10021002+ project_event(&db, &save).await.unwrap();
10031003+ assert_eq!(stat(&db, subject, "save_count").await, 1);
10041004+10051005+ project_event(
10061006+ &db,
10071007+ &delete_event(2, "did:plc:a", "space.polymodel.graph.save", "s1"),
10081008+ )
10091009+ .await
10101010+ .unwrap();
10111011+ assert_eq!(stat(&db, subject, "save_count").await, 0);
10121012+ }
10131013+10141014+ #[tokio::test]
10151015+ async fn thing_update_keeps_fts_consistent() {
10161016+ // Regression for UPSERT vs INSERT OR REPLACE: updating a thing must fire
10171017+ // the AFTER UPDATE FTS trigger so the old term is removed and the new one
10181018+ // added, with no orphaned FTS row. INSERT OR REPLACE leaked one orphan per
10191019+ // update because REPLACE is DELETE+new-rowid (UPDATE trigger never fires
10201020+ // and, with recursive_triggers off, the DELETE trigger doesn't either).
10211021+ let db = test_db().await;
10221022+ let did = "did:plc:abc";
10231023+ let uri = "at://did:plc:abc/space.polymodel.library.thing/r1";
10241024+10251025+ let create = ProjectionEvent {
10261026+ seq: 1,
10271027+ did: did.to_string(),
10281028+ collection: "space.polymodel.library.thing".to_string(),
10291029+ rkey: "r1".to_string(),
10301030+ action: "create".to_string(),
10311031+ record: Some(json!({
10321032+ "name": "Calibration Cube",
10331033+ "tags": ["cube"],
10341034+ "createdAt": "2024-01-01T00:00:00.000Z",
10351035+ "models": [],
10361036+ })),
10371037+ cid: Some("cid1".to_string()),
10381038+ };
10391039+ project_event(&db, &create).await.unwrap();
10401040+ assert_eq!(fts_uris(&db, "cube").await, vec![uri.to_string()]);
10411041+10421042+ let mut update = create.clone();
10431043+ update.seq = 2;
10441044+ update.action = "update".to_string();
10451045+ update.cid = Some("cid2".to_string());
10461046+ update.record = Some(json!({
10471047+ "name": "Friendly Sphere",
10481048+ "tags": ["sphere"],
10491049+ "createdAt": "2024-01-02T00:00:00.000Z",
10501050+ "models": [],
10511051+ }));
10521052+ project_event(&db, &update).await.unwrap();
10531053+10541054+ // Old term must be gone (no orphan), new term present, single live row.
10551055+ assert!(
10561056+ fts_uris(&db, "cube").await.is_empty(),
10571057+ "stale FTS entry for the old term leaked across update"
10581058+ );
10591059+ assert_eq!(fts_uris(&db, "sphere").await, vec![uri.to_string()]);
10601060+ let live = sqlx::query_as::<_, CountRow>(
10611061+ "SELECT COUNT(*) AS count FROM things WHERE did = ? AND rkey = ?",
10621062+ )
10631063+ .bind(did)
10641064+ .bind("r1")
10651065+ .fetch_one(&db)
10661066+ .await
10671067+ .unwrap()
10681068+ .count;
10691069+ assert_eq!(live, 1);
10701070+ }
10711071+}
+38
src/indexing/setup.rs
···11+//! Hydrant ATProto indexer setup and configuration.
22+//!
33+//! Initializes Hydrant from environment and applies a signal-mode filter so only
44+//! repos that publish `space.polymodel.*` records are indexed. The returned
55+//! handle is cheaply cloneable; the firehose/backfill driver is started later via
66+//! [`Hydrant::run`] in [`super::start_indexing`].
77+88+use hydrant::FilterMode;
99+use hydrant::control::Hydrant;
1010+1111+/// NSID patterns to index: the whole Polymodel namespace, used for both
1212+/// signal-mode discovery and the collection allowlist.
1313+const COLLECTIONS: [&str; 1] = ["space.polymodel.*"];
1414+1515+/// Initialize Hydrant from environment and configure signal-mode filtering.
1616+///
1717+/// Does NOT start the firehose. The caller drives [`Hydrant::run`] (started in
1818+/// [`super::start_indexing`]) to begin ingest and backfill.
1919+pub async fn init_hydrant() -> anyhow::Result<Hydrant> {
2020+ let hydrant = Hydrant::from_env()
2121+ .await
2222+ .map_err(|e| anyhow::anyhow!("hydrant init failed: {e}"))?;
2323+2424+ // Signal mode: only index repos that emit space.polymodel.* records, and
2525+ // only store those collections. Code takes precedence over HYDRANT_FILTER_*
2626+ // env vars, but the env vars are still documented in .env.example.
2727+ hydrant
2828+ .filter
2929+ .set_mode(FilterMode::Filter)
3030+ .set_signals(COLLECTIONS)
3131+ .set_collections(COLLECTIONS)
3232+ .apply()
3333+ .await
3434+ .map_err(|e| anyhow::anyhow!("hydrant filter apply failed: {e}"))?;
3535+3636+ tracing::info!("hydrant initialized with signal filter for space.polymodel.*");
3737+ Ok(hydrant)
3838+}
+31
src/main.rs
···33mod mesh;
44mod viewer;
5566+// Frontend compile-time configuration, generated by build.rs from POLYMODEL_*.
77+mod env;
88+// Server-only indexing infrastructure (Hydrant firehose + SQLite projection).
99+#[cfg(feature = "server")]
1010+mod indexing;
1111+612use viewer::{ViewerPage, demo_model_for_route, demo_models};
713814const FAVICON: Asset = asset!("/assets/favicon.jpg");
···1218const BROWSE_CSS: Asset = asset!("/assets/styling/browse.css");
1319const POLYMODEL_CSS: Asset = asset!("/assets/styling/polymodel.css");
14202121+#[cfg(feature = "server")]
2222+fn main() {
2323+ // The closure runs inside the server's Tokio runtime at startup, so the
2424+ // indexing tasks spawned here live for the process lifetime alongside Axum.
2525+ dioxus::serve(|| async {
2626+ // Server runtime config: load .env, then read DATABASE_URL (and any
2727+ // future server env) in one place. Hydrant reads its own HYDRANT_*.
2828+ let _ = dotenvy::dotenv();
2929+ let _ = tracing_subscriber::fmt()
3030+ .with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
3131+ .try_init();
3232+3333+ let cfg = indexing::config::ServerConfig::load();
3434+ let db = indexing::db::init_db(&cfg)
3535+ .await
3636+ .expect("failed to initialize sqlite projection database");
3737+ indexing::start_indexing(db)
3838+ .await
3939+ .expect("failed to start indexing pipeline");
4040+4141+ Ok(dioxus::server::router(App))
4242+ });
4343+}
4444+4545+#[cfg(not(feature = "server"))]
1546fn main() {
1647 #[cfg(all(target_family = "wasm", target_os = "unknown"))]
1748 {