mirror of
https://tangled.org/evan.jarrett.net/at-container-registry
synced 2026-09-03 16:56:56 +00:00
753 lines
21 KiB
Go
753 lines
21 KiB
Go
package pds
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"database/sql"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log/slog"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
atproto "github.com/bluesky-social/indigo/api/atproto"
|
|
"github.com/bluesky-social/indigo/events"
|
|
lexutil "github.com/bluesky-social/indigo/lex/util"
|
|
"github.com/bluesky-social/indigo/repo"
|
|
"github.com/gorilla/websocket"
|
|
"github.com/ipfs/go-cid"
|
|
"github.com/ipld/go-car"
|
|
carutil "github.com/ipld/go-car/util"
|
|
_ "github.com/tursodatabase/go-libsql"
|
|
)
|
|
|
|
// EventBroadcaster manages WebSocket connections and broadcasts repo events
|
|
type EventBroadcaster struct {
|
|
mu sync.RWMutex
|
|
subscribers map[*Subscriber]bool
|
|
eventSeq int64
|
|
eventHistory []HistoricalEvent // Ring buffer for cursor backfill (deprecated, kept for compatibility)
|
|
maxHistory int
|
|
holdDID string // DID of the hold for setting repo field
|
|
db *sql.DB // Database for persistent event storage
|
|
dbPath string // Path to database file
|
|
ownsDB bool // true when this broadcaster opened the connection itself
|
|
}
|
|
|
|
// Subscriber represents a WebSocket client subscribed to the firehose
|
|
type Subscriber struct {
|
|
conn *websocket.Conn
|
|
send chan *RepoCommitEvent
|
|
cursor int64 // Last sequence number this subscriber has seen
|
|
}
|
|
|
|
// HistoricalEvent stores past events for cursor-based backfill
|
|
type HistoricalEvent struct {
|
|
Seq int64
|
|
Event *RepoCommitEvent
|
|
}
|
|
|
|
// RepoCommitEvent represents a #commit event in subscribeRepos
|
|
type RepoCommitEvent struct {
|
|
Seq int64 `json:"seq" cborgen:"seq"`
|
|
Repo string `json:"repo" cborgen:"repo"`
|
|
Commit string `json:"commit" cborgen:"commit"` // CID string
|
|
Rev string `json:"rev" cborgen:"rev"`
|
|
Since *string `json:"since,omitempty" cborgen:"since,omitempty"`
|
|
Blocks []byte `json:"blocks" cborgen:"blocks"` // CAR slice bytes
|
|
Ops []*atproto.SyncSubscribeRepos_RepoOp `json:"ops" cborgen:"ops"`
|
|
Time string `json:"time" cborgen:"time"`
|
|
Type string `json:"$type" cborgen:"$type"` // Always "#commit"
|
|
}
|
|
|
|
// NewEventBroadcaster creates a new event broadcaster with persistent storage
|
|
// dbPath should point to the carstore database file (e.g., "/path/to/pds/db.sqlite3")
|
|
func NewEventBroadcaster(holdDID string, maxHistory int, dbPath string) *EventBroadcaster {
|
|
if maxHistory <= 0 {
|
|
maxHistory = 100 // Default to keeping 100 events
|
|
}
|
|
|
|
broadcaster := &EventBroadcaster{
|
|
subscribers: make(map[*Subscriber]bool),
|
|
eventSeq: 0,
|
|
eventHistory: make([]HistoricalEvent, 0, maxHistory),
|
|
maxHistory: maxHistory,
|
|
holdDID: holdDID,
|
|
dbPath: dbPath,
|
|
}
|
|
|
|
// Initialize database connection and schema
|
|
if dbPath != "" && dbPath != ":memory:" {
|
|
if err := broadcaster.initDatabase(); err != nil {
|
|
slog.Warn("Failed to initialize event database", "error", err)
|
|
slog.Warn("Events will not persist across restarts")
|
|
}
|
|
}
|
|
|
|
return broadcaster
|
|
}
|
|
|
|
// NewEventBroadcasterWithDB creates an event broadcaster using an existing *sql.DB connection.
|
|
// The caller is responsible for the DB lifecycle.
|
|
func NewEventBroadcasterWithDB(holdDID string, maxHistory int, db *sql.DB) *EventBroadcaster {
|
|
if maxHistory <= 0 {
|
|
maxHistory = 100
|
|
}
|
|
|
|
broadcaster := &EventBroadcaster{
|
|
subscribers: make(map[*Subscriber]bool),
|
|
eventSeq: 0,
|
|
eventHistory: make([]HistoricalEvent, 0, maxHistory),
|
|
maxHistory: maxHistory,
|
|
holdDID: holdDID,
|
|
db: db,
|
|
ownsDB: false,
|
|
}
|
|
|
|
if db != nil {
|
|
if err := broadcaster.initSchema(); err != nil {
|
|
slog.Warn("Failed to initialize event schema", "error", err)
|
|
slog.Warn("Events will not persist across restarts")
|
|
}
|
|
}
|
|
|
|
return broadcaster
|
|
}
|
|
|
|
// initDatabase opens database connection, creates table, and loads last sequence
|
|
func (b *EventBroadcaster) initDatabase() error {
|
|
// Open database connection
|
|
dsn := b.dbPath
|
|
if b.dbPath != ":memory:" && !strings.HasPrefix(b.dbPath, "file:") {
|
|
dsn = "file:" + b.dbPath
|
|
}
|
|
db, err := sql.Open("libsql", dsn)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Test connection
|
|
if err := db.Ping(); err != nil {
|
|
db.Close()
|
|
return err
|
|
}
|
|
|
|
b.db = db
|
|
b.ownsDB = true
|
|
|
|
if err := b.initSchema(); err != nil {
|
|
db.Close()
|
|
b.db = nil
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// initSchema creates the events table and loads last sequence number
|
|
func (b *EventBroadcaster) initSchema() error {
|
|
// Create events table if it doesn't exist
|
|
// Execute statements individually for go-libsql compatibility
|
|
stmts := []string{
|
|
`CREATE TABLE IF NOT EXISTS firehose_events (
|
|
seq INTEGER PRIMARY KEY,
|
|
commit_cid TEXT NOT NULL,
|
|
rev TEXT NOT NULL,
|
|
since_rev TEXT,
|
|
repo_slice BLOB NOT NULL,
|
|
ops_json TEXT NOT NULL,
|
|
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
|
|
)`,
|
|
`CREATE INDEX IF NOT EXISTS idx_firehose_events_rev ON firehose_events(rev)`,
|
|
}
|
|
|
|
for _, stmt := range stmts {
|
|
if _, err := b.db.Exec(stmt); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
|
|
// Load last sequence number from database
|
|
var lastSeq sql.NullInt64
|
|
err := b.db.QueryRow("SELECT MAX(seq) FROM firehose_events").Scan(&lastSeq)
|
|
if err != nil {
|
|
slog.Warn("Failed to load last event sequence", "error", err)
|
|
} else if lastSeq.Valid {
|
|
b.eventSeq = lastSeq.Int64
|
|
slog.Info("Loaded event sequence from database", "seq", b.eventSeq)
|
|
} else {
|
|
// Database is empty but might have existing repo records
|
|
// This happens on first deployment after adding persistent events
|
|
slog.Info("No events in database - will bootstrap from repo if needed")
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// extractCreatedAt attempts to extract createdAt timestamp from various record types
|
|
func extractCreatedAt(recordValue any) time.Time {
|
|
// Try known types with createdAt fields
|
|
switch v := recordValue.(type) {
|
|
case map[string]any:
|
|
// Generic map fallback (shouldn't happen with lexutil but just in case)
|
|
if createdAtStr, ok := v["createdAt"].(string); ok {
|
|
if parsed, err := time.Parse(time.RFC3339, createdAtStr); err == nil {
|
|
return parsed
|
|
}
|
|
}
|
|
}
|
|
|
|
// Try to extract via JSON marshal/unmarshal as last resort
|
|
// This works for any struct with a createdAt field
|
|
jsonBytes, err := json.Marshal(recordValue)
|
|
if err == nil {
|
|
var generic map[string]any
|
|
if err := json.Unmarshal(jsonBytes, &generic); err == nil {
|
|
if createdAtStr, ok := generic["createdAt"].(string); ok {
|
|
if parsed, err := time.Parse(time.RFC3339, createdAtStr); err == nil {
|
|
return parsed
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Default to current time if no createdAt found
|
|
return time.Now()
|
|
}
|
|
|
|
// BootstrapFromRepo generates synthetic events from all current records in the repo
|
|
// This is called once when deploying persistent events to an existing repo
|
|
func (b *EventBroadcaster) BootstrapFromRepo(pds *HoldPDS) error {
|
|
if b.db == nil {
|
|
return fmt.Errorf("database not initialized")
|
|
}
|
|
|
|
// Check if we already have events
|
|
var count int64
|
|
if err := b.db.QueryRow("SELECT COUNT(*) FROM firehose_events").Scan(&count); err != nil {
|
|
return fmt.Errorf("failed to check event count: %w", err)
|
|
}
|
|
|
|
if count > 0 {
|
|
slog.Info("Database already has events, skipping bootstrap", "count", count)
|
|
return nil
|
|
}
|
|
|
|
ctx := context.Background()
|
|
|
|
// Get current repo state
|
|
session, err := pds.carstore.ReadOnlySession(pds.uid)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to create session: %w", err)
|
|
}
|
|
|
|
head, err := pds.carstore.GetUserRepoHead(ctx, pds.uid)
|
|
if err != nil || !head.Defined() {
|
|
// Empty repo, nothing to bootstrap
|
|
slog.Info("Empty repo, no events to bootstrap")
|
|
return nil
|
|
}
|
|
|
|
repoHandle, err := repo.OpenRepo(ctx, session, head)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to open repo: %w", err)
|
|
}
|
|
|
|
// Get current rev
|
|
rev, err := pds.repomgr.GetRepoRev(ctx, pds.uid)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to get repo rev: %w", err)
|
|
}
|
|
|
|
slog.Info("Bootstrapping firehose events from current repo state",
|
|
"head", head.String(),
|
|
"rev", rev)
|
|
|
|
var recordCount int64
|
|
|
|
// Walk all records in the repo and create synthetic events
|
|
err = repoHandle.ForEach(ctx, "", func(path string, recordCID cid.Cid) error {
|
|
// Get record value
|
|
_, recBytes, err := repoHandle.GetRecordBytes(ctx, path)
|
|
if err != nil {
|
|
slog.Warn("Failed to get record bytes", "path", path, "error", err)
|
|
return nil // Skip this record but continue
|
|
}
|
|
|
|
recordValue, err := lexutil.CborDecodeValue(*recBytes)
|
|
if err != nil {
|
|
slog.Warn("Failed to decode record", "path", path, "error", err)
|
|
return nil
|
|
}
|
|
|
|
// Parse collection and rkey from path
|
|
parts := strings.Split(path, "/")
|
|
if len(parts) < 2 {
|
|
return nil // Invalid path
|
|
}
|
|
|
|
collection := strings.Join(parts[:len(parts)-1], "/")
|
|
rkey := parts[len(parts)-1]
|
|
|
|
// Extract createdAt timestamp from record if it exists
|
|
// Try to extract from known record types with createdAt fields
|
|
recordTime := extractCreatedAt(recordValue)
|
|
|
|
// Create synthetic RepoOp
|
|
ops := []RepoOp{
|
|
{
|
|
Kind: EvtKindCreateRecord,
|
|
Collection: collection,
|
|
Rkey: rkey,
|
|
RecCid: &recordCID,
|
|
Record: recordValue,
|
|
},
|
|
}
|
|
|
|
// Get CAR slice for this record (minimal - just the record block)
|
|
var carBuf bytes.Buffer
|
|
carHeader := &car.CarHeader{
|
|
Roots: []cid.Cid{head},
|
|
Version: 1,
|
|
}
|
|
if err := car.WriteHeader(carHeader, &carBuf); err != nil {
|
|
slog.Warn("Failed to write CAR header", "error", err)
|
|
return nil
|
|
}
|
|
|
|
// Write the record block
|
|
if err := carutil.LdWrite(&carBuf, recordCID.Bytes(), *recBytes); err != nil {
|
|
slog.Warn("Failed to write record block", "error", err)
|
|
return nil
|
|
}
|
|
|
|
// Create synthetic RepoEvent
|
|
repoEvent := &RepoEvent{
|
|
NewRoot: head,
|
|
Rev: rev,
|
|
Since: nil, // No "since" for bootstrap events
|
|
RepoSlice: carBuf.Bytes(),
|
|
Ops: ops,
|
|
}
|
|
|
|
// Convert to commit event and persist
|
|
b.mu.Lock()
|
|
b.eventSeq++
|
|
seq := b.eventSeq
|
|
commitEvent := b.convertToCommitEvent(repoEvent, seq)
|
|
|
|
// Override event time with record's original createdAt
|
|
commitEvent.Time = recordTime.Format(time.RFC3339)
|
|
|
|
// Persist to database
|
|
if err := b.persistEvent(commitEvent); err != nil {
|
|
b.mu.Unlock()
|
|
return fmt.Errorf("failed to persist bootstrap event seq=%d: %w", seq, err)
|
|
}
|
|
|
|
// Also add to in-memory history for immediate use
|
|
b.addToHistory(seq, commitEvent)
|
|
b.mu.Unlock()
|
|
|
|
recordCount++
|
|
return nil
|
|
})
|
|
|
|
if err != nil {
|
|
return fmt.Errorf("failed to walk repo: %w", err)
|
|
}
|
|
|
|
slog.Info("Bootstrapped events from repo",
|
|
"recordCount", recordCount,
|
|
"seq", b.eventSeq)
|
|
return nil
|
|
}
|
|
|
|
// Close closes the database connection
|
|
func (b *EventBroadcaster) Close() error {
|
|
if b.db != nil && b.ownsDB {
|
|
return b.db.Close()
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Subscribe adds a new WebSocket subscriber
|
|
func (b *EventBroadcaster) Subscribe(conn *websocket.Conn, cursor int64, userAgent string) *Subscriber {
|
|
sub := &Subscriber{
|
|
conn: conn,
|
|
send: make(chan *RepoCommitEvent, 10), // Buffer 10 events
|
|
cursor: cursor,
|
|
}
|
|
|
|
b.mu.Lock()
|
|
b.subscribers[sub] = true
|
|
currentSeq := b.eventSeq
|
|
b.mu.Unlock()
|
|
|
|
slog.Info("New firehose subscriber", "remote", conn.RemoteAddr(), "cursor", cursor, "currentSeq", currentSeq, "userAgent", userAgent)
|
|
|
|
// Handle cursor-based backfill:
|
|
// - cursor < 0: No backfill, stream new events only
|
|
// - cursor >= 0: Backfill events from cursor onwards
|
|
// - cursor=0: Replay all events from beginning
|
|
// - cursor < currentSeq: Normal backfill
|
|
// - cursor >= currentSeq: Relay reconnecting after our restart, backfill from database
|
|
if cursor >= 0 {
|
|
if cursor < currentSeq {
|
|
// Normal case: relay is behind, backfill missing events
|
|
go b.backfillSubscriber(sub, cursor)
|
|
} else if cursor > currentSeq {
|
|
// Relay has cursor ahead of us - server was restarted
|
|
// Database should have the events if we had them before
|
|
slog.Info("Relay cursor ahead of current seq, attempting database backfill",
|
|
"cursor", cursor,
|
|
"currentSeq", currentSeq)
|
|
go b.backfillSubscriber(sub, cursor)
|
|
}
|
|
// else cursor == currentSeq: relay is caught up, just stream new events
|
|
}
|
|
|
|
// Start goroutine to handle sending events to this subscriber
|
|
go b.handleSubscriber(sub)
|
|
|
|
return sub
|
|
}
|
|
|
|
// Unsubscribe removes a WebSocket subscriber
|
|
func (b *EventBroadcaster) Unsubscribe(sub *Subscriber) {
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
|
|
if _, ok := b.subscribers[sub]; ok {
|
|
delete(b.subscribers, sub)
|
|
close(sub.send)
|
|
}
|
|
}
|
|
|
|
// Broadcast sends an event to all subscribers
|
|
func (b *EventBroadcaster) Broadcast(ctx context.Context, event *RepoEvent) {
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
|
|
// Increment sequence
|
|
b.eventSeq++
|
|
seq := b.eventSeq
|
|
|
|
// Convert RepoEvent to RepoCommitEvent
|
|
commitEvent := b.convertToCommitEvent(event, seq)
|
|
|
|
// Persist event to database
|
|
if b.db != nil {
|
|
if err := b.persistEvent(commitEvent); err != nil {
|
|
slog.Warn("Failed to persist event to database",
|
|
"seq", seq,
|
|
"error", err)
|
|
}
|
|
}
|
|
|
|
// Store in history for backfill (deprecated, but kept for compatibility)
|
|
b.addToHistory(seq, commitEvent)
|
|
|
|
// Broadcast to all subscribers
|
|
for sub := range b.subscribers {
|
|
select {
|
|
case sub.send <- commitEvent:
|
|
// Sent successfully
|
|
default:
|
|
// Subscriber's buffer is full, skip (they'll get disconnected for being too slow)
|
|
slog.Warn("Subscriber buffer full, skipping event", "seq", seq)
|
|
}
|
|
}
|
|
}
|
|
|
|
// persistEvent stores an event in the database
|
|
func (b *EventBroadcaster) persistEvent(event *RepoCommitEvent) error {
|
|
// Serialize ops to JSON
|
|
opsJSON, err := json.Marshal(event.Ops)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
// Get since_rev value (may be nil)
|
|
var sinceRev sql.NullString
|
|
if event.Since != nil {
|
|
sinceRev = sql.NullString{String: *event.Since, Valid: true}
|
|
}
|
|
|
|
// Insert event
|
|
query := `
|
|
INSERT INTO firehose_events (seq, commit_cid, rev, since_rev, repo_slice, ops_json)
|
|
VALUES (?, ?, ?, ?, ?, ?)
|
|
`
|
|
_, err = b.db.Exec(query, event.Seq, event.Commit, event.Rev, sinceRev, event.Blocks, opsJSON)
|
|
return err
|
|
}
|
|
|
|
// convertToCommitEvent converts a RepoEvent to a RepoCommitEvent
|
|
func (b *EventBroadcaster) convertToCommitEvent(event *RepoEvent, seq int64) *RepoCommitEvent {
|
|
// Convert RepoOps to atproto.SyncSubscribeRepos_RepoOp
|
|
ops := make([]*atproto.SyncSubscribeRepos_RepoOp, len(event.Ops))
|
|
for i, op := range event.Ops {
|
|
action := string(op.Kind) // "create", "update", "delete"
|
|
path := op.Collection + "/" + op.Rkey
|
|
|
|
// Convert CID to LexLink if present
|
|
var cidLink *lexutil.LexLink
|
|
if op.RecCid != nil {
|
|
link := lexutil.LexLink(*op.RecCid)
|
|
cidLink = &link
|
|
}
|
|
|
|
ops[i] = &atproto.SyncSubscribeRepos_RepoOp{
|
|
Action: action,
|
|
Path: path,
|
|
Cid: cidLink,
|
|
}
|
|
}
|
|
|
|
// Event.NewRoot is a cid.Cid, convert to string
|
|
commitCID := event.NewRoot.String()
|
|
|
|
return &RepoCommitEvent{
|
|
Seq: seq,
|
|
Repo: b.holdDID, // Set to hold's DID
|
|
Commit: commitCID,
|
|
Rev: event.Rev,
|
|
Since: event.Since,
|
|
Blocks: event.RepoSlice, // CAR slice bytes
|
|
Ops: ops,
|
|
Time: time.Now().Format(time.RFC3339),
|
|
Type: "#commit",
|
|
}
|
|
}
|
|
|
|
// addToHistory adds an event to the history ring buffer
|
|
func (b *EventBroadcaster) addToHistory(seq int64, event *RepoCommitEvent) {
|
|
he := HistoricalEvent{
|
|
Seq: seq,
|
|
Event: event,
|
|
}
|
|
|
|
// Simple ring buffer: keep last N events
|
|
if len(b.eventHistory) >= b.maxHistory {
|
|
// Remove oldest event
|
|
b.eventHistory = b.eventHistory[1:]
|
|
}
|
|
b.eventHistory = append(b.eventHistory, he)
|
|
}
|
|
|
|
// backfillSubscriber sends historical events to a subscriber
|
|
// Query events from database where seq > cursor
|
|
func (b *EventBroadcaster) backfillSubscriber(sub *Subscriber, cursor int64) {
|
|
// If database is available, use it for backfill
|
|
if b.db != nil {
|
|
if err := b.backfillFromDatabase(sub, cursor); err != nil {
|
|
slog.Warn("Database backfill failed, falling back to in-memory", "error", err)
|
|
b.backfillFromMemory(sub, cursor)
|
|
}
|
|
return
|
|
}
|
|
|
|
// Fall back to in-memory backfill
|
|
b.backfillFromMemory(sub, cursor)
|
|
}
|
|
|
|
// backfillFromDatabase queries events from database and sends to subscriber
|
|
func (b *EventBroadcaster) backfillFromDatabase(sub *Subscriber, cursor int64) error {
|
|
// Query events where seq > cursor, ordered by seq
|
|
// Include created_at to preserve original event timestamp
|
|
query := `
|
|
SELECT seq, commit_cid, rev, since_rev, repo_slice, ops_json, created_at
|
|
FROM firehose_events
|
|
WHERE seq > ?
|
|
ORDER BY seq ASC
|
|
`
|
|
|
|
rows, err := b.db.Query(query, cursor)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer rows.Close()
|
|
|
|
for rows.Next() {
|
|
var (
|
|
seq int64
|
|
commitCID string
|
|
rev string
|
|
sinceRev sql.NullString
|
|
repoSlice []byte
|
|
opsJSON []byte
|
|
createdAt time.Time
|
|
)
|
|
|
|
if err := rows.Scan(&seq, &commitCID, &rev, &sinceRev, &repoSlice, &opsJSON, &createdAt); err != nil {
|
|
slog.Error("Error scanning event row", "error", err)
|
|
continue
|
|
}
|
|
|
|
// Deserialize ops from JSON
|
|
var ops []*atproto.SyncSubscribeRepos_RepoOp
|
|
if err := json.Unmarshal(opsJSON, &ops); err != nil {
|
|
slog.Error("Error unmarshaling ops", "seq", seq, "error", err)
|
|
continue
|
|
}
|
|
|
|
// Reconstruct event
|
|
var since *string
|
|
if sinceRev.Valid {
|
|
since = &sinceRev.String
|
|
}
|
|
|
|
event := &RepoCommitEvent{
|
|
Seq: seq,
|
|
Repo: b.holdDID,
|
|
Commit: commitCID,
|
|
Rev: rev,
|
|
Since: since,
|
|
Blocks: repoSlice,
|
|
Ops: ops,
|
|
Time: createdAt.Format(time.RFC3339), // Use original event time from database
|
|
Type: "#commit",
|
|
}
|
|
|
|
// Send to subscriber
|
|
select {
|
|
case sub.send <- event:
|
|
// Sent successfully
|
|
case <-time.After(5 * time.Second):
|
|
// Timeout, subscriber too slow
|
|
slog.Warn("Backfill timeout for subscriber", "seq", seq)
|
|
return nil
|
|
}
|
|
}
|
|
|
|
return rows.Err()
|
|
}
|
|
|
|
// backfillFromMemory sends events from in-memory ring buffer (fallback)
|
|
func (b *EventBroadcaster) backfillFromMemory(sub *Subscriber, cursor int64) {
|
|
b.mu.RLock()
|
|
defer b.mu.RUnlock()
|
|
|
|
for _, he := range b.eventHistory {
|
|
if he.Seq > cursor {
|
|
select {
|
|
case sub.send <- he.Event:
|
|
// Sent
|
|
case <-time.After(5 * time.Second):
|
|
// Timeout, subscriber too slow
|
|
slog.Warn("Backfill timeout for subscriber", "seq", he.Seq)
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// handleSubscriber handles sending events to a subscriber over WebSocket
|
|
func (b *EventBroadcaster) handleSubscriber(sub *Subscriber) {
|
|
defer func() {
|
|
b.Unsubscribe(sub)
|
|
sub.conn.Close()
|
|
}()
|
|
|
|
for event := range sub.send {
|
|
// Create event header (ATProto firehose format)
|
|
header := events.EventHeader{
|
|
Op: events.EvtKindMessage,
|
|
MsgType: "#commit",
|
|
}
|
|
|
|
// Get a writer for this message
|
|
wc, err := sub.conn.NextWriter(websocket.BinaryMessage)
|
|
if err != nil {
|
|
slog.Error("Failed to get websocket writer", "error", err)
|
|
return
|
|
}
|
|
|
|
// Write header as CBOR
|
|
if err := header.MarshalCBOR(wc); err != nil {
|
|
slog.Error("Failed to write event header", "error", err)
|
|
wc.Close()
|
|
return
|
|
}
|
|
|
|
// Convert our RepoCommitEvent to indigo's SyncSubscribeRepos_Commit
|
|
indigoEvent := convertToIndigoCommit(event)
|
|
|
|
// Write the event as CBOR
|
|
var obj lexutil.CBOR = indigoEvent
|
|
if err := obj.MarshalCBOR(wc); err != nil {
|
|
slog.Error("Failed to write event body", "error", err)
|
|
wc.Close()
|
|
return
|
|
}
|
|
|
|
// Close the writer to flush the message
|
|
if err := wc.Close(); err != nil {
|
|
slog.Error("Failed to close websocket writer", "error", err)
|
|
return
|
|
}
|
|
|
|
// Update cursor
|
|
sub.cursor = event.Seq
|
|
}
|
|
}
|
|
|
|
// convertToIndigoCommit converts our RepoCommitEvent to indigo's SyncSubscribeRepos_Commit
|
|
// which has proper CBOR marshaling methods generated
|
|
func convertToIndigoCommit(event *RepoCommitEvent) *atproto.SyncSubscribeRepos_Commit {
|
|
// Parse commit CID string to cid.Cid, then convert to LexLink
|
|
commitCID, err := cid.Decode(event.Commit)
|
|
if err != nil {
|
|
slog.Warn("Failed to parse commit CID",
|
|
"cid", event.Commit,
|
|
"error", err)
|
|
// Create an empty CID as fallback
|
|
commitCID = cid.Undef
|
|
}
|
|
|
|
// Convert cid.Cid to LexLink
|
|
commitLink := lexutil.LexLink(commitCID)
|
|
|
|
// Convert blocks to LexBytes
|
|
blocks := lexutil.LexBytes(event.Blocks)
|
|
|
|
return &atproto.SyncSubscribeRepos_Commit{
|
|
Seq: event.Seq,
|
|
Repo: event.Repo,
|
|
Commit: commitLink,
|
|
Rev: event.Rev,
|
|
Since: event.Since,
|
|
Blocks: blocks,
|
|
Ops: event.Ops,
|
|
Time: event.Time,
|
|
Blobs: []lexutil.LexLink{}, // Empty for now, we don't track blob refs in our simplified model
|
|
Rebase: false, // DEPRECATED field
|
|
TooBig: false, // Not implementing tooBig for now
|
|
}
|
|
}
|
|
|
|
// encodeCBOR encodes an event as CBOR (DEPRECATED - kept for tests)
|
|
func encodeCBOR(event *RepoCommitEvent) ([]byte, error) {
|
|
// For backward compatibility with tests, encode as JSON
|
|
// Production code uses convertToIndigoCommit + CBOR marshaling in handleSubscriber
|
|
return json.Marshal(event)
|
|
}
|
|
|
|
// SetRepoEventHandler creates a callback to be registered with RepoManager
|
|
func (b *EventBroadcaster) SetRepoEventHandler() func(context.Context, *RepoEvent) {
|
|
return func(ctx context.Context, event *RepoEvent) {
|
|
// Broadcast the event to all subscribers
|
|
// The holdDID is already set in the broadcaster
|
|
b.Broadcast(ctx, event)
|
|
}
|
|
}
|
|
|
|
// GetCurrentSeq returns the current event sequence number
|
|
func (b *EventBroadcaster) GetCurrentSeq() int64 {
|
|
b.mu.RLock()
|
|
defer b.mu.RUnlock()
|
|
return b.eventSeq
|
|
}
|