diff --git a/BACKLOG.md b/BACKLOG.md index 0c7bc51..100391e 100644 --- a/BACKLOG.md +++ b/BACKLOG.md @@ -26,9 +26,6 @@ Each should be addressed one at a time, and the item should be removed after imp ## Fixes -- Adding new gear (grinders, etc.) from profile page redirects to brews on profile page after (should stay on current page) -- Brews on profile page are stored in chronological order, should be reverse chronological (newest first) - - [Future work]: adjust timing of caching in feed, maybe use firehose and a sqlite database since we are only storing a few anyway - Goal: reduce pings to server when idling diff --git a/cmd/server/main.go b/cmd/server/main.go index 7545622..33ca9c9 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -1,15 +1,20 @@ package main import ( + "context" + "flag" "fmt" "net/http" "os" + "os/signal" "path/filepath" + "syscall" "time" "arabica/internal/atproto" "arabica/internal/database/boltstore" "arabica/internal/feed" + "arabica/internal/firehose" "arabica/internal/handlers" "arabica/internal/routing" @@ -18,6 +23,10 @@ import ( ) func main() { + // Parse command-line flags + useFirehose := flag.Bool("firehose", false, "Enable firehose-based feed (Jetstream consumer)") + flag.Parse() + // Configure zerolog // Set log level from environment (default: info) logLevel := os.Getenv("LOG_LEVEL") @@ -46,7 +55,7 @@ func main() { }) } - log.Info().Msg("Starting Arabica Coffee Tracker") + log.Info().Bool("firehose", *useFirehose).Msg("Starting Arabica Coffee Tracker") // Get port from env or use default port := os.Getenv("PORT") @@ -123,10 +132,84 @@ func main() { Int("registered_users", feedRegistry.Count()). Msg("Feed service initialized with persistent registry") + // Setup context for graceful shutdown + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + // Handle shutdown signals + sigCh := make(chan os.Signal, 1) + signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) + + // Initialize firehose consumer if enabled + var firehoseConsumer *firehose.Consumer + if *useFirehose { + // Determine feed index path + feedIndexPath := os.Getenv("ARABICA_FEED_INDEX_PATH") + if feedIndexPath == "" { + dataDir := os.Getenv("XDG_DATA_HOME") + if dataDir == "" { + home, err := os.UserHomeDir() + if err != nil { + log.Fatal().Err(err).Msg("Failed to get home directory for feed index") + } + dataDir = filepath.Join(home, ".local", "share") + } + feedIndexPath = filepath.Join(dataDir, "arabica", "feed-index.db") + } + + // Create firehose config + firehoseConfig := firehose.DefaultConfig() + firehoseConfig.IndexPath = feedIndexPath + + // Parse profile cache TTL from env if set + if ttlStr := os.Getenv("ARABICA_PROFILE_CACHE_TTL"); ttlStr != "" { + if ttl, err := time.ParseDuration(ttlStr); err == nil { + firehoseConfig.ProfileCacheTTL = int64(ttl.Seconds()) + } + } + + // Create feed index + feedIndex, err := firehose.NewFeedIndex(feedIndexPath, time.Duration(firehoseConfig.ProfileCacheTTL)*time.Second) + if err != nil { + log.Fatal().Err(err).Str("path", feedIndexPath).Msg("Failed to create feed index") + } + + log.Info().Str("path", feedIndexPath).Msg("Feed index opened") + + // Create and start consumer + firehoseConsumer = firehose.NewConsumer(firehoseConfig, feedIndex) + firehoseConsumer.Start(ctx) + + // Wire up the feed service to use the firehose index + adapter := firehose.NewFeedIndexAdapter(feedIndex) + feedService.SetFirehoseIndex(adapter) + + log.Info().Msg("Firehose consumer started") + + // Backfill registered users in background + go func() { + time.Sleep(5 * time.Second) // Wait for initial connection + for _, did := range feedRegistry.List() { + if err := firehoseConsumer.BackfillDID(ctx, did); err != nil { + log.Warn().Err(err).Str("did", did).Msg("Failed to backfill user") + } + } + log.Info().Int("count", feedRegistry.Count()).Msg("Backfill of registered users complete") + }() + } + // Register users in the feed when they authenticate // This ensures users are added to the feed even if they had an existing session oauthManager.SetOnAuthSuccess(func(did string) { feedRegistry.Register(did) + // If firehose is enabled, backfill the user's records + if firehoseConsumer != nil { + go func() { + if err := firehoseConsumer.BackfillDID(context.Background(), did); err != nil { + log.Warn().Err(err).Str("did", did).Msg("Failed to backfill new user") + } + }() + } }) if clientID == "" { @@ -175,15 +258,44 @@ func main() { Logger: log.Logger, }) - // Start HTTP server - log.Info(). - Str("address", "0.0.0.0:"+port). - Str("url", "http://localhost:"+port). - Bool("secure_cookies", secureCookies). - Str("database", dbPath). - Msg("Starting HTTP server") - - if err := http.ListenAndServe("0.0.0.0:"+port, handler); err != nil { - log.Fatal().Err(err).Msg("Server failed to start") + // Create HTTP server + server := &http.Server{ + Addr: "0.0.0.0:" + port, + Handler: handler, + } + + // Start HTTP server in goroutine + go func() { + log.Info(). + Str("address", "0.0.0.0:"+port). + Str("url", "http://localhost:"+port). + Bool("secure_cookies", secureCookies). + Bool("firehose", *useFirehose). + Str("database", dbPath). + Msg("Starting HTTP server") + + if err := server.ListenAndServe(); err != nil && err != http.ErrServerClosed { + log.Fatal().Err(err).Msg("Server failed to start") + } + }() + + // Wait for shutdown signal + <-sigCh + log.Info().Msg("Shutdown signal received") + + // Stop firehose consumer first + if firehoseConsumer != nil { + log.Info().Msg("Stopping firehose consumer...") + firehoseConsumer.Stop() } + + // Graceful shutdown of HTTP server + shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 10*time.Second) + defer shutdownCancel() + + if err := server.Shutdown(shutdownCtx); err != nil { + log.Error().Err(err).Msg("HTTP server shutdown error") + } + + log.Info().Msg("Server stopped") } diff --git a/docs/firehose-plan.md b/docs/firehose-plan.md new file mode 100644 index 0000000..1a93ded --- /dev/null +++ b/docs/firehose-plan.md @@ -0,0 +1,641 @@ +# Firehose Integration Plan for Arabica + +## Executive Summary + +This document proposes refactoring Arabica's home page feed to consume events from the AT Protocol firehose via Jetstream, replacing the current polling-based approach. This will provide real-time updates, dramatically reduce API calls, and improve scalability. + +**Recommendation:** Implement Jetstream consumer with local BoltDB index as Phase 1, with optional Slingshot/Constellation integration in Phase 2. + +--- + +## Problem Statement + +### Current Architecture + +The feed service (`internal/feed/service.go`) polls each registered user's PDS directly: + +``` +For N registered users: +- N profile fetches +- N × 5 collection fetches (brew, bean, roaster, grinder, brewer) +- N × 4 reference resolution fetches +- Total: ~10N API calls per refresh +``` + +### Issues + +| Problem | Impact | +| ------------------------ | ----------------------------------- | +| High API call volume | Risk of rate limiting as users grow | +| 5-minute cache staleness | Users don't see recent activity | +| N+1 query pattern | Linear scaling, O(N) per refresh | +| PDS dependency | Feed fails if any PDS is slow/down | +| No real-time updates | Requires manual refresh | + +--- + +## Proposed Solution: Jetstream Consumer + +### Architecture Overview + +``` +┌─────────────────┐ ┌──────────────────┐ ┌─────────────────┐ +│ AT Protocol │ │ Jetstream │ │ Arabica │ +│ Firehose │────▶│ (Public/Self) │────▶│ Consumer │ +│ (all records) │ │ JSON over WS │ │ (background) │ +└─────────────────┘ └──────────────────┘ └────────┬────────┘ + │ + ▼ + ┌─────────────────┐ + │ Feed Index │ + │ (BoltDB) │ + └────────┬────────┘ + │ + ▼ + ┌─────────────────┐ + │ HTTP Handler │ + │ (instant) │ + └─────────────────┘ +``` + +### How It Works + +1. **Background Consumer** connects to Jetstream WebSocket +2. **Filters** for `social.arabica.alpha.*` collections +3. **Indexes** incoming events into local BoltDB +4. **Serves** feed requests instantly from local index +5. **Fallback** to direct polling if consumer disconnects + +### Benefits + +| Metric | Current | With Jetstream | +| --------------------- | ---------------- | ----------------- | +| API calls per refresh | ~10N | 0 | +| Feed latency | 5 min cache | Real-time (<1s) | +| PDS dependency | High | None (after sync) | +| User discovery | Manual registry | Automatic | +| Scalability | O(N) per request | O(1) per request | + +--- + +## Technical Design + +### 1. Jetstream Client Configuration + +```go +// internal/firehose/config.go + +type JetstreamConfig struct { + // Public endpoints (fallback rotation) + Endpoints []string + + // Filter to Arabica collections only + WantedCollections []string + + // Optional: filter to registered DIDs only + // Leave empty to discover all Arabica users + WantedDids []string + + // Enable zstd compression (~56% bandwidth reduction) + Compress bool + + // Cursor file path for restart recovery + CursorFile string +} + +func DefaultConfig() *JetstreamConfig { + return &JetstreamConfig{ + Endpoints: []string{ + "wss://jetstream1.us-east.bsky.network/subscribe", + "wss://jetstream2.us-east.bsky.network/subscribe", + "wss://jetstream1.us-west.bsky.network/subscribe", + "wss://jetstream2.us-west.bsky.network/subscribe", + }, + WantedCollections: []string{ + "social.arabica.alpha.brew", + "social.arabica.alpha.bean", + "social.arabica.alpha.roaster", + "social.arabica.alpha.grinder", + "social.arabica.alpha.brewer", + }, + Compress: true, + CursorFile: "jetstream-cursor.txt", + } +} +``` + +### 2. Event Processing + +```go +// internal/firehose/consumer.go + +type Consumer struct { + config *JetstreamConfig + index *FeedIndex + client *jetstream.Client + cursor atomic.Int64 + connected atomic.Bool +} + +func (c *Consumer) handleEvent(ctx context.Context, event *models.Event) error { + if event.Kind != "commit" || event.Commit == nil { + return nil + } + + commit := event.Commit + + // Only process Arabica collections + if !strings.HasPrefix(commit.Collection, "social.arabica.alpha.") { + return nil + } + + switch commit.Operation { + case "create", "update": + return c.index.UpsertRecord(ctx, event.Did, commit) + case "delete": + return c.index.DeleteRecord(ctx, event.Did, commit.Collection, commit.RKey) + } + + // Update cursor for recovery + c.cursor.Store(event.TimeUS) + + return nil +} +``` + +### 3. Feed Index Schema (BoltDB) + +```go +// internal/firehose/index.go + +// BoltDB Buckets: +// - "records" : {at-uri} -> {record JSON + metadata} +// - "by_time" : {timestamp:at-uri} -> {} (for chronological queries) +// - "by_did" : {did:at-uri} -> {} (for user-specific queries) +// - "by_type" : {collection:timestamp:at-uri} -> {} (for type filtering) +// - "profiles" : {did} -> {profile JSON} (cached profiles) +// - "cursor" : "jetstream" -> {cursor value} + +type FeedIndex struct { + db *bbolt.DB +} + +type IndexedRecord struct { + URI string `json:"uri"` + DID string `json:"did"` + Collection string `json:"collection"` + RKey string `json:"rkey"` + Record json.RawMessage `json:"record"` + CID string `json:"cid"` + IndexedAt time.Time `json:"indexed_at"` +} + +func (idx *FeedIndex) GetRecentFeed(ctx context.Context, limit int) ([]*FeedItem, error) { + // Query by_time bucket in reverse order + // Hydrate with profile data from profiles bucket + // Return feed items instantly from local data +} +``` + +### 4. Profile Resolution + +Profiles are not part of Arabica's lexicons, so we need a strategy: + +**Option A: Lazy Loading (Recommended for Phase 1)** + +```go +func (idx *FeedIndex) resolveProfile(ctx context.Context, did string) (*Profile, error) { + // Check local cache first + if profile := idx.getCachedProfile(did); profile != nil { + return profile, nil + } + + // Fetch from public API and cache + profile, err := publicClient.GetProfile(ctx, did) + if err != nil { + return nil, err + } + + idx.cacheProfile(did, profile, 1*time.Hour) + return profile, nil +} +``` + +**Option B: Slingshot Integration (Phase 2)** + +```go +// Use Slingshot's resolveMiniDoc for faster profile resolution +func (idx *FeedIndex) resolveProfileViaSlingshot(ctx context.Context, did string) (*Profile, error) { + url := fmt.Sprintf("https://slingshot.microcosm.blue/xrpc/com.bad-example.identity.resolveMiniDoc?identifier=%s", did) + // Returns {did, handle, pds} in one call +} +``` + +### 5. Reference Resolution + +Brews reference beans, grinders, and brewers. The index already has these records: + +```go +func (idx *FeedIndex) resolveBrew(ctx context.Context, brew *IndexedRecord) (*FeedItem, error) { + var record map[string]interface{} + json.Unmarshal(brew.Record, &record) + + item := &FeedItem{RecordType: "brew"} + + // Resolve bean reference from local index + if beanRef, ok := record["beanRef"].(string); ok { + if bean := idx.getRecord(beanRef); bean != nil { + item.Bean = recordToBean(bean) + } + } + + // Similar for grinder, brewer references + // All from local index - no API calls + + return item, nil +} +``` + +### 6. Fallback and Resilience + +```go +// internal/firehose/consumer.go + +func (c *Consumer) Run(ctx context.Context) error { + for { + select { + case <-ctx.Done(): + return ctx.Err() + default: + if err := c.connectAndConsume(ctx); err != nil { + log.Warn().Err(err).Msg("jetstream connection lost, reconnecting...") + + // Exponential backoff + time.Sleep(c.backoff.NextBackOff()) + + // Rotate to next endpoint + c.rotateEndpoint() + continue + } + } + } +} + +func (c *Consumer) connectAndConsume(ctx context.Context) error { + cursor := c.loadCursor() + + // Rewind cursor slightly to handle duplicates safely + if cursor > 0 { + cursor -= 5 * time.Second.Microseconds() + } + + return c.client.ConnectAndRead(ctx, &cursor) +} +``` + +### 7. Feed Service Integration + +```go +// internal/feed/service.go (modified) + +type Service struct { + registry *Registry + publicClient *atproto.PublicClient + cache *publicFeedCache + + // New: firehose index + firehoseIndex *firehose.FeedIndex + useFirehose bool +} + +func (s *Service) GetRecentRecords(ctx context.Context, limit int) ([]*FeedItem, error) { + // Prefer firehose index if available and populated + if s.useFirehose && s.firehoseIndex.IsReady() { + return s.firehoseIndex.GetRecentFeed(ctx, limit) + } + + // Fallback to polling (existing code) + return s.getRecentRecordsViaPolling(ctx, limit) +} +``` + +--- + +## Implementation Phases + +### Phase 1: Core Jetstream Consumer (2 weeks) + +**Goal:** Replace polling with firehose consumption for the feed. + +**Tasks:** + +1. Create `internal/firehose/` package + - `config.go` - Jetstream configuration + - `consumer.go` - WebSocket consumer with reconnection + - `index.go` - BoltDB-backed feed index + - `scheduler.go` - Event processing scheduler + +2. Integrate with existing feed service + - Add feature flag: `ARABICA_USE_FIREHOSE=true` (just use a cli flag) + - Keep polling as fallback + +3. Handle profile resolution + - Cache profiles locally with 1-hour TTL + - Lazy fetch on first access + - Background refresh for active users + +4. Cursor management + - Persist cursor to survive restarts + - Rewind on reconnection for safety + +**Deliverables:** + +- Real-time feed updates +- Reduced API calls to near-zero +- Automatic user discovery (anyone using Arabica lexicons) + +### Phase 2: Slingshot Optimization (1 week) + +**Goal:** Faster profile and record hydration. + +**Tasks:** + +1. Add Slingshot client (`internal/atproto/slingshot.go`) +2. Use `resolveMiniDoc` for profile resolution +3. Use Slingshot as fallback for missing records + +**Deliverables:** + +- Faster profile loading +- Resilience to slow PDS endpoints + +### Phase 3: Constellation for Social (1 week) + +**Goal:** Enable like/comment counts when social features are added. + +**Tasks:** + +1. Add Constellation client (`internal/atproto/constellation.go`) +2. Query backlinks for interaction counts +3. Display counts on feed items + +**Deliverables:** + +- Like count on brews +- Comment count on brews +- Foundation for social features + +### Phase 4: Spacedust for Real-time Notifications (Future) + +**Goal:** Push notifications for interactions. + +**Tasks:** + +1. Subscribe to Spacedust for user's content interactions +2. Build notification storage and API +3. WebSocket to frontend for live updates + +--- + +## Data Flow Comparison + +### Before (Polling) + +``` +User Request → Check Cache → [Cache Miss] → Poll N PDSes → Build Feed → Return + ↓ + ~10N API calls + 5-10 second latency +``` + +### After (Jetstream) + +``` +Jetstream → Consumer → Index (BoltDB) + ↓ +User Request → Query Index → Return + ↓ + 0 API calls + <10ms latency +``` + +--- + +## Automatic User Discovery + +A major benefit of firehose consumption is automatic user discovery: + +**Current:** Users must explicitly register via `/api/feed/register` + +**With Jetstream:** Any user who creates an Arabica record is automatically indexed + +```go +// When we see a new DID creating Arabica records +func (c *Consumer) handleNewUser(did string) { + // Auto-register for feed + c.registry.Register(did) + + // Fetch and cache their profile + go c.index.fetchAndCacheProfile(did) + + // Backfill their existing records + go c.backfillUser(did) +} +``` + +This could replace the manual registry entirely, or supplement it for "featured" users. + +--- + +## Backfill Strategy + +When starting fresh or discovering a new user, we need historical data: + +**Option A: Direct PDS Fetch (Simple)** + +```go +func (c *Consumer) backfillUser(ctx context.Context, did string) error { + for _, collection := range arabicaCollections { + records, _ := publicClient.ListRecords(ctx, did, collection, 100) + for _, record := range records { + c.index.UpsertFromPDS(record) + } + } + return nil +} +``` + +**Option B: Slingshot Fetch (Faster)** + +```go +func (c *Consumer) backfillUserViaSlingshot(ctx context.Context, did string) error { + // Single endpoint, pre-cached records + // Same API as PDS but faster +} +``` + +**Option C: Jetstream Cursor Rewind (Events Only)** + +- Rewind cursor to desired point in time +- Replay events (no records available before cursor) +- Limited to ~24h of history typically + +**Recommendation:** Use Option A for Phase 1, add Option B in Phase 2. + +--- + +## Configuration + +```bash +# Environment variables + +# Enable firehose-based feed (default: false during rollout) +ARABICA_USE_FIREHOSE=true + +# Jetstream endpoint (default: public Bluesky instances) +JETSTREAM_URL=wss://jetstream1.us-east.bsky.network/subscribe + +# Optional: self-hosted Jetstream +# JETSTREAM_URL=ws://localhost:6008/subscribe + +# Feed index database path +ARABICA_FEED_INDEX_PATH=~/.local/share/arabica/feed-index.db + +# Profile cache TTL (default: 1h) +ARABICA_PROFILE_CACHE_TTL=1h + +# Optional: Slingshot endpoint for Phase 2 +# SLINGSHOT_URL=https://slingshot.microcosm.blue + +# Optional: Constellation endpoint for Phase 3 +# CONSTELLATION_URL=https://constellation.microcosm.blue +``` + +--- + +## Monitoring and Metrics + +```go +// Prometheus metrics to track firehose health + +var ( + eventsReceived = prometheus.NewCounterVec( + prometheus.CounterOpts{ + Name: "arabica_firehose_events_total", + Help: "Total events received from Jetstream", + }, + []string{"collection", "operation"}, + ) + + indexSize = prometheus.NewGauge( + prometheus.GaugeOpts{ + Name: "arabica_feed_index_records", + Help: "Number of records in feed index", + }, + ) + + consumerLag = prometheus.NewGauge( + prometheus.GaugeOpts{ + Name: "arabica_firehose_lag_seconds", + Help: "Lag between event time and processing time", + }, + ) + + connectionState = prometheus.NewGauge( + prometheus.GaugeOpts{ + Name: "arabica_firehose_connected", + Help: "1 if connected to Jetstream, 0 otherwise", + }, + ) +) +``` + +--- + +## Risk Assessment + +| Risk | Mitigation | +| ----------------------- | --------------------------------------------- | +| Jetstream unavailable | Fallback to polling, rotate endpoints | +| Index corruption | Rebuild from backfill, periodic snapshots | +| Duplicate events | Idempotent upserts using AT-URI as key | +| Missing historical data | Backfill on startup and new user discovery | +| High event volume | Filter to Arabica collections only (~0 noise) | +| Profile resolution lag | Local cache with background refresh | + +--- + +## Open Questions + +1. **Should we remove the registry entirely?** + - Pro: Simpler, automatic discovery + - Con: Lose ability to curate "featured" users + - Recommendation: Keep registry for admin features, but don't require it for feed inclusion + +2. **Self-host Jetstream or use public?** + - Public is free and reliable + - Self-host gives control and removes dependency + - Recommendation: Start with public, evaluate self-hosting if issues arise + +3. **How long to keep historical data?** + - Option: Rolling 30-day window + - Option: Keep everything (disk is cheap) + - Recommendation: Keep 90 days, prune older records + +4. **Real-time feed updates to frontend?** + - Could push new items via WebSocket/SSE + - Or just reduce cache TTL to ~30 seconds + - Recommendation: Phase 1 just reduces staleness; real-time push is future enhancement + +--- + +## Alternatives Considered + +### 1. Tap (Bluesky's Full Sync Tool) + +**Pros:** Full verification, automatic backfill, collection signal mode +**Cons:** Heavy operational overhead, overkill for current scale +**Verdict:** Revisit when user base exceeds 500+ + +### 2. Direct Firehose Consumption + +**Pros:** No Jetstream dependency +**Cons:** Complex CBOR/CAR parsing, high bandwidth +**Verdict:** Jetstream provides the simplicity we need + +### 3. Slingshot as Primary Data Source + +**Pros:** Pre-cached records, single endpoint +**Cons:** Still polling-based, no real-time +**Verdict:** Use as optimization layer, not primary + +### 4. Spacedust Instead of Jetstream + +**Pros:** Link-focused, lightweight +**Cons:** Only links, no full records +**Verdict:** Use for notifications, not feed content + +--- + +## Success Criteria + +| Metric | Target | +| -------------------------- | ----------------------- | +| Feed latency | <100ms (from >5s) | +| API calls per feed request | 0 (from ~10N) | +| Time to see new content | <5s (from 5min) | +| Feed availability | 99.9% (with fallback) | +| New user discovery | Automatic (from manual) | + +--- + +## References + +- [Jetstream GitHub](https://github.com/bluesky-social/jetstream) +- [Jetstream Blog Post](https://docs.bsky.app/blog/jetstream) +- [Jetstream Go Client](https://pkg.go.dev/github.com/bluesky-social/jetstream/pkg/client) +- [Microcosm.blue Services](https://microcosm.blue/) +- [Constellation API](https://constellation.microcosm.blue/) +- [Slingshot API](https://slingshot.microcosm.blue/) +- [Existing Evaluation: Jetstream/Tap](./jetstream-tap-evaluation.md) +- [Existing Evaluation: Microcosm Tools](./microcosm-tools-evaluation.md) diff --git a/go.mod b/go.mod index 0a33626..4928114 100644 --- a/go.mod +++ b/go.mod @@ -18,7 +18,9 @@ require ( github.com/golang-jwt/jwt/v5 v5.2.2 // indirect github.com/google/go-cmp v0.6.0 // indirect github.com/google/go-querystring v1.1.0 // indirect + github.com/gorilla/websocket v1.5.3 // indirect github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect + github.com/klauspost/compress v1.18.3 // indirect github.com/mattn/go-colorable v0.1.13 // indirect github.com/mattn/go-isatty v0.0.20 // indirect github.com/matttproud/golang_protobuf_extensions/v2 v2.0.0 // indirect diff --git a/go.sum b/go.sum index c4bd6cf..e1fb056 100644 --- a/go.sum +++ b/go.sum @@ -17,10 +17,14 @@ github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-querystring v1.1.0 h1:AnCroh3fv4ZBgVIf1Iwtovgjaw/GiKJo8M8yD/fhyJ8= github.com/google/go-querystring v1.1.0/go.mod h1:Kcdr2DB4koayq7X8pmAG4sNG59So17icRSOU623lUBU= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/ipfs/go-cid v0.4.1 h1:A/T3qGvxi4kpKWWcPC/PgbvDA2bjVLO7n4UeVwnbs/s= github.com/ipfs/go-cid v0.4.1/go.mod h1:uQHwDeX4c6CtyrFwdqyhpNcxVewur1M7l7fNU7LKwZk= +github.com/klauspost/compress v1.18.3 h1:9PJRvfbmTabkOX8moIpXPbMMbYN60bWImDDU7L+/6zw= +github.com/klauspost/compress v1.18.3/go.mod h1:R0h/fSBs8DE4ENlcrlib3PsXS61voFxhIs2DeRhCvJ4= github.com/klauspost/cpuid/v2 v2.2.7 h1:ZWSB3igEs+d0qvnxR/ZBzXVmxkgt8DdzP6m9pfuVLDM= github.com/klauspost/cpuid/v2 v2.2.7/go.mod h1:Lcz8mBdAVJIBVzewtcLocK12l3Y+JytZYpaMropDUws= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= diff --git a/internal/feed/service.go b/internal/feed/service.go index f716612..64f5a99 100644 --- a/internal/feed/service.go +++ b/internal/feed/service.go @@ -45,11 +45,35 @@ type publicFeedCache struct { mu sync.RWMutex } +// FirehoseIndex is the interface for the firehose feed index +// This allows the feed service to use firehose data when available +type FirehoseIndex interface { + IsReady() bool + GetRecentFeed(ctx context.Context, limit int) ([]*FirehoseFeedItem, error) +} + +// FirehoseFeedItem matches the FeedItem structure from firehose package +// This avoids import cycles +type FirehoseFeedItem struct { + RecordType string + Action string + Brew *models.Brew + Bean *models.Bean + Roaster *models.Roaster + Grinder *models.Grinder + Brewer *models.Brewer + Author *atproto.Profile + Timestamp time.Time + TimeAgo string +} + // Service fetches and aggregates brews from registered users type Service struct { - registry *Registry - publicClient *atproto.PublicClient - cache *publicFeedCache + registry *Registry + publicClient *atproto.PublicClient + cache *publicFeedCache + firehoseIndex FirehoseIndex + useFirehose bool } // NewService creates a new feed service @@ -61,6 +85,13 @@ func NewService(registry *Registry) *Service { } } +// SetFirehoseIndex configures the service to use firehose-based feed when available +func (s *Service) SetFirehoseIndex(index FirehoseIndex) { + s.firehoseIndex = index + s.useFirehose = true + log.Info().Msg("feed: firehose index configured") +} + // GetCachedPublicFeed returns cached feed items for unauthenticated users. // It returns up to PublicFeedLimit items from the cache, refreshing if expired. func (s *Service) GetCachedPublicFeed(ctx context.Context) ([]*FeedItem, error) { @@ -115,13 +146,54 @@ func (s *Service) refreshPublicFeedCache(ctx context.Context) ([]*FeedItem, erro // GetRecentRecords fetches recent activity (brews and other records) from all registered users // Returns up to `limit` items sorted by most recent first func (s *Service) GetRecentRecords(ctx context.Context, limit int) ([]*FeedItem, error) { + // Try firehose index first if available and ready + if s.useFirehose && s.firehoseIndex != nil && s.firehoseIndex.IsReady() { + log.Debug().Msg("feed: using firehose index") + return s.getRecentRecordsFromFirehose(ctx, limit) + } + + // Fallback to polling + return s.getRecentRecordsViaPolling(ctx, limit) +} + +// getRecentRecordsFromFirehose fetches feed items from the firehose index +func (s *Service) getRecentRecordsFromFirehose(ctx context.Context, limit int) ([]*FeedItem, error) { + firehoseItems, err := s.firehoseIndex.GetRecentFeed(ctx, limit) + if err != nil { + log.Warn().Err(err).Msg("feed: firehose index error, falling back to polling") + return s.getRecentRecordsViaPolling(ctx, limit) + } + + // Convert FirehoseFeedItem to FeedItem + items := make([]*FeedItem, len(firehoseItems)) + for i, fi := range firehoseItems { + items[i] = &FeedItem{ + RecordType: fi.RecordType, + Action: fi.Action, + Brew: fi.Brew, + Bean: fi.Bean, + Roaster: fi.Roaster, + Grinder: fi.Grinder, + Brewer: fi.Brewer, + Author: fi.Author, + Timestamp: fi.Timestamp, + TimeAgo: fi.TimeAgo, + } + } + + log.Debug().Int("count", len(items)).Msg("feed: returning items from firehose index") + return items, nil +} + +// getRecentRecordsViaPolling fetches feed items by polling each user's PDS +func (s *Service) getRecentRecordsViaPolling(ctx context.Context, limit int) ([]*FeedItem, error) { dids := s.registry.List() if len(dids) == 0 { log.Debug().Msg("feed: no registered users") return nil, nil } - log.Debug().Int("user_count", len(dids)).Msg("feed: fetching activity from registered users") + log.Debug().Int("user_count", len(dids)).Msg("feed: fetching activity from registered users (polling)") // Fetch all records from all users in parallel type userActivity struct { diff --git a/internal/firehose/adapter.go b/internal/firehose/adapter.go new file mode 100644 index 0000000..20b76d3 --- /dev/null +++ b/internal/firehose/adapter.go @@ -0,0 +1,51 @@ +package firehose + +import ( + "context" + + "arabica/internal/feed" +) + +// FeedIndexAdapter wraps FeedIndex to implement feed.FirehoseIndex interface +// This avoids import cycles between feed and firehose packages +type FeedIndexAdapter struct { + index *FeedIndex +} + +// NewFeedIndexAdapter creates a new adapter for the FeedIndex +func NewFeedIndexAdapter(index *FeedIndex) *FeedIndexAdapter { + return &FeedIndexAdapter{index: index} +} + +// IsReady returns true if the index is ready to serve queries +func (a *FeedIndexAdapter) IsReady() bool { + return a.index.IsReady() +} + +// GetRecentFeed returns recent feed items from the index +// Converts FeedItem to feed.FirehoseFeedItem to satisfy the interface +func (a *FeedIndexAdapter) GetRecentFeed(ctx context.Context, limit int) ([]*feed.FirehoseFeedItem, error) { + items, err := a.index.GetRecentFeed(ctx, limit) + if err != nil { + return nil, err + } + + // Convert to the type expected by feed.Service + result := make([]*feed.FirehoseFeedItem, len(items)) + for i, item := range items { + result[i] = &feed.FirehoseFeedItem{ + RecordType: item.RecordType, + Action: item.Action, + Brew: item.Brew, + Bean: item.Bean, + Roaster: item.Roaster, + Grinder: item.Grinder, + Brewer: item.Brewer, + Author: item.Author, + Timestamp: item.Timestamp, + TimeAgo: item.TimeAgo, + } + } + + return result, nil +} diff --git a/internal/firehose/config.go b/internal/firehose/config.go new file mode 100644 index 0000000..7edc286 --- /dev/null +++ b/internal/firehose/config.go @@ -0,0 +1,53 @@ +// Package firehose provides real-time AT Protocol event consumption via Jetstream. +// It indexes Arabica records into a local BoltDB database for fast feed queries. +package firehose + +import ( + "arabica/internal/atproto" +) + +// Default Jetstream public endpoints +var DefaultJetstreamEndpoints = []string{ + "wss://jetstream1.us-east.bsky.network/subscribe", + "wss://jetstream2.us-east.bsky.network/subscribe", + "wss://jetstream1.us-west.bsky.network/subscribe", + "wss://jetstream2.us-west.bsky.network/subscribe", +} + +// ArabicaCollections lists all Arabica lexicon collections to filter for +var ArabicaCollections = []string{ + atproto.NSIDBrew, + atproto.NSIDBean, + atproto.NSIDRoaster, + atproto.NSIDGrinder, + atproto.NSIDBrewer, +} + +// Config holds configuration for the Jetstream consumer +type Config struct { + // Endpoints is a list of Jetstream WebSocket URLs to connect to (with fallback rotation) + Endpoints []string + + // WantedCollections filters events to specific collection NSIDs + WantedCollections []string + + // Compress enables zstd compression (~56% bandwidth reduction) + Compress bool + + // IndexPath is the path to the BoltDB feed index database + IndexPath string + + // ProfileCacheTTL is how long to cache profile data + ProfileCacheTTL int64 // seconds +} + +// DefaultConfig returns a configuration with sensible defaults +func DefaultConfig() *Config { + return &Config{ + Endpoints: DefaultJetstreamEndpoints, + WantedCollections: ArabicaCollections, + Compress: false, // Disabled: Jetstream uses custom zstd dictionary + IndexPath: "", // Will be set based on data directory + ProfileCacheTTL: 3600, // 1 hour + } +} diff --git a/internal/firehose/consumer.go b/internal/firehose/consumer.go new file mode 100644 index 0000000..1114962 --- /dev/null +++ b/internal/firehose/consumer.go @@ -0,0 +1,380 @@ +package firehose + +import ( + "context" + "encoding/json" + "fmt" + "net/url" + "strings" + "sync" + "sync/atomic" + "time" + + "github.com/gorilla/websocket" + "github.com/klauspost/compress/zstd" + "github.com/rs/zerolog/log" +) + +// JetstreamEvent represents an event from Jetstream +type JetstreamEvent struct { + DID string `json:"did"` + TimeUS int64 `json:"time_us"` + Kind string `json:"kind"` // "commit", "identity", "account" + Commit *struct { + Rev string `json:"rev"` + Operation string `json:"operation"` // "create", "update", "delete" + Collection string `json:"collection"` + RKey string `json:"rkey"` + Record json.RawMessage `json:"record,omitempty"` + CID string `json:"cid"` + } `json:"commit,omitempty"` +} + +// Consumer consumes events from Jetstream and indexes them +type Consumer struct { + config *Config + index *FeedIndex + + // Connection state + conn *websocket.Conn + connMu sync.Mutex + currentEndpointIdx int + + // Zstd decoder for compressed messages + zstdDecoder *zstd.Decoder + + // Cursor for resume + cursor atomic.Int64 + + // Stats + eventsReceived atomic.Int64 + bytesReceived atomic.Int64 + + // Control + connected atomic.Bool + stopCh chan struct{} + wg sync.WaitGroup +} + +// NewConsumer creates a new Jetstream consumer +func NewConsumer(config *Config, index *FeedIndex) *Consumer { + // Create zstd decoder for compressed messages + decoder, err := zstd.NewReader(nil, zstd.WithDecoderConcurrency(1)) + if err != nil { + log.Fatal().Err(err).Msg("firehose: failed to create zstd decoder") + } + + c := &Consumer{ + config: config, + index: index, + stopCh: make(chan struct{}), + zstdDecoder: decoder, + } + + // Load cursor from index + if cursor, err := index.GetCursor(); err == nil && cursor > 0 { + c.cursor.Store(cursor) + log.Info().Int64("cursor", cursor).Msg("firehose: loaded cursor from index") + } + + return c +} + +// Start begins consuming events in a background goroutine +func (c *Consumer) Start(ctx context.Context) { + c.wg.Add(1) + go func() { + defer c.wg.Done() + c.run(ctx) + }() +} + +// Stop gracefully stops the consumer +func (c *Consumer) Stop() { + close(c.stopCh) + c.connMu.Lock() + if c.conn != nil { + c.conn.Close() + } + c.connMu.Unlock() + c.wg.Wait() + + // Close zstd decoder + if c.zstdDecoder != nil { + c.zstdDecoder.Close() + } +} + +// IsConnected returns true if currently connected to Jetstream +func (c *Consumer) IsConnected() bool { + return c.connected.Load() +} + +// Stats returns consumer statistics +func (c *Consumer) Stats() (eventsReceived, bytesReceived int64) { + return c.eventsReceived.Load(), c.bytesReceived.Load() +} + +func (c *Consumer) run(ctx context.Context) { + backoff := time.Second + maxBackoff := 30 * time.Second + + for { + select { + case <-ctx.Done(): + log.Info().Msg("firehose: context cancelled, stopping consumer") + return + case <-c.stopCh: + log.Info().Msg("firehose: stop requested, stopping consumer") + return + default: + } + + endpoint := c.config.Endpoints[c.currentEndpointIdx] + err := c.connectAndConsume(ctx, endpoint) + + if err != nil { + c.connected.Store(false) + log.Warn().Err(err).Str("endpoint", endpoint).Msg("firehose: connection error") + + // Rotate to next endpoint + c.currentEndpointIdx = (c.currentEndpointIdx + 1) % len(c.config.Endpoints) + + // Backoff before retry + select { + case <-ctx.Done(): + return + case <-c.stopCh: + return + case <-time.After(backoff): + } + + // Increase backoff + backoff *= 2 + if backoff > maxBackoff { + backoff = maxBackoff + } + } else { + // Reset backoff on successful connection + backoff = time.Second + } + } +} + +func (c *Consumer) connectAndConsume(ctx context.Context, endpoint string) error { + // Build WebSocket URL with query parameters + wsURL, err := c.buildWebSocketURL(endpoint) + if err != nil { + return fmt.Errorf("failed to build WebSocket URL: %w", err) + } + + log.Info().Str("url", wsURL).Msg("firehose: connecting to Jetstream") + + // Connect + dialer := websocket.Dialer{ + HandshakeTimeout: 10 * time.Second, + } + + conn, _, err := dialer.DialContext(ctx, wsURL, nil) + if err != nil { + return fmt.Errorf("failed to connect: %w", err) + } + + c.connMu.Lock() + c.conn = conn + c.connMu.Unlock() + + c.connected.Store(true) + log.Info().Str("endpoint", endpoint).Msg("firehose: connected to Jetstream") + + // Mark index as ready once connected + c.index.SetReady(true) + + defer func() { + c.connMu.Lock() + if c.conn != nil { + c.conn.Close() + c.conn = nil + } + c.connMu.Unlock() + c.connected.Store(false) + }() + + // Read events + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-c.stopCh: + return nil + default: + } + + // Set read deadline + conn.SetReadDeadline(time.Now().Add(60 * time.Second)) + + _, message, err := conn.ReadMessage() + if err != nil { + return fmt.Errorf("read error: %w", err) + } + + c.bytesReceived.Add(int64(len(message))) + + if err := c.processMessage(message); err != nil { + log.Warn().Err(err).Msg("firehose: failed to process message") + } + } +} + +func (c *Consumer) buildWebSocketURL(endpoint string) (string, error) { + u, err := url.Parse(endpoint) + if err != nil { + return "", err + } + + q := u.Query() + + // Add wanted collections + for _, coll := range c.config.WantedCollections { + q.Add("wantedCollections", coll) + } + + // Add compression + if c.config.Compress { + q.Set("compress", "true") + } + + // Add cursor if we have one (rewind slightly for safety) + cursor := c.cursor.Load() + if cursor > 0 { + // Rewind by 5 seconds to handle any gaps + cursor -= 5 * time.Second.Microseconds() + q.Set("cursor", fmt.Sprintf("%d", cursor)) + } + + u.RawQuery = q.Encode() + return u.String(), nil +} + +func (c *Consumer) processMessage(data []byte) error { + // Try to decompress if compression is enabled and data looks compressed + // Zstd compressed data starts with magic number 0x28 0xB5 0x2F 0xFD + if c.config.Compress && len(data) >= 4 && data[0] == 0x28 && data[1] == 0xB5 && data[2] == 0x2F && data[3] == 0xFD { + decompressed, err := c.zstdDecoder.DecodeAll(data, nil) + if err != nil { + return fmt.Errorf("failed to decompress message: %w", err) + } + data = decompressed + } else if c.config.Compress && len(data) > 0 && data[0] != '{' { + // Try decompression anyway if it doesn't look like JSON + decompressed, err := c.zstdDecoder.DecodeAll(data, nil) + if err == nil { + data = decompressed + } + // If decompression fails, try parsing as-is (maybe it's uncompressed) + } + + var event JetstreamEvent + if err := json.Unmarshal(data, &event); err != nil { + // Log the first few bytes for debugging + preview := data + if len(preview) > 50 { + preview = preview[:50] + } + return fmt.Errorf("failed to unmarshal event (first bytes: %q): %w", preview, err) + } + + c.eventsReceived.Add(1) + + // Update cursor + if event.TimeUS > 0 { + c.cursor.Store(event.TimeUS) + + // Persist cursor periodically (every 1000 events) + if c.eventsReceived.Load()%1000 == 0 { + if err := c.index.SetCursor(event.TimeUS); err != nil { + log.Warn().Err(err).Msg("firehose: failed to persist cursor") + } + } + } + + // Only process commit events + if event.Kind != "commit" || event.Commit == nil { + return nil + } + + commit := event.Commit + + // Verify it's an Arabica collection + if !strings.HasPrefix(commit.Collection, "social.arabica.alpha.") { + return nil + } + + log.Debug(). + Str("did", event.DID). + Str("collection", commit.Collection). + Str("operation", commit.Operation). + Str("rkey", commit.RKey). + Msg("firehose: processing event") + + switch commit.Operation { + case "create", "update": + if commit.Record == nil { + return nil + } + if err := c.index.UpsertRecord( + event.DID, + commit.Collection, + commit.RKey, + commit.CID, + commit.Record, + event.TimeUS, + ); err != nil { + return fmt.Errorf("failed to upsert record: %w", err) + } + + case "delete": + if err := c.index.DeleteRecord( + event.DID, + commit.Collection, + commit.RKey, + ); err != nil { + return fmt.Errorf("failed to delete record: %w", err) + } + } + + return nil +} + +// BackfillKnownUsers backfills records for all known DIDs +// This is useful on startup to ensure we have all existing records +func (c *Consumer) BackfillKnownUsers(ctx context.Context) error { + dids, err := c.index.GetKnownDIDs() + if err != nil { + return fmt.Errorf("failed to get known DIDs: %w", err) + } + + log.Info().Int("count", len(dids)).Msg("firehose: backfilling known users") + + for _, did := range dids { + select { + case <-ctx.Done(): + return ctx.Err() + default: + } + + if err := c.index.BackfillUser(ctx, did); err != nil { + log.Warn().Err(err).Str("did", did).Msg("firehose: failed to backfill user") + } + + // Small delay to avoid hammering PDS servers + time.Sleep(100 * time.Millisecond) + } + + return nil +} + +// BackfillDID backfills records for a specific DID +func (c *Consumer) BackfillDID(ctx context.Context, did string) error { + return c.index.BackfillUser(ctx, did) +} diff --git a/internal/firehose/index.go b/internal/firehose/index.go new file mode 100644 index 0000000..acacfec --- /dev/null +++ b/internal/firehose/index.go @@ -0,0 +1,726 @@ +package firehose + +import ( + "context" + "encoding/binary" + "encoding/json" + "fmt" + "os" + "path/filepath" + "sort" + "strings" + "sync" + "time" + + "arabica/internal/atproto" + "arabica/internal/models" + + "github.com/rs/zerolog/log" + bolt "go.etcd.io/bbolt" +) + +// Bucket names for the feed index +var ( + // BucketRecords stores full record data: {at-uri} -> {IndexedRecord JSON} + BucketRecords = []byte("records") + + // BucketByTime stores records by timestamp for chronological queries: {timestamp:at-uri} -> {} + BucketByTime = []byte("by_time") + + // BucketByDID stores records by DID for user-specific queries: {did:at-uri} -> {} + BucketByDID = []byte("by_did") + + // BucketByCollection stores records by type: {collection:timestamp:at-uri} -> {} + BucketByCollection = []byte("by_collection") + + // BucketProfiles stores cached profile data: {did} -> {CachedProfile JSON} + BucketProfiles = []byte("profiles") + + // BucketMeta stores metadata like cursor position: {key} -> {value} + BucketMeta = []byte("meta") + + // BucketKnownDIDs stores all DIDs we've seen with Arabica records + BucketKnownDIDs = []byte("known_dids") +) + +// IndexedRecord represents a record stored in the index +type IndexedRecord struct { + URI string `json:"uri"` + DID string `json:"did"` + Collection string `json:"collection"` + RKey string `json:"rkey"` + Record json.RawMessage `json:"record"` + CID string `json:"cid"` + IndexedAt time.Time `json:"indexed_at"` + CreatedAt time.Time `json:"created_at"` // Parsed from record +} + +// CachedProfile stores profile data with TTL +type CachedProfile struct { + Profile *atproto.Profile `json:"profile"` + CachedAt time.Time `json:"cached_at"` + ExpiresAt time.Time `json:"expires_at"` +} + +// FeedIndex provides persistent storage for firehose events +type FeedIndex struct { + db *bolt.DB + publicClient *atproto.PublicClient + profileTTL time.Duration + + // In-memory cache for hot data + profileCache map[string]*CachedProfile + profileCacheMu sync.RWMutex + + ready bool + readyMu sync.RWMutex +} + +// NewFeedIndex creates a new feed index backed by BoltDB +func NewFeedIndex(path string, profileTTL time.Duration) (*FeedIndex, error) { + if path == "" { + return nil, fmt.Errorf("index path is required") + } + + // Ensure parent directory exists + dir := filepath.Dir(path) + if dir != "" && dir != "." { + if err := os.MkdirAll(dir, 0755); err != nil { + return nil, fmt.Errorf("failed to create index directory: %w", err) + } + } + + db, err := bolt.Open(path, 0600, &bolt.Options{ + Timeout: 5 * time.Second, + }) + if err != nil { + return nil, fmt.Errorf("failed to open index database: %w", err) + } + + // Create buckets + err = db.Update(func(tx *bolt.Tx) error { + buckets := [][]byte{ + BucketRecords, + BucketByTime, + BucketByDID, + BucketByCollection, + BucketProfiles, + BucketMeta, + BucketKnownDIDs, + } + for _, bucket := range buckets { + if _, err := tx.CreateBucketIfNotExists(bucket); err != nil { + return fmt.Errorf("failed to create bucket %s: %w", bucket, err) + } + } + return nil + }) + if err != nil { + db.Close() + return nil, err + } + + idx := &FeedIndex{ + db: db, + publicClient: atproto.NewPublicClient(), + profileTTL: profileTTL, + profileCache: make(map[string]*CachedProfile), + } + + return idx, nil +} + +// Close closes the index database +func (idx *FeedIndex) Close() error { + if idx.db != nil { + return idx.db.Close() + } + return nil +} + +// SetReady marks the index as ready to serve queries +func (idx *FeedIndex) SetReady(ready bool) { + idx.readyMu.Lock() + defer idx.readyMu.Unlock() + idx.ready = ready +} + +// IsReady returns true if the index is populated and ready +func (idx *FeedIndex) IsReady() bool { + idx.readyMu.RLock() + defer idx.readyMu.RUnlock() + return idx.ready +} + +// GetCursor returns the last processed cursor (microseconds timestamp) +func (idx *FeedIndex) GetCursor() (int64, error) { + var cursor int64 + err := idx.db.View(func(tx *bolt.Tx) error { + b := tx.Bucket(BucketMeta) + v := b.Get([]byte("cursor")) + if v != nil && len(v) == 8 { + cursor = int64(binary.BigEndian.Uint64(v)) + } + return nil + }) + return cursor, err +} + +// SetCursor stores the cursor position +func (idx *FeedIndex) SetCursor(cursor int64) error { + return idx.db.Update(func(tx *bolt.Tx) error { + b := tx.Bucket(BucketMeta) + buf := make([]byte, 8) + binary.BigEndian.PutUint64(buf, uint64(cursor)) + return b.Put([]byte("cursor"), buf) + }) +} + +// UpsertRecord adds or updates a record in the index +func (idx *FeedIndex) UpsertRecord(did, collection, rkey, cid string, record json.RawMessage, eventTime int64) error { + uri := atproto.BuildATURI(did, collection, rkey) + + // Parse createdAt from record + var recordData map[string]interface{} + createdAt := time.Now() + if err := json.Unmarshal(record, &recordData); err == nil { + if createdAtStr, ok := recordData["createdAt"].(string); ok { + if t, err := time.Parse(time.RFC3339, createdAtStr); err == nil { + createdAt = t + } + } + } + + indexed := &IndexedRecord{ + URI: uri, + DID: did, + Collection: collection, + RKey: rkey, + Record: record, + CID: cid, + IndexedAt: time.Now(), + CreatedAt: createdAt, + } + + data, err := json.Marshal(indexed) + if err != nil { + return fmt.Errorf("failed to marshal record: %w", err) + } + + return idx.db.Update(func(tx *bolt.Tx) error { + // Store the record + records := tx.Bucket(BucketRecords) + if err := records.Put([]byte(uri), data); err != nil { + return err + } + + // Index by time (use createdAt for sorting, not event time) + byTime := tx.Bucket(BucketByTime) + timeKey := makeTimeKey(createdAt, uri) + if err := byTime.Put(timeKey, nil); err != nil { + return err + } + + // Index by DID + byDID := tx.Bucket(BucketByDID) + didKey := []byte(did + ":" + uri) + if err := byDID.Put(didKey, nil); err != nil { + return err + } + + // Index by collection + byCollection := tx.Bucket(BucketByCollection) + collKey := []byte(collection + ":" + string(timeKey)) + if err := byCollection.Put(collKey, nil); err != nil { + return err + } + + // Track known DID + knownDIDs := tx.Bucket(BucketKnownDIDs) + if err := knownDIDs.Put([]byte(did), []byte("1")); err != nil { + return err + } + + return nil + }) +} + +// DeleteRecord removes a record from the index +func (idx *FeedIndex) DeleteRecord(did, collection, rkey string) error { + uri := atproto.BuildATURI(did, collection, rkey) + + return idx.db.Update(func(tx *bolt.Tx) error { + // Get the existing record to find its timestamp + records := tx.Bucket(BucketRecords) + existingData := records.Get([]byte(uri)) + if existingData == nil { + // Record doesn't exist, nothing to delete + return nil + } + + var existing IndexedRecord + if err := json.Unmarshal(existingData, &existing); err != nil { + // Can't parse, just delete the main record + return records.Delete([]byte(uri)) + } + + // Delete from records + if err := records.Delete([]byte(uri)); err != nil { + return err + } + + // Delete from by_time index + byTime := tx.Bucket(BucketByTime) + timeKey := makeTimeKey(existing.CreatedAt, uri) + if err := byTime.Delete(timeKey); err != nil { + return err + } + + // Delete from by_did index + byDID := tx.Bucket(BucketByDID) + didKey := []byte(did + ":" + uri) + if err := byDID.Delete(didKey); err != nil { + return err + } + + // Delete from by_collection index + byCollection := tx.Bucket(BucketByCollection) + collKey := []byte(collection + ":" + string(timeKey)) + if err := byCollection.Delete(collKey); err != nil { + return err + } + + return nil + }) +} + +// GetRecord retrieves a single record by URI +func (idx *FeedIndex) GetRecord(uri string) (*IndexedRecord, error) { + var record *IndexedRecord + err := idx.db.View(func(tx *bolt.Tx) error { + b := tx.Bucket(BucketRecords) + data := b.Get([]byte(uri)) + if data == nil { + return nil + } + record = &IndexedRecord{} + return json.Unmarshal(data, record) + }) + return record, err +} + +// FeedItem represents an item in the feed (matches feed.FeedItem structure) +type FeedItem struct { + RecordType string + Action string + + Brew *models.Brew + Bean *models.Bean + Roaster *models.Roaster + Grinder *models.Grinder + Brewer *models.Brewer + + Author *atproto.Profile + Timestamp time.Time + TimeAgo string +} + +// GetRecentFeed returns recent feed items from the index +func (idx *FeedIndex) GetRecentFeed(ctx context.Context, limit int) ([]*FeedItem, error) { + var records []*IndexedRecord + + err := idx.db.View(func(tx *bolt.Tx) error { + byTime := tx.Bucket(BucketByTime) + recordsBucket := tx.Bucket(BucketRecords) + + c := byTime.Cursor() + + // Iterate in reverse (newest first) + count := 0 + for k, _ := c.Last(); k != nil && count < limit*2; k, _ = c.Prev() { + // Extract URI from key (format: timestamp:uri) + uri := extractURIFromTimeKey(k) + if uri == "" { + continue + } + + data := recordsBucket.Get([]byte(uri)) + if data == nil { + continue + } + + var record IndexedRecord + if err := json.Unmarshal(data, &record); err != nil { + continue + } + + records = append(records, &record) + count++ + } + + return nil + }) + if err != nil { + return nil, err + } + + // Build lookup maps for reference resolution + recordsByURI := make(map[string]*IndexedRecord) + for _, r := range records { + recordsByURI[r.URI] = r + } + + // Also load additional records we might need for references + err = idx.db.View(func(tx *bolt.Tx) error { + recordsBucket := tx.Bucket(BucketRecords) + return recordsBucket.ForEach(func(k, v []byte) error { + uri := string(k) + if _, exists := recordsByURI[uri]; exists { + return nil + } + var record IndexedRecord + if err := json.Unmarshal(v, &record); err != nil { + return nil + } + // Only load beans, roasters, grinders, brewers for reference resolution + switch record.Collection { + case atproto.NSIDBean, atproto.NSIDRoaster, atproto.NSIDGrinder, atproto.NSIDBrewer: + recordsByURI[uri] = &record + } + return nil + }) + }) + if err != nil { + return nil, err + } + + // Convert to FeedItems + items := make([]*FeedItem, 0, len(records)) + for _, record := range records { + item, err := idx.recordToFeedItem(ctx, record, recordsByURI) + if err != nil { + log.Warn().Err(err).Str("uri", record.URI).Msg("failed to convert record to feed item") + continue + } + items = append(items, item) + } + + // Sort by timestamp descending + sort.Slice(items, func(i, j int) bool { + return items[i].Timestamp.After(items[j].Timestamp) + }) + + // Apply limit + if len(items) > limit { + items = items[:limit] + } + + return items, nil +} + +// recordToFeedItem converts an IndexedRecord to a FeedItem +func (idx *FeedIndex) recordToFeedItem(ctx context.Context, record *IndexedRecord, refMap map[string]*IndexedRecord) (*FeedItem, error) { + var recordData map[string]interface{} + if err := json.Unmarshal(record.Record, &recordData); err != nil { + return nil, err + } + + item := &FeedItem{ + Timestamp: record.CreatedAt, + TimeAgo: formatTimeAgo(record.CreatedAt), + } + + // Get author profile + profile, err := idx.GetProfile(ctx, record.DID) + if err != nil { + log.Warn().Err(err).Str("did", record.DID).Msg("failed to get profile") + // Use a placeholder profile + profile = &atproto.Profile{ + DID: record.DID, + Handle: record.DID, // Use DID as handle if we can't resolve + } + } + item.Author = profile + + switch record.Collection { + case atproto.NSIDBrew: + brew, err := atproto.RecordToBrew(recordData, record.URI) + if err != nil { + return nil, err + } + + // Resolve bean reference + if beanRef, ok := recordData["beanRef"].(string); ok && beanRef != "" { + if beanRecord, found := refMap[beanRef]; found { + var beanData map[string]interface{} + if err := json.Unmarshal(beanRecord.Record, &beanData); err == nil { + bean, _ := atproto.RecordToBean(beanData, beanRef) + brew.Bean = bean + + // Resolve roaster reference for bean + if roasterRef, ok := beanData["roasterRef"].(string); ok && roasterRef != "" { + if roasterRecord, found := refMap[roasterRef]; found { + var roasterData map[string]interface{} + if err := json.Unmarshal(roasterRecord.Record, &roasterData); err == nil { + roaster, _ := atproto.RecordToRoaster(roasterData, roasterRef) + brew.Bean.Roaster = roaster + } + } + } + } + } + } + + // Resolve grinder reference + if grinderRef, ok := recordData["grinderRef"].(string); ok && grinderRef != "" { + if grinderRecord, found := refMap[grinderRef]; found { + var grinderData map[string]interface{} + if err := json.Unmarshal(grinderRecord.Record, &grinderData); err == nil { + grinder, _ := atproto.RecordToGrinder(grinderData, grinderRef) + brew.GrinderObj = grinder + } + } + } + + // Resolve brewer reference + if brewerRef, ok := recordData["brewerRef"].(string); ok && brewerRef != "" { + if brewerRecord, found := refMap[brewerRef]; found { + var brewerData map[string]interface{} + if err := json.Unmarshal(brewerRecord.Record, &brewerData); err == nil { + brewer, _ := atproto.RecordToBrewer(brewerData, brewerRef) + brew.BrewerObj = brewer + } + } + } + + item.RecordType = "brew" + item.Action = "added a new brew" + item.Brew = brew + + case atproto.NSIDBean: + bean, err := atproto.RecordToBean(recordData, record.URI) + if err != nil { + return nil, err + } + + // Resolve roaster reference + if roasterRef, ok := recordData["roasterRef"].(string); ok && roasterRef != "" { + if roasterRecord, found := refMap[roasterRef]; found { + var roasterData map[string]interface{} + if err := json.Unmarshal(roasterRecord.Record, &roasterData); err == nil { + roaster, _ := atproto.RecordToRoaster(roasterData, roasterRef) + bean.Roaster = roaster + } + } + } + + item.RecordType = "bean" + item.Action = "added a new bean" + item.Bean = bean + + case atproto.NSIDRoaster: + roaster, err := atproto.RecordToRoaster(recordData, record.URI) + if err != nil { + return nil, err + } + item.RecordType = "roaster" + item.Action = "added a new roaster" + item.Roaster = roaster + + case atproto.NSIDGrinder: + grinder, err := atproto.RecordToGrinder(recordData, record.URI) + if err != nil { + return nil, err + } + item.RecordType = "grinder" + item.Action = "added a new grinder" + item.Grinder = grinder + + case atproto.NSIDBrewer: + brewer, err := atproto.RecordToBrewer(recordData, record.URI) + if err != nil { + return nil, err + } + item.RecordType = "brewer" + item.Action = "added a new brewer" + item.Brewer = brewer + + default: + return nil, fmt.Errorf("unknown collection: %s", record.Collection) + } + + return item, nil +} + +// GetProfile fetches a profile, using cache when possible +func (idx *FeedIndex) GetProfile(ctx context.Context, did string) (*atproto.Profile, error) { + // Check in-memory cache first + idx.profileCacheMu.RLock() + if cached, ok := idx.profileCache[did]; ok && time.Now().Before(cached.ExpiresAt) { + idx.profileCacheMu.RUnlock() + return cached.Profile, nil + } + idx.profileCacheMu.RUnlock() + + // Check persistent cache + var cached *CachedProfile + err := idx.db.View(func(tx *bolt.Tx) error { + b := tx.Bucket(BucketProfiles) + data := b.Get([]byte(did)) + if data == nil { + return nil + } + cached = &CachedProfile{} + return json.Unmarshal(data, cached) + }) + if err == nil && cached != nil && time.Now().Before(cached.ExpiresAt) { + // Update in-memory cache + idx.profileCacheMu.Lock() + idx.profileCache[did] = cached + idx.profileCacheMu.Unlock() + return cached.Profile, nil + } + + // Fetch from API + profile, err := idx.publicClient.GetProfile(ctx, did) + if err != nil { + return nil, err + } + + // Cache the result + now := time.Now() + cached = &CachedProfile{ + Profile: profile, + CachedAt: now, + ExpiresAt: now.Add(idx.profileTTL), + } + + // Update in-memory cache + idx.profileCacheMu.Lock() + idx.profileCache[did] = cached + idx.profileCacheMu.Unlock() + + // Persist to database + data, _ := json.Marshal(cached) + _ = idx.db.Update(func(tx *bolt.Tx) error { + b := tx.Bucket(BucketProfiles) + return b.Put([]byte(did), data) + }) + + return profile, nil +} + +// GetKnownDIDs returns all DIDs that have created Arabica records +func (idx *FeedIndex) GetKnownDIDs() ([]string, error) { + var dids []string + err := idx.db.View(func(tx *bolt.Tx) error { + b := tx.Bucket(BucketKnownDIDs) + return b.ForEach(func(k, v []byte) error { + dids = append(dids, string(k)) + return nil + }) + }) + return dids, err +} + +// RecordCount returns the total number of indexed records +func (idx *FeedIndex) RecordCount() int { + var count int + _ = idx.db.View(func(tx *bolt.Tx) error { + b := tx.Bucket(BucketRecords) + count = b.Stats().KeyN + return nil + }) + return count +} + +// Helper functions + +func makeTimeKey(t time.Time, uri string) []byte { + // Format: inverted timestamp (for reverse chronological order) + ":" + uri + // Use nanoseconds for uniqueness + inverted := ^uint64(t.UnixNano()) + buf := make([]byte, 8) + binary.BigEndian.PutUint64(buf, inverted) + return append(buf, []byte(":"+uri)...) +} + +func extractURIFromTimeKey(key []byte) string { + if len(key) < 10 { // 8 bytes timestamp + ":" + at least 1 char + return "" + } + // Skip 8 bytes timestamp + 1 byte ":" + return string(key[9:]) +} + +func formatTimeAgo(t time.Time) string { + now := time.Now() + diff := now.Sub(t) + + switch { + case diff < time.Minute: + return "just now" + case diff < time.Hour: + mins := int(diff.Minutes()) + if mins == 1 { + return "1 minute ago" + } + return fmt.Sprintf("%d minutes ago", mins) + case diff < 24*time.Hour: + hours := int(diff.Hours()) + if hours == 1 { + return "1 hour ago" + } + return fmt.Sprintf("%d hours ago", hours) + case diff < 48*time.Hour: + return "yesterday" + case diff < 7*24*time.Hour: + days := int(diff.Hours() / 24) + return fmt.Sprintf("%d days ago", days) + case diff < 30*24*time.Hour: + weeks := int(diff.Hours() / 24 / 7) + if weeks == 1 { + return "1 week ago" + } + return fmt.Sprintf("%d weeks ago", weeks) + default: + months := int(diff.Hours() / 24 / 30) + if months == 1 { + return "1 month ago" + } + return fmt.Sprintf("%d months ago", months) + } +} + +// BackfillUser fetches all existing records for a DID and adds them to the index +func (idx *FeedIndex) BackfillUser(ctx context.Context, did string) error { + log.Info().Str("did", did).Msg("backfilling user records") + + for _, collection := range ArabicaCollections { + records, err := idx.publicClient.ListRecords(ctx, did, collection, 100) + if err != nil { + log.Warn().Err(err).Str("did", did).Str("collection", collection).Msg("failed to list records for backfill") + continue + } + + for _, record := range records.Records { + // Extract rkey from URI + parts := strings.Split(record.URI, "/") + if len(parts) < 3 { + continue + } + rkey := parts[len(parts)-1] + + recordJSON, err := json.Marshal(record.Value) + if err != nil { + continue + } + + if err := idx.UpsertRecord(did, collection, rkey, record.CID, recordJSON, 0); err != nil { + log.Warn().Err(err).Str("uri", record.URI).Msg("failed to upsert record during backfill") + } + } + } + + return nil +} diff --git a/justfile b/justfile index 48fd720..9b468fb 100644 --- a/justfile +++ b/justfile @@ -4,6 +4,9 @@ run: run-production: @LOG_FORMAT=json SECURE_COOKIES=true go run cmd/server/main.go +run-firehose: + @LOG_LEVEL=debug LOG_FORMAT=console go run cmd/server/main.go -firehose + test: @go test ./... -cover -coverprofile=cover.out diff --git a/module.nix b/module.nix index 1e58169..481a5aa 100644 --- a/module.nix +++ b/module.nix @@ -36,6 +36,16 @@ in { default = true; description = "Whether to set the Secure flag on cookies. Should be true when using HTTPS."; }; + + firehose = lib.mkOption { + type = lib.types.bool; + default = false; + description = '' + Enable firehose-based feed using Jetstream. + This provides real-time feed updates with zero API calls per request, + instead of polling each user's PDS. + ''; + }; }; oauth = { @@ -103,7 +113,7 @@ in { Type = "simple"; User = cfg.user; Group = cfg.group; - ExecStart = "${cfg.package}/bin/arabica"; + ExecStart = "${cfg.package}/bin/arabica${lib.optionalString cfg.settings.firehose " -firehose"}"; Restart = "on-failure"; RestartSec = "10s";