nandi/gleanpublic⑂ Fork 0
⑂ a33920b
Commits
⬇ Clone ▾
git clone https://git.rickub.com/nandi/glean.git
git clone ssh://git@rickub.com/nandi/glean.git

Host key fingerprint (ed25519): SHA256:iycHnxEyq0Q7uyVpB7JlznP0G7JrTPXLYRcAU5CSLhc — verify it before your first connect.

Add configurable sync and cluster intervalsUnverified

Julien Robert committed 2026-04-21T10:01:10+02:00 Browse files
a33920b parent: aea953a
modified .env.example +2 -0
@@ -1,6 +1,8 @@
11 GLEAN_ADDR=:8080
22 GLEAN_DB=glean.db
33 GLEAN_JETSTREAM=wss://jetstream2.fr.hose.cam
4+GLEAN_SYNC_INTERVAL=1h
5+GLEAN_CLUSTER_INTERVAL=6h
46 # Leave empty for localhost OAuth (development)
57 # GLEAN_OAUTH_CLIENT_ID=https://glean.at/oauth/client-metadata
68 # GLEAN_OAUTH_REDIRECT_URL=https://glean.at/auth/callback
@@ -1,6 +1,8 @@
1 GLEAN_ADDR=:80801 GLEAN_ADDR=:8080
2 GLEAN_DB=glean.db2 GLEAN_DB=glean.db
3 GLEAN_JETSTREAM=wss://jetstream2.fr.hose.cam3 GLEAN_JETSTREAM=wss://jetstream2.fr.hose.cam
4+GLEAN_SYNC_INTERVAL=1h
5+GLEAN_CLUSTER_INTERVAL=6h
4 # Leave empty for localhost OAuth (development)6 # Leave empty for localhost OAuth (development)
5 # GLEAN_OAUTH_CLIENT_ID=https://glean.at/oauth/client-metadata7 # GLEAN_OAUTH_CLIENT_ID=https://glean.at/oauth/client-metadata
6 # GLEAN_OAUTH_REDIRECT_URL=https://glean.at/auth/callback8 # GLEAN_OAUTH_REDIRECT_URL=https://glean.at/auth/callback
modified docs/specs.md +80 -80
@@ -2,22 +2,22 @@
22
33 ## 1. Overview
44
5-Glean is a social RSS reader built on the AT Protocol. It operates as an **AppView** for the `at.glean.*` lexicon namespace: it indexes records from the relay firehose, serves XRPC query endpoints, and provides the web UI at [glean.at](https://glean.at).
5+Glean is a social RSS reader built on the AT Protocol. It operates as an **AppView** for the `at.glean.*` lexicon namespace: it indexes records from Jetstream, serves XRPC query endpoints, and provides the web UI at [glean.at](https://glean.at).
66
7-Users store their RSS feed subscriptions as individual lexicon records on their PDS (one record per feed). Glean's AppView consumes the firehose, indexes those records, fetches the referenced RSS feeds, and serves both the reader UI and public XRPC APIs for the `at.glean.*` namespace.
7+Users store their RSS feed subscriptions as individual lexicon records on their PDS (one record per feed). Glean's AppView consumes Jetstream, indexes those records, fetches the referenced RSS feeds, and serves both the reader UI and public XRPC APIs for the `at.glean.*` namespace.
88
99 The core idea: your RSS subscriptions are a strong signal about your interests. When enough people expose theirs, you can discover both **people** (who reads the same things) and **content** (what similar readers follow that you don't).
1010
1111 ## 2. Stack
1212
13-| Layer | Technology |
14-| ---------------- | ---------------------------------- |
15-| Backend | Go |
16-| Database | SQLite (via `mattn/go-sqlite3`) |
17-| Frontend | htmx + TailwindCSS |
18-| Auth | AT Protocol OAuth / DID resolution |
19-| AT Protocol role | AppView for `at.glean.*` lexicons |
20-| Data source | AT Relay firehose → SQLite index |
13+| Layer | Technology |
14+| ---------------- | ------------------------------------ |
15+| Backend | Go |
16+| Database | SQLite (via `mattn/go-sqlite3`) |
17+| Frontend | htmx + TailwindCSS |
18+| Auth | AT Protocol OAuth / DID resolution |
19+| AT Protocol role | AppView for `at.glean.*` lexicons |
20+| Data source | AT Protocol Jetstream → SQLite index |
2121
2222 ## 3. AT Protocol Lexicons
2323
@@ -217,9 +217,9 @@ Output:
217217 people: [{ did, handle, displayName, avatar, jaccard, commonFeeds }]
218218 ```
219219
220-### 3.5 AppView Firehose Consumption
220+### 3.5 AppView Jetstream Consumption
221221
222-Glean subscribes to the AT Relay firehose (`wss://bsky.network`) for all `at.glean.*` records:
222+Glean subscribes to a Jetstream endpoint (`GLEAN_JETSTREAM`, default `wss://jetstream2.fr.hose.cam`) for all `at.glean.*` records:
223223
224224 ```
225225 SUBSCRIBE collections: ["at.glean.subscription", "at.glean.annotation", "at.glean.like"]
@@ -231,7 +231,7 @@ On each event:
231231 - **delete**: Tombstone the record (soft delete to preserve foreign key integrity)
232232 - **update**: Replace the record's CID and value
233233
234-The AppView does not handle writes. Users write records to their own PDS. Glean only reads them from the firehose.
234+The AppView does not handle writes. Users write records to their own PDS. Glean only reads them from Jetstream.
235235
236236 ## 4. RSS Reader
237237
@@ -384,43 +384,43 @@ Beyond the clustering system, Glean also discovers new feeds from article conten
384384
385385 ## 5. System Architecture
386386
387-Glean runs as a single Go binary that fills three roles: **AppView** (indexing `at.glean.*` records from the firehose, serving XRPC queries), **RSS reader** (fetching and storing feed content), and **web UI** (htmx frontend).
388-
389-```
390- AT Relay (bsky.network)
391- │ firehose
392- ▼
393- ┌─────────────────────┐
394- │ Go Server (glean.at)│
395- │ │
396- Browser ──HTTP──► │ ┌────────────────┐ │ ──XRPC queries──► Other AT apps
397- (htmx + TW) │ │ Router │ │
398- │ │ ┌───────────┐ │ │
399- │ │ │ Handlers │ │ │
400- │ │ │ (UI + XRPC)│ │ │
401- │ │ └─────┬─────┘ │ │
402- │ └────────┼────────┘ │
403- │ │ │
404- │ ┌────────▼────────┐ │ ┌──────────────────┐
405- │ │ Service Layer │ │ │ Feed Scheduler │
406- │ │ │──┼──sync──►│ (goroutine) │
407- │ └────────┬────────┘ │ │ Fetcher + Parser│
408- │ │ │ └────────┬─────────┘
409- │ ┌────────▼────────┐ │ │
410- │ │ SQLite │ │ RSS/Atom/JSON feeds
411- │ │ (firehose idx, │ │
412- │ │ articles, │ │ ┌──────────────────┐
413- │ │ read state, │ │ │ Cluster Engine │
414- │ │ clustering) │◄─┼────────►│ (periodic cron) │
415- │ └─────────────────┘ │ └──────────────────┘
416- └──────────────────────┘
417-
418- AppView responsibilities:
419- • Subscribe to firehose for at.glean.subscription, at.glean.annotation, at.glean.like
420- • Index records into SQLite
421- • Serve XRPC query endpoints (at.glean.listSubscriptions, etc.)
422- • Host the web UI at glean.at
423- • Write to user PDS on behalf of user (when user acts through UI)
387+Glean runs as a single Go binary that fills three roles: **AppView** (indexing `at.glean.*` records from Jetstream, serving XRPC queries), **RSS reader** (fetching and storing feed content), and **web UI** (htmx frontend).
388+
389+```
390+ Jetstream (GLEAN_JETSTREAM)
391+ │ subscribe
392+ ▼
393+ ┌─────────────────────┐
394+ │ Go Server (glean.at)│
395+ │ │
396+ Browser ──HTTP──► │ ┌────────────────┐ │ ──XRPC queries──► Other AT apps
397+ (htmx + TW) │ │ Router │ │
398+ │ │ ┌───────────┐ │ │
399+ │ │ │ Handlers │ │ │
400+ │ │ │ (UI + XRPC)│ │ │
401+ │ │ └─────┬─────┘ │ │
402+ │ └────────┼────────┘ │
403+ │ │ │
404+ │ ┌────────▼────────┐ │ ┌──────────────────┐
405+ │ │ Service Layer │ │ │ Feed Scheduler │
406+ │ │ │──┼──sync──►│ (goroutine) │
407+ │ └────────┬────────┘ │ │ Fetcher + Parser│
408+ │ │ │ └────────┬─────────┘
409+ │ ┌────────▼────────┐ │ │
410+ │ │ SQLite │ │ RSS/Atom/JSON feeds
411+ │ │ (jetstream idx, │ │
412+ │ │ articles, │ │ ┌──────────────────┐
413+ │ │ read state, │ │ │ Cluster Engine │
414+ │ │ clustering) │◄─┼────────►│ (periodic cron) │
415+ │ └─────────────────┘ │ └──────────────────┘
416+ └──────────────────────┘
417+
418+ AppView responsibilities:
419+ • Subscribe to Jetstream for at.glean.subscription, at.glean.annotation, at.glean.like
420+ • Index records into SQLite
421+ • Serve XRPC query endpoints (at.glean.listSubscriptions, etc.)
422+ • Host the web UI at glean.at
423+ • Write to user PDS on behalf of user (when user acts through UI)
424424 ```
425425
426426 ## 6. Database Schema (SQLite)
@@ -525,6 +525,7 @@ CREATE INDEX idx_read_state_starred ON read_state(user_did, is_starred) WHERE is
525525 ```
526526
527527 ### 6.6 Annotations, Likes
528+
528529 Local mirror of AT Protocol lexicon records for fast querying.
529530
530531 ```sql
@@ -654,9 +655,9 @@ For larger scale, move to MinHash + LSH (banded hashing) to approximate Jaccard
654655
655656 ### 7.5 Clustering Engine (Cron)
656657
657-A background goroutine runs on a schedule (e.g., every 6 hours):
658+A background goroutine runs on a configurable schedule (`GLEAN_CLUSTER_INTERVAL`, default 6h):
658659
659-1. **Firehose ingestion**: Subscribe to AT Protocol firehose for `at.glean.*` records
660+1. **Jetstream ingestion**: Subscribe to Jetstream for `at.glean.*` records
660661 2. **Index new records**: Parse lexicon records, upsert into SQLite
661662 3. **Compute similarities**: Batch-update the `feed_similarity`, `user_similarity`, and `article_co_like` tables
662663 4. **Generate recommendations**: Materialize top recommendations per user into cache tables
@@ -686,28 +687,28 @@ The server renders HTML fragments that htmx swaps into the page. No JSON API nee
686687
687688 ### 8.1 Pages
688689
689-| Route | Method | Description |
690-| ---------------------- | ------ | -------------------------------------------------------- |
691-| `/` | GET | Landing page / auth redirect |
692-| `/dashboard` | GET | Main dashboard: unread articles, recommendations sidebar |
693-| `/feeds` | GET | Manage RSS subscriptions (OPML import for onboarding) |
694-| `/feeds/opml/upload` | POST | Upload OPML file to bulk-import subscriptions |
695-| `/feeds/opml/download` | GET | Export subscriptions as OPML (offboarding) |
696-| `/feeds/add` | POST | Add a single feed URL |
697-| `/feeds/remove` | DELETE | Remove a feed |
698-| `/feeds/refresh` | POST | Refresh all subscribed feeds |
699-| `/feeds/clear` | POST | Clear all subscriptions |
700-| `/articles` | GET | Read articles (paginated, filterable by feed) |
701-| `/articles/{id}` | GET | Article detail view |
702-| `/articles/{id}/read` | POST | Mark article as read |
703-| `/articles/{id}/unread`| POST | Mark article as unread |
704-| `/articles/{id}/like` | POST | Like an article |
705-| `/articles/mark-all-read` | POST | Mark all articles as read |
706-| `/trending` | GET | Community feed: articles ranked by likes |
707-| `/library` | GET | Liked articles and annotations |
708-| `/library/create` | POST | Create annotation on an article |
709-| `/library/{id}/delete` | POST | Delete an annotation |
710-| `/profile/{did}` | GET | Public profile: their feeds, likes, annotations |
690+| Route | Method | Description |
691+| ------------------------- | ------ | -------------------------------------------------------- |
692+| `/` | GET | Landing page / auth redirect |
693+| `/dashboard` | GET | Main dashboard: unread articles, recommendations sidebar |
694+| `/feeds` | GET | Manage RSS subscriptions (OPML import for onboarding) |
695+| `/feeds/opml/upload` | POST | Upload OPML file to bulk-import subscriptions |
696+| `/feeds/opml/download` | GET | Export subscriptions as OPML (offboarding) |
697+| `/feeds/add` | POST | Add a single feed URL |
698+| `/feeds/remove` | DELETE | Remove a feed |
699+| `/feeds/refresh` | POST | Refresh all subscribed feeds |
700+| `/feeds/clear` | POST | Clear all subscriptions |
701+| `/articles` | GET | Read articles (paginated, filterable by feed) |
702+| `/articles/{id}` | GET | Article detail view |
703+| `/articles/{id}/read` | POST | Mark article as read |
704+| `/articles/{id}/unread` | POST | Mark article as unread |
705+| `/articles/{id}/like` | POST | Like an article |
706+| `/articles/mark-all-read` | POST | Mark all articles as read |
707+| `/trending` | GET | Community feed: articles ranked by likes |
708+| `/library` | GET | Liked articles and annotations |
709+| `/library/create` | POST | Create annotation on an article |
710+| `/library/{id}/delete` | POST | Delete an annotation |
711+| `/profile/{did}` | GET | Public profile: their feeds, likes, annotations |
711712
712713 ### 8.2 htmx Patterns
713714
@@ -729,7 +730,8 @@ glean/
729730 │ ├── atproto/
730731 │ │ ├── auth.go # DID resolution, OAuth flow
731732 │ │ ├── client.go # XRPC client (write to user PDS)
732-│ │ ├── firehose.go # Subscribe to AT Relay firehose
733+│ │ ├── jetstream.go # Subscribe to Jetstream via official client
734+│ │ ├── stream_handler.go # Stream event → DB handler
733735 │ │ ├── sync.go # PDS record reconciliation
734736 │ │ └── xrpc.go # XRPC query handlers (AppView endpoints)
735737 │ ├── db/
@@ -862,7 +864,7 @@ Glean exposes a `/metrics` endpoint for monitoring. Key metrics:
862864 - **`glean_feeds_fetched_total`** — Feed fetch attempts labeled by status (`success`, `error`, `not_modified`)
863865 - **`glean_feed_fetch_duration_seconds`** — Histogram of feed fetch latency
864866 - **`glean_articles_upserted_total`** — Counter of articles stored from feeds
865-- **`glean_firehose_events_total`** — Firehose events labeled by collection and action
867+- **`glean_jetstream_events_total`** — Jetstream events labeled by collection and action
866868 - **`glean_http_requests_total`** — HTTP request counts labeled by method, path, and status
867869 - **`glean_cluster_runs_total`** / **`glean_cluster_duration_seconds`** — Recommendation engine runs and timing
868870
@@ -878,8 +880,8 @@ Glean exposes a `/metrics` endpoint for monitoring. Key metrics:
878880
879881 Glean operates as an AT Protocol AppView. This means:
880882
881-- **Read path**: All `at.glean.*` data is consumed from the Relay firehose, not by polling individual PDS instances. The firehose handler runs as a persistent goroutine, upserting records into SQLite as they arrive.
882-- **Write path**: Users write records to their own PDS (via standard AT Protocol `com.atproto.repo.createRecord` / `deleteRecord`). Glean never stores user data directly — it only indexes what the firehose delivers.
883+- **Read path**: All `at.glean.*` data is consumed from Jetstream, not by polling individual PDS instances. The Jetstream consumer runs as a persistent goroutine, upserting records into SQLite as they arrive.
884+- **Write path**: Users write records to their own PDS (via standard AT Protocol `com.atproto.repo.createRecord` / `deleteRecord`). Glean never stores user data directly — it only indexes what Jetstream delivers.
883885 - **Query path**: Other AT Protocol apps can query Glean's XRPC endpoints to access indexed data (subscriptions, annotations, likes, recommendations) without building their own indexer.
884886 - **Trade-off**: Article content (fetched from RSS feeds) is local-only and not part of the AT Protocol layer. Only individual feed subscription records (`at.glean.subscription`) live on the PDS.
885887
@@ -891,6 +893,4 @@ All PDS records are public. There is no notion of private data on the AT Protoco
891893
892894 - **MinHash/LSH**: Replace brute-force Jaccard when user count exceeds ~50k
893895 - **Full-text search**: Add FTS5 virtual table on articles for search
894-- **Feed groups / reading lists**: Allow users to create curated lists (separate lexicon)
895896 - **Email digest**: Periodic email with top articles from subscribed feeds
896-- **Multi-AppView scaling**: Distribute firehose consumption across multiple instances behind a load balancer
@@ -2,22 +2,22 @@
2 2
3 ## 1. Overview3 ## 1. Overview
4 4
5-Glean is a social RSS reader built on the AT Protocol. It operates as an **AppView** for the `at.glean.*` lexicon namespace: it indexes records from the relay firehose, serves XRPC query endpoints, and provides the web UI at [glean.at](https://glean.at).5+Glean is a social RSS reader built on the AT Protocol. It operates as an **AppView** for the `at.glean.*` lexicon namespace: it indexes records from Jetstream, serves XRPC query endpoints, and provides the web UI at [glean.at](https://glean.at).
6 6
7-Users store their RSS feed subscriptions as individual lexicon records on their PDS (one record per feed). Glean's AppView consumes the firehose, indexes those records, fetches the referenced RSS feeds, and serves both the reader UI and public XRPC APIs for the `at.glean.*` namespace.7+Users store their RSS feed subscriptions as individual lexicon records on their PDS (one record per feed). Glean's AppView consumes Jetstream, indexes those records, fetches the referenced RSS feeds, and serves both the reader UI and public XRPC APIs for the `at.glean.*` namespace.
8 8
9 The core idea: your RSS subscriptions are a strong signal about your interests. When enough people expose theirs, you can discover both **people** (who reads the same things) and **content** (what similar readers follow that you don't).9 The core idea: your RSS subscriptions are a strong signal about your interests. When enough people expose theirs, you can discover both **people** (who reads the same things) and **content** (what similar readers follow that you don't).
10 10
11 ## 2. Stack11 ## 2. Stack
12 12
13-| Layer | Technology |13+| Layer | Technology |
14-| ---------------- | ---------------------------------- |14+| ---------------- | ------------------------------------ |
15-| Backend | Go |15+| Backend | Go |
16-| Database | SQLite (via `mattn/go-sqlite3`) |16+| Database | SQLite (via `mattn/go-sqlite3`) |
17-| Frontend | htmx + TailwindCSS |17+| Frontend | htmx + TailwindCSS |
18-| Auth | AT Protocol OAuth / DID resolution |18+| Auth | AT Protocol OAuth / DID resolution |
19-| AT Protocol role | AppView for `at.glean.*` lexicons |19+| AT Protocol role | AppView for `at.glean.*` lexicons |
20-| Data source | AT Relay firehose → SQLite index |20+| Data source | AT Protocol Jetstream → SQLite index |
21 21
22 ## 3. AT Protocol Lexicons22 ## 3. AT Protocol Lexicons
23 23
@@ -217,9 +217,9 @@ Output:
217 people: [{ did, handle, displayName, avatar, jaccard, commonFeeds }]217 people: [{ did, handle, displayName, avatar, jaccard, commonFeeds }]
218 ```218 ```
219 219
220-### 3.5 AppView Firehose Consumption220+### 3.5 AppView Jetstream Consumption
221 221
222-Glean subscribes to the AT Relay firehose (`wss://bsky.network`) for all `at.glean.*` records:222+Glean subscribes to a Jetstream endpoint (`GLEAN_JETSTREAM`, default `wss://jetstream2.fr.hose.cam`) for all `at.glean.*` records:
223 223
224 ```224 ```
225 SUBSCRIBE collections: ["at.glean.subscription", "at.glean.annotation", "at.glean.like"]225 SUBSCRIBE collections: ["at.glean.subscription", "at.glean.annotation", "at.glean.like"]
@@ -231,7 +231,7 @@ On each event:
231 - **delete**: Tombstone the record (soft delete to preserve foreign key integrity)231 - **delete**: Tombstone the record (soft delete to preserve foreign key integrity)
232 - **update**: Replace the record's CID and value232 - **update**: Replace the record's CID and value
233 233
234-The AppView does not handle writes. Users write records to their own PDS. Glean only reads them from the firehose.234+The AppView does not handle writes. Users write records to their own PDS. Glean only reads them from Jetstream.
235 235
236 ## 4. RSS Reader236 ## 4. RSS Reader
237 237
@@ -384,43 +384,43 @@ Beyond the clustering system, Glean also discovers new feeds from article conten
384 384
385 ## 5. System Architecture385 ## 5. System Architecture
386 386
387-Glean runs as a single Go binary that fills three roles: **AppView** (indexing `at.glean.*` records from the firehose, serving XRPC queries), **RSS reader** (fetching and storing feed content), and **web UI** (htmx frontend).387+Glean runs as a single Go binary that fills three roles: **AppView** (indexing `at.glean.*` records from Jetstream, serving XRPC queries), **RSS reader** (fetching and storing feed content), and **web UI** (htmx frontend).
388-388+
389-```389+```
390- AT Relay (bsky.network)390+ Jetstream (GLEAN_JETSTREAM)
391- │ firehose391+ │ subscribe
392- ▼392+ ▼
393- ┌─────────────────────┐393+ ┌─────────────────────┐
394- │ Go Server (glean.at)│394+ │ Go Server (glean.at)│
395- │ │395+ │ │
396- Browser ──HTTP──► │ ┌────────────────┐ │ ──XRPC queries──► Other AT apps396+ Browser ──HTTP──► │ ┌────────────────┐ │ ──XRPC queries──► Other AT apps
397- (htmx + TW) │ │ Router │ │397+ (htmx + TW) │ │ Router │ │
398- │ │ ┌───────────┐ │ │398+ │ │ ┌───────────┐ │ │
399- │ │ │ Handlers │ │ │399+ │ │ │ Handlers │ │ │
400- │ │ │ (UI + XRPC)│ │ │400+ │ │ │ (UI + XRPC)│ │ │
401- │ │ └─────┬─────┘ │ │401+ │ │ └─────┬─────┘ │ │
402- │ └────────┼────────┘ │402+ │ └────────┼────────┘ │
403- │ │ │403+ │ │ │
404- │ ┌────────▼────────┐ │ ┌──────────────────┐404+ │ ┌────────▼────────┐ │ ┌──────────────────┐
405- │ │ Service Layer │ │ │ Feed Scheduler │405+ │ │ Service Layer │ │ │ Feed Scheduler │
406- │ │ │──┼──sync──►│ (goroutine) │406+ │ │ │──┼──sync──►│ (goroutine) │
407- │ └────────┬────────┘ │ │ Fetcher + Parser│407+ │ └────────┬────────┘ │ │ Fetcher + Parser│
408- │ │ │ └────────┬─────────┘408+ │ │ │ └────────┬─────────┘
409- │ ┌────────▼────────┐ │ │409+ │ ┌────────▼────────┐ │ │
410- │ │ SQLite │ │ RSS/Atom/JSON feeds410+ │ │ SQLite │ │ RSS/Atom/JSON feeds
411- │ │ (firehose idx, │ │411+ │ │ (jetstream idx, │ │
412- │ │ articles, │ │ ┌──────────────────┐412+ │ │ articles, │ │ ┌──────────────────┐
413- │ │ read state, │ │ │ Cluster Engine │413+ │ │ read state, │ │ │ Cluster Engine │
414- │ │ clustering) │◄─┼────────►│ (periodic cron) │414+ │ │ clustering) │◄─┼────────►│ (periodic cron) │
415- │ └─────────────────┘ │ └──────────────────┘415+ │ └─────────────────┘ │ └──────────────────┘
416- └──────────────────────┘416+ └──────────────────────┘
417-417+
418- AppView responsibilities:418+ AppView responsibilities:
419- • Subscribe to firehose for at.glean.subscription, at.glean.annotation, at.glean.like419+ • Subscribe to Jetstream for at.glean.subscription, at.glean.annotation, at.glean.like
420- • Index records into SQLite420+ • Index records into SQLite
421- • Serve XRPC query endpoints (at.glean.listSubscriptions, etc.)421+ • Serve XRPC query endpoints (at.glean.listSubscriptions, etc.)
422- • Host the web UI at glean.at422+ • Host the web UI at glean.at
423- • Write to user PDS on behalf of user (when user acts through UI)423+ • Write to user PDS on behalf of user (when user acts through UI)
424 ```424 ```
425 425
426 ## 6. Database Schema (SQLite)426 ## 6. Database Schema (SQLite)
@@ -525,6 +525,7 @@ CREATE INDEX idx_read_state_starred ON read_state(user_did, is_starred) WHERE is
525 ```525 ```
526 526
527 ### 6.6 Annotations, Likes527 ### 6.6 Annotations, Likes
528+
528 Local mirror of AT Protocol lexicon records for fast querying.529 Local mirror of AT Protocol lexicon records for fast querying.
529 530
530 ```sql531 ```sql
@@ -654,9 +655,9 @@ For larger scale, move to MinHash + LSH (banded hashing) to approximate Jaccard
654 655
655 ### 7.5 Clustering Engine (Cron)656 ### 7.5 Clustering Engine (Cron)
656 657
657-A background goroutine runs on a schedule (e.g., every 6 hours):658+A background goroutine runs on a configurable schedule (`GLEAN_CLUSTER_INTERVAL`, default 6h):
658 659
659-1. **Firehose ingestion**: Subscribe to AT Protocol firehose for `at.glean.*` records660+1. **Jetstream ingestion**: Subscribe to Jetstream for `at.glean.*` records
660 2. **Index new records**: Parse lexicon records, upsert into SQLite661 2. **Index new records**: Parse lexicon records, upsert into SQLite
661 3. **Compute similarities**: Batch-update the `feed_similarity`, `user_similarity`, and `article_co_like` tables662 3. **Compute similarities**: Batch-update the `feed_similarity`, `user_similarity`, and `article_co_like` tables
662 4. **Generate recommendations**: Materialize top recommendations per user into cache tables663 4. **Generate recommendations**: Materialize top recommendations per user into cache tables
@@ -686,28 +687,28 @@ The server renders HTML fragments that htmx swaps into the page. No JSON API nee
686 687
687 ### 8.1 Pages688 ### 8.1 Pages
688 689
689-| Route | Method | Description |690+| Route | Method | Description |
690-| ---------------------- | ------ | -------------------------------------------------------- |691+| ------------------------- | ------ | -------------------------------------------------------- |
691-| `/` | GET | Landing page / auth redirect |692+| `/` | GET | Landing page / auth redirect |
692-| `/dashboard` | GET | Main dashboard: unread articles, recommendations sidebar |693+| `/dashboard` | GET | Main dashboard: unread articles, recommendations sidebar |
693-| `/feeds` | GET | Manage RSS subscriptions (OPML import for onboarding) |694+| `/feeds` | GET | Manage RSS subscriptions (OPML import for onboarding) |
694-| `/feeds/opml/upload` | POST | Upload OPML file to bulk-import subscriptions |695+| `/feeds/opml/upload` | POST | Upload OPML file to bulk-import subscriptions |
695-| `/feeds/opml/download` | GET | Export subscriptions as OPML (offboarding) |696+| `/feeds/opml/download` | GET | Export subscriptions as OPML (offboarding) |
696-| `/feeds/add` | POST | Add a single feed URL |697+| `/feeds/add` | POST | Add a single feed URL |
697-| `/feeds/remove` | DELETE | Remove a feed |698+| `/feeds/remove` | DELETE | Remove a feed |
698-| `/feeds/refresh` | POST | Refresh all subscribed feeds |699+| `/feeds/refresh` | POST | Refresh all subscribed feeds |
699-| `/feeds/clear` | POST | Clear all subscriptions |700+| `/feeds/clear` | POST | Clear all subscriptions |
700-| `/articles` | GET | Read articles (paginated, filterable by feed) |701+| `/articles` | GET | Read articles (paginated, filterable by feed) |
701-| `/articles/{id}` | GET | Article detail view |702+| `/articles/{id}` | GET | Article detail view |
702-| `/articles/{id}/read` | POST | Mark article as read |703+| `/articles/{id}/read` | POST | Mark article as read |
703-| `/articles/{id}/unread`| POST | Mark article as unread |704+| `/articles/{id}/unread` | POST | Mark article as unread |
704-| `/articles/{id}/like` | POST | Like an article |705+| `/articles/{id}/like` | POST | Like an article |
705-| `/articles/mark-all-read` | POST | Mark all articles as read |706+| `/articles/mark-all-read` | POST | Mark all articles as read |
706-| `/trending` | GET | Community feed: articles ranked by likes |707+| `/trending` | GET | Community feed: articles ranked by likes |
707-| `/library` | GET | Liked articles and annotations |708+| `/library` | GET | Liked articles and annotations |
708-| `/library/create` | POST | Create annotation on an article |709+| `/library/create` | POST | Create annotation on an article |
709-| `/library/{id}/delete` | POST | Delete an annotation |710+| `/library/{id}/delete` | POST | Delete an annotation |
710-| `/profile/{did}` | GET | Public profile: their feeds, likes, annotations |711+| `/profile/{did}` | GET | Public profile: their feeds, likes, annotations |
711 712
712 ### 8.2 htmx Patterns713 ### 8.2 htmx Patterns
713 714
@@ -729,7 +730,8 @@ glean/
729 │ ├── atproto/730 │ ├── atproto/
730 │ │ ├── auth.go # DID resolution, OAuth flow731 │ │ ├── auth.go # DID resolution, OAuth flow
731 │ │ ├── client.go # XRPC client (write to user PDS)732 │ │ ├── client.go # XRPC client (write to user PDS)
732-│ │ ├── firehose.go # Subscribe to AT Relay firehose733+│ │ ├── jetstream.go # Subscribe to Jetstream via official client
734+│ │ ├── stream_handler.go # Stream event → DB handler
733 │ │ ├── sync.go # PDS record reconciliation735 │ │ ├── sync.go # PDS record reconciliation
734 │ │ └── xrpc.go # XRPC query handlers (AppView endpoints)736 │ │ └── xrpc.go # XRPC query handlers (AppView endpoints)
735 │ ├── db/737 │ ├── db/
@@ -862,7 +864,7 @@ Glean exposes a `/metrics` endpoint for monitoring. Key metrics:
862 - **`glean_feeds_fetched_total`** — Feed fetch attempts labeled by status (`success`, `error`, `not_modified`)864 - **`glean_feeds_fetched_total`** — Feed fetch attempts labeled by status (`success`, `error`, `not_modified`)
863 - **`glean_feed_fetch_duration_seconds`** — Histogram of feed fetch latency865 - **`glean_feed_fetch_duration_seconds`** — Histogram of feed fetch latency
864 - **`glean_articles_upserted_total`** — Counter of articles stored from feeds866 - **`glean_articles_upserted_total`** — Counter of articles stored from feeds
865-- **`glean_firehose_events_total`** — Firehose events labeled by collection and action867+- **`glean_jetstream_events_total`** — Jetstream events labeled by collection and action
866 - **`glean_http_requests_total`** — HTTP request counts labeled by method, path, and status868 - **`glean_http_requests_total`** — HTTP request counts labeled by method, path, and status
867 - **`glean_cluster_runs_total`** / **`glean_cluster_duration_seconds`** — Recommendation engine runs and timing869 - **`glean_cluster_runs_total`** / **`glean_cluster_duration_seconds`** — Recommendation engine runs and timing
868 870
@@ -878,8 +880,8 @@ Glean exposes a `/metrics` endpoint for monitoring. Key metrics:
878 880
879 Glean operates as an AT Protocol AppView. This means:881 Glean operates as an AT Protocol AppView. This means:
880 882
881-- **Read path**: All `at.glean.*` data is consumed from the Relay firehose, not by polling individual PDS instances. The firehose handler runs as a persistent goroutine, upserting records into SQLite as they arrive.883+- **Read path**: All `at.glean.*` data is consumed from Jetstream, not by polling individual PDS instances. The Jetstream consumer runs as a persistent goroutine, upserting records into SQLite as they arrive.
882-- **Write path**: Users write records to their own PDS (via standard AT Protocol `com.atproto.repo.createRecord` / `deleteRecord`). Glean never stores user data directly — it only indexes what the firehose delivers.884+- **Write path**: Users write records to their own PDS (via standard AT Protocol `com.atproto.repo.createRecord` / `deleteRecord`). Glean never stores user data directly — it only indexes what Jetstream delivers.
883 - **Query path**: Other AT Protocol apps can query Glean's XRPC endpoints to access indexed data (subscriptions, annotations, likes, recommendations) without building their own indexer.885 - **Query path**: Other AT Protocol apps can query Glean's XRPC endpoints to access indexed data (subscriptions, annotations, likes, recommendations) without building their own indexer.
884 - **Trade-off**: Article content (fetched from RSS feeds) is local-only and not part of the AT Protocol layer. Only individual feed subscription records (`at.glean.subscription`) live on the PDS.886 - **Trade-off**: Article content (fetched from RSS feeds) is local-only and not part of the AT Protocol layer. Only individual feed subscription records (`at.glean.subscription`) live on the PDS.
885 887
@@ -891,6 +893,4 @@ All PDS records are public. There is no notion of private data on the AT Protoco
891 893
892 - **MinHash/LSH**: Replace brute-force Jaccard when user count exceeds ~50k894 - **MinHash/LSH**: Replace brute-force Jaccard when user count exceeds ~50k
893 - **Full-text search**: Add FTS5 virtual table on articles for search895 - **Full-text search**: Add FTS5 virtual table on articles for search
894-- **Feed groups / reading lists**: Allow users to create curated lists (separate lexicon)
895 - **Email digest**: Periodic email with top articles from subscribed feeds896 - **Email digest**: Periodic email with top articles from subscribed feeds
896-- **Multi-AppView scaling**: Distribute firehose consumption across multiple instances behind a load balancer
modified main.go +13 -2
@@ -22,6 +22,8 @@ func main() {
2222 addr := flag.String("addr", envOr("GLEAN_ADDR", ":8080"), "listen address")
2323 dbPath := flag.String("db", envOr("GLEAN_DB", "glean.db"), "database path")
2424 jetstreamURL := flag.String("jetstream", envOr("GLEAN_JETSTREAM", "wss://jetstream2.fr.hose.cam"), "Jetstream URL")
25+ syncInterval := flag.Duration("sync-interval", envDuration("GLEAN_SYNC_INTERVAL", 1*time.Hour), "PDS sync interval")
26+ clusterInterval := flag.Duration("cluster-interval", envDuration("GLEAN_CLUSTER_INTERVAL", 6*time.Hour), "cluster recomputation interval")
2527 flag.Parse()
2628
2729 logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
@@ -42,7 +44,7 @@ func main() {
4244 srv := server.New(database, clientID, callbackURL, *addr, scheduler, logger)
4345
4446 engine := cluster.NewEngine(database.DB, logger)
45- cron := cluster.NewCron(engine, 6*time.Hour, logger)
47+ cron := cluster.NewCron(engine, *clusterInterval, logger)
4648
4749 handler := atproto.NewStreamDBHandler(database, logger)
4850 jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger)
@@ -61,7 +63,7 @@ func main() {
6163 }
6264 }()
6365 go func() {
64- srv.PeriodicSync(ctx, 1*time.Hour)
66+ srv.PeriodicSync(ctx, *syncInterval)
6567 }()
6668 go func() {
6769 if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil {
@@ -107,3 +109,12 @@ func envOr(key, fallback string) string {
107109 }
108110 return fallback
109111 }
112+
113+func envDuration(key string, fallback time.Duration) time.Duration {
114+ if v := os.Getenv(key); v != "" {
115+ if d, err := time.ParseDuration(v); err == nil {
116+ return d
117+ }
118+ }
119+ return fallback
120+}
@@ -22,6 +22,8 @@ func main() {
22 addr := flag.String("addr", envOr("GLEAN_ADDR", ":8080"), "listen address")22 addr := flag.String("addr", envOr("GLEAN_ADDR", ":8080"), "listen address")
23 dbPath := flag.String("db", envOr("GLEAN_DB", "glean.db"), "database path")23 dbPath := flag.String("db", envOr("GLEAN_DB", "glean.db"), "database path")
24 jetstreamURL := flag.String("jetstream", envOr("GLEAN_JETSTREAM", "wss://jetstream2.fr.hose.cam"), "Jetstream URL")24 jetstreamURL := flag.String("jetstream", envOr("GLEAN_JETSTREAM", "wss://jetstream2.fr.hose.cam"), "Jetstream URL")
25+ syncInterval := flag.Duration("sync-interval", envDuration("GLEAN_SYNC_INTERVAL", 1*time.Hour), "PDS sync interval")
26+ clusterInterval := flag.Duration("cluster-interval", envDuration("GLEAN_CLUSTER_INTERVAL", 6*time.Hour), "cluster recomputation interval")
25 flag.Parse()27 flag.Parse()
26 28
27 logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))29 logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: slog.LevelInfo}))
@@ -42,7 +44,7 @@ func main() {
42 srv := server.New(database, clientID, callbackURL, *addr, scheduler, logger)44 srv := server.New(database, clientID, callbackURL, *addr, scheduler, logger)
43 45
44 engine := cluster.NewEngine(database.DB, logger)46 engine := cluster.NewEngine(database.DB, logger)
45- cron := cluster.NewCron(engine, 6*time.Hour, logger)47+ cron := cluster.NewCron(engine, *clusterInterval, logger)
46 48
47 handler := atproto.NewStreamDBHandler(database, logger)49 handler := atproto.NewStreamDBHandler(database, logger)
48 jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger)50 jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger)
@@ -61,7 +63,7 @@ func main() {
61 }63 }
62 }()64 }()
63 go func() {65 go func() {
64- srv.PeriodicSync(ctx, 1*time.Hour)66+ srv.PeriodicSync(ctx, *syncInterval)
65 }()67 }()
66 go func() {68 go func() {
67 if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil {69 if err := jetstream.Start(ctx); err != nil && ctx.Err() == nil {
@@ -107,3 +109,12 @@ func envOr(key, fallback string) string {
107 }109 }
108 return fallback110 return fallback
109 }111 }
112+
113+func envDuration(key string, fallback time.Duration) time.Duration {
114+ if v := os.Getenv(key); v != "" {
115+ if d, err := time.ParseDuration(v); err == nil {
116+ return d
117+ }
118+ }
119+ return fallback
120+}
modified readme.md +2 -0
@@ -41,6 +41,8 @@ Then open `http://localhost:8080`.
4141 | `GLEAN_ADDR` | `:8080` | Listen address |
4242 | `GLEAN_DB` | `glean.db` | SQLite database path |
4343 | `GLEAN_JETSTREAM` | `wss://jetstream2.fr.hose.cam` | Jetstream WebSocket URL |
44+| `GLEAN_SYNC_INTERVAL` | `1h` | PDS sync interval (Go duration: `30m`, `2h30m`, etc.) |
45+| `GLEAN_CLUSTER_INTERVAL` | `6h` | Cluster recomputation interval (Go duration) |
4446 | `GLEAN_OAUTH_CLIENT_ID` | _(empty)_ | OAuth client metadata URL (leave empty for localhost dev) |
4547 | `GLEAN_OAUTH_REDIRECT_URL` | _(empty)_ | OAuth redirect URL (leave empty for localhost dev) |
4648
@@ -41,6 +41,8 @@ Then open `http://localhost:8080`.
41 | `GLEAN_ADDR` | `:8080` | Listen address |41 | `GLEAN_ADDR` | `:8080` | Listen address |
42 | `GLEAN_DB` | `glean.db` | SQLite database path |42 | `GLEAN_DB` | `glean.db` | SQLite database path |
43 | `GLEAN_JETSTREAM` | `wss://jetstream2.fr.hose.cam` | Jetstream WebSocket URL |43 | `GLEAN_JETSTREAM` | `wss://jetstream2.fr.hose.cam` | Jetstream WebSocket URL |
44+| `GLEAN_SYNC_INTERVAL` | `1h` | PDS sync interval (Go duration: `30m`, `2h30m`, etc.) |
45+| `GLEAN_CLUSTER_INTERVAL` | `6h` | Cluster recomputation interval (Go duration) |
44 | `GLEAN_OAUTH_CLIENT_ID` | _(empty)_ | OAuth client metadata URL (leave empty for localhost dev) |46 | `GLEAN_OAUTH_CLIENT_ID` | _(empty)_ | OAuth client metadata URL (leave empty for localhost dev) |
45 | `GLEAN_OAUTH_REDIRECT_URL` | _(empty)_ | OAuth redirect URL (leave empty for localhost dev) |47 | `GLEAN_OAUTH_REDIRECT_URL` | _(empty)_ | OAuth redirect URL (leave empty for localhost dev) |
46 48