From e56634d6c89538afcf0d78f26ed4b5fece68e529 Mon Sep 17 00:00:00 2001 From: Evan Jarrett Date: Mon, 6 Oct 2025 17:39:21 -0500 Subject: [PATCH] more jetstream and ui improvements --- Dockerfile | 4 +- cmd/registry/serve.go | 71 +++- docker-compose.yml | 2 + go.mod | 3 +- go.sum | 4 + pkg/appview/db/queries.go | 23 +- pkg/appview/handlers/home.go | 34 +- pkg/appview/handlers/images.go | 7 +- pkg/appview/handlers/settings.go | 62 ++- pkg/appview/jetstream/worker.go | 383 ++++++++++++++++++ pkg/appview/middleware/auth.go | 10 +- pkg/appview/session/session.go | 105 ++++- pkg/appview/static/js/app.js | 14 + pkg/appview/templates/pages/images.html | 35 +- pkg/appview/templates/partials/push-list.html | 6 +- pkg/auth/oauth/refresher.go | 25 +- pkg/auth/oauth/server.go | 6 +- 17 files changed, 696 insertions(+), 98 deletions(-) create mode 100644 pkg/appview/jetstream/worker.go diff --git a/Dockerfile b/Dockerfile index e0d9001..6550ae0 100644 --- a/Dockerfile +++ b/Dockerfile @@ -22,8 +22,8 @@ RUN CGO_ENABLED=1 GOOS=linux go build -a -o atcr-registry ./cmd/registry # Runtime stage FROM alpine:latest -# Install CA certificates for HTTPS and SQLite runtime libraries -RUN apk --no-cache add ca-certificates sqlite-libs +# Install CA certificates for HTTPS, SQLite runtime libraries, and sqlite CLI for debugging +RUN apk --no-cache add ca-certificates sqlite-libs sqlite # Set working directory WORKDIR /app diff --git a/cmd/registry/serve.go b/cmd/registry/serve.go index eeb93e2..c36b104 100644 --- a/cmd/registry/serve.go +++ b/cmd/registry/serve.go @@ -27,6 +27,7 @@ import ( "atcr.io/pkg/appview" "atcr.io/pkg/appview/db" uihandlers "atcr.io/pkg/appview/handlers" + "atcr.io/pkg/appview/jetstream" appmiddleware "atcr.io/pkg/appview/middleware" appsession "atcr.io/pkg/appview/session" "github.com/gorilla/mux" @@ -127,7 +128,7 @@ func serveRegistry(cmd *cobra.Command, args []string) error { middleware.SetGlobalRefresher(refresher) // 6. Initialize UI components (get session store for OAuth integration) - uiDatabase, uiSessionStore, uiTemplates, uiRouter := initializeUI(config) + uiDatabase, uiSessionStore, uiTemplates, uiRouter := initializeUI(config, refresher, baseURL) // 7. Create OAuth server oauthServer := oauth.NewServer(refreshStorage, sessionManager, baseURL) @@ -338,7 +339,7 @@ func extractDefaultHoldEndpoint(config *configuration.Configuration) string { } // initializeUI initializes the web UI components -func initializeUI(config *configuration.Configuration) (*sql.DB, *appsession.Store, *template.Template, *mux.Router) { +func initializeUI(config *configuration.Configuration, refresher *oauth.Refresher, baseURL string) (*sql.DB, *appsession.Store, *template.Template, *mux.Router) { // Check if UI is enabled (optional configuration) uiEnabled := os.Getenv("ATCR_UI_ENABLED") if uiEnabled == "false" { @@ -367,8 +368,12 @@ func initializeUI(config *configuration.Configuration) (*sql.DB, *appsession.Sto fmt.Printf("UI database initialized at %s\n", dbPath) - // Create session store - sessionStore := appsession.NewStore() + // Create session store with file persistence + sessionStorePath := os.Getenv("ATCR_UI_SESSION_PATH") + if sessionStorePath == "" { + sessionStorePath = "/var/lib/atcr/ui-sessions.json" + } + sessionStore := appsession.NewStore(sessionStorePath) // Start cleanup goroutine go func() { @@ -399,15 +404,17 @@ func initializeUI(config *configuration.Configuration) (*sql.DB, *appsession.Sto // Public routes (with optional auth for navbar) router.Handle("/", appmiddleware.OptionalAuth(sessionStore)( &uihandlers.HomeHandler{ - DB: database, - Templates: templates, + DB: database, + Templates: templates, + RegistryURL: baseURL, }, )).Methods("GET") router.Handle("/api/recent-pushes", appmiddleware.OptionalAuth(sessionStore)( &uihandlers.RecentPushesHandler{ - DB: database, - Templates: templates, + DB: database, + Templates: templates, + RegistryURL: baseURL, }, )).Methods("GET") @@ -416,15 +423,19 @@ func initializeUI(config *configuration.Configuration) (*sql.DB, *appsession.Sto authRouter.Use(appmiddleware.RequireAuth(sessionStore)) authRouter.Handle("/images", &uihandlers.ImagesHandler{ - DB: database, - Templates: templates, + DB: database, + Templates: templates, + RegistryURL: baseURL, }).Methods("GET") authRouter.Handle("/settings", &uihandlers.SettingsHandler{ Templates: templates, + Refresher: refresher, }).Methods("GET") - authRouter.Handle("/api/profile/default-hold", &uihandlers.UpdateDefaultHoldHandler{}).Methods("POST") + authRouter.Handle("/api/profile/default-hold", &uihandlers.UpdateDefaultHoldHandler{ + Refresher: refresher, + }).Methods("POST") authRouter.Handle("/api/images/{repository}/tags/{tag}", &uihandlers.DeleteTagHandler{ DB: database, @@ -443,5 +454,43 @@ func initializeUI(config *configuration.Configuration) (*sql.DB, *appsession.Sto http.Redirect(w, r, "/", http.StatusFound) }).Methods("POST") + // Start Jetstream worker + jetstreamURL := os.Getenv("JETSTREAM_URL") + if jetstreamURL == "" { + jetstreamURL = "wss://jetstream2.us-west.bsky.network/subscribe" + } + + // Parse cursor for backfilling historical data + // Set to Unix microseconds timestamp to replay from that point + // Examples: + // - 2 weeks ago: use `date -d '2 weeks ago' +%s` * 1000000 + // - Leave unset (or 0) to start from now + var startCursor int64 + if cursorStr := os.Getenv("JETSTREAM_START_CURSOR"); cursorStr != "" { + if cursor, err := time.Parse(time.RFC3339, cursorStr); err == nil { + // Support RFC3339 format: "2025-09-23T00:00:00Z" + startCursor = cursor.UnixMicro() + fmt.Printf("Jetstream: Starting from %s (%d microseconds)\n", cursorStr, startCursor) + } else if cursor, err := time.ParseDuration(cursorStr); err == nil { + // Support duration format: "-336h" (2 weeks ago) + startCursor = time.Now().Add(cursor).UnixMicro() + fmt.Printf("Jetstream: Starting from %s ago (%d microseconds)\n", cursorStr, startCursor) + } else { + fmt.Printf("Warning: Invalid JETSTREAM_START_CURSOR format: %s\n", cursorStr) + } + } + + worker := jetstream.NewWorker(database, jetstreamURL, startCursor) + go func() { + for { + if err := worker.Start(context.Background()); err != nil { + fmt.Printf("Jetstream worker error: %v, reconnecting in 10s...\n", err) + time.Sleep(10 * time.Second) + } + } + }() + + fmt.Println("Jetstream worker started") + return database, sessionStore, templates, router } diff --git a/docker-compose.yml b/docker-compose.yml index 63b653f..51eb54c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -10,6 +10,8 @@ services: environment: - ATCR_TOKEN_STORAGE_PATH=/var/lib/atcr/tokens/oauth-tokens.json - ATCR_UI_ENABLED=true + # Jetstream backfill: Replay 5 days of historical events + # - JETSTREAM_START_CURSOR=-120h volumes: # Auth keys (JWT signing keys) - atcr-auth:/var/lib/atcr/auth diff --git a/go.mod b/go.mod index 42da30f..e03e8b6 100644 --- a/go.mod +++ b/go.mod @@ -30,12 +30,13 @@ require ( github.com/go-logr/stdr v1.2.2 // indirect github.com/gorilla/handlers v1.5.2 // indirect github.com/gorilla/mux v1.8.1 // indirect + github.com/gorilla/websocket v1.5.3 // indirect github.com/grpc-ecosystem/grpc-gateway/v2 v2.23.0 // indirect github.com/hashicorp/golang-lru/arc/v2 v2.0.6 // indirect github.com/hashicorp/golang-lru/v2 v2.0.7 // indirect github.com/inconshreveable/mousetrap v1.1.0 // indirect github.com/jmespath/go-jmespath v0.4.0 // indirect - github.com/klauspost/compress v1.17.11 // indirect + github.com/klauspost/compress v1.18.0 // indirect github.com/mattn/go-sqlite3 v1.14.32 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/opencontainers/image-spec v1.1.0 // indirect diff --git a/go.sum b/go.sum index ad60e80..3116e22 100644 --- a/go.sum +++ b/go.sum @@ -77,6 +77,8 @@ github.com/gorilla/handlers v1.5.2 h1:cLTUSsNkgcwhgRqvCNmdbRWG0A3N4F+M2nWKdScwyE github.com/gorilla/handlers v1.5.2/go.mod h1:dX+xVpaxdSw+q0Qek8SSsl3dfMk3jNddUkMzo0GtH0w= github.com/gorilla/mux v1.8.1 h1:TuBL49tXwgrFYWhqrNgrUNEY92u81SPhu7sTdzQEiWY= github.com/gorilla/mux v1.8.1/go.mod h1:AKf9I4AEqPTmMytcMc0KkNouC66V3BtZ4qD5fmWSiMQ= +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/grpc-ecosystem/grpc-gateway/v2 v2.23.0 h1:ad0vkEBuk23VJzZR9nkLVG0YAoN9coASF1GusYX6AlU= github.com/grpc-ecosystem/grpc-gateway/v2 v2.23.0/go.mod h1:igFoXX2ELCW06bol23DWPB5BEWfZISOzSP5K2sbLea0= github.com/hashicorp/golang-lru/arc/v2 v2.0.6 h1:4NU7uP5vSoK6TbaMj3NtY478TTAWLso/vL1gpNrInHg= @@ -94,6 +96,8 @@ github.com/json-iterator/go v1.1.7/go.mod h1:KdQUCv79m/52Kvf8AW2vK1V8akMuk1QjK/u github.com/julienschmidt/httprouter v1.2.0/go.mod h1:SYymIcj16QtmaHHD7aYtjjsJG7VTCxuUUipMqKk8s4w= github.com/klauspost/compress v1.17.11 h1:In6xLpyWOi1+C7tXUUWv2ot1QvBjxevKAaI6IXrJmUc= github.com/klauspost/compress v1.17.11/go.mod h1:pMDklpSncoRMuLFrf1W9Ss9KT+0rH90U12bZKk7uwG0= +github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= +github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= github.com/konsorten/go-windows-terminal-sequences v1.0.1/go.mod h1:T0+1ngSBFLxvqU3pZ+m/2kptfBszLMUkC4ZK/EgS/cQ= github.com/kr/logfmt v0.0.0-20140226030751-b84e30acd515/go.mod h1:+0opPa2QZZtGFBFZlji/RkVcI2GknAs/DXo4wKdlNEc= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= diff --git a/pkg/appview/db/queries.go b/pkg/appview/db/queries.go index ea6d0c6..5f067ed 100644 --- a/pkg/appview/db/queries.go +++ b/pkg/appview/db/queries.go @@ -2,6 +2,7 @@ package db import ( "database/sql" + "time" ) // GetRecentPushes fetches recent pushes with pagination @@ -81,10 +82,30 @@ func GetUserRepositories(db *sql.DB, did string) ([]Repository, error) { var repos []Repository for rows.Next() { var r Repository - if err := rows.Scan(&r.Name, &r.TagCount, &r.ManifestCount, &r.LastPush); err != nil { + var lastPushStr string + if err := rows.Scan(&r.Name, &r.TagCount, &r.ManifestCount, &lastPushStr); err != nil { return nil, err } + // Parse the timestamp string into time.Time + if lastPushStr != "" { + // Try multiple timestamp formats + formats := []string{ + time.RFC3339Nano, // 2006-01-02T15:04:05.999999999Z07:00 + "2006-01-02 15:04:05.999999999-07:00", // SQLite with microseconds and timezone + "2006-01-02 15:04:05.999999999", // SQLite with microseconds + time.RFC3339, // 2006-01-02T15:04:05Z07:00 + "2006-01-02 15:04:05", // SQLite default + } + + for _, format := range formats { + if t, err := time.Parse(format, lastPushStr); err == nil { + r.LastPush = t + break + } + } + } + // Get tags for this repo tagRows, err := db.Query(` SELECT id, tag, digest, created_at diff --git a/pkg/appview/handlers/home.go b/pkg/appview/handlers/home.go index 4ab457d..60ac541 100644 --- a/pkg/appview/handlers/home.go +++ b/pkg/appview/handlers/home.go @@ -12,17 +12,20 @@ import ( // HomeHandler handles the home page type HomeHandler struct { - DB *sql.DB - Templates *template.Template + DB *sql.DB + Templates *template.Template + RegistryURL string } func (h *HomeHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { data := struct { - User *db.User - Query string + User *db.User + Query string + RegistryURL string }{ - User: middleware.GetUser(r), - Query: r.URL.Query().Get("q"), + User: middleware.GetUser(r), + Query: r.URL.Query().Get("q"), + RegistryURL: h.RegistryURL, } if err := h.Templates.ExecuteTemplate(w, "home", data); err != nil { @@ -33,8 +36,9 @@ func (h *HomeHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { // RecentPushesHandler handles the HTMX request for recent pushes type RecentPushesHandler struct { - DB *sql.DB - Templates *template.Template + DB *sql.DB + Templates *template.Template + RegistryURL string } func (h *RecentPushesHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { @@ -57,13 +61,15 @@ func (h *RecentPushesHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) } data := struct { - Pushes []db.Push - HasMore bool - NextOffset int + Pushes []db.Push + HasMore bool + NextOffset int + RegistryURL string }{ - Pushes: pushes, - HasMore: offset+limit < total, - NextOffset: offset + limit, + Pushes: pushes, + HasMore: offset+limit < total, + NextOffset: offset + limit, + RegistryURL: h.RegistryURL, } if err := h.Templates.ExecuteTemplate(w, "push-list.html", data); err != nil { diff --git a/pkg/appview/handlers/images.go b/pkg/appview/handlers/images.go index e6ceda5..00160ea 100644 --- a/pkg/appview/handlers/images.go +++ b/pkg/appview/handlers/images.go @@ -12,8 +12,9 @@ import ( // ImagesHandler handles the images management page type ImagesHandler struct { - DB *sql.DB - Templates *template.Template + DB *sql.DB + Templates *template.Template + RegistryURL string } func (h *ImagesHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { @@ -34,10 +35,12 @@ func (h *ImagesHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { User *db.User Repositories []db.Repository Query string + RegistryURL string }{ User: user, Repositories: repos, Query: r.URL.Query().Get("q"), + RegistryURL: h.RegistryURL, } if err := h.Templates.ExecuteTemplate(w, "images", data); err != nil { diff --git a/pkg/appview/handlers/settings.go b/pkg/appview/handlers/settings.go index ba58ec4..bedbf85 100644 --- a/pkg/appview/handlers/settings.go +++ b/pkg/appview/handlers/settings.go @@ -1,29 +1,50 @@ package handlers import ( + "fmt" "html/template" "net/http" "time" "atcr.io/pkg/appview/db" "atcr.io/pkg/appview/middleware" + "atcr.io/pkg/atproto" + "atcr.io/pkg/auth/oauth" ) // SettingsHandler handles the settings page type SettingsHandler struct { Templates *template.Template - // TODO: Add ATProto client when implementing profile fetching + Refresher *oauth.Refresher } func (h *SettingsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { user := middleware.GetUser(r) if user == nil { - http.Redirect(w, r, "/auth/oauth/login?return_to=/ui/settings", http.StatusFound) + http.Redirect(w, r, "/auth/oauth/login?return_to=/settings", http.StatusFound) return } - // TODO: Fetch actual profile from PDS using ATProto client - // For now, using mock data from session + // Get access token and DPoP transport for the user + accessToken, _, dpopTransport, err := h.Refresher.GetAccessToken(r.Context(), user.DID) + if err != nil { + http.Error(w, "Failed to get access token: "+err.Error(), http.StatusInternalServerError) + return + } + + // Create ATProto client with DPoP transport + client := atproto.NewClientWithDPoP(user.PDSEndpoint, user.DID, accessToken, nil, dpopTransport) + + // Fetch sailor profile + profile, err := atproto.GetProfile(r.Context(), client) + if err != nil { + // Log error but don't fail - profile might not exist yet + fmt.Printf("WARNING [settings]: Failed to fetch profile for %s: %v\n", user.DID, err) + profile = &atproto.SailorProfileRecord{} + } else { + fmt.Printf("DEBUG [settings]: Fetched profile for %s: defaultHold=%s\n", user.DID, profile.DefaultHold) + } + data := struct { User *db.User Profile struct { @@ -43,7 +64,7 @@ func (h *SettingsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { data.Profile.Handle = user.Handle data.Profile.DID = user.DID data.Profile.PDSEndpoint = user.PDSEndpoint - // data.Profile.DefaultHold will be empty for now + data.Profile.DefaultHold = profile.DefaultHold if err := h.Templates.ExecuteTemplate(w, "settings", data); err != nil { http.Error(w, err.Error(), http.StatusInternalServerError) @@ -53,7 +74,7 @@ func (h *SettingsHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { // UpdateDefaultHoldHandler handles updating the default hold type UpdateDefaultHoldHandler struct { - // TODO: Add ATProto client for updating profile + Refresher *oauth.Refresher } func (h *UpdateDefaultHoldHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) { @@ -65,9 +86,32 @@ func (h *UpdateDefaultHoldHandler) ServeHTTP(w http.ResponseWriter, r *http.Requ holdEndpoint := r.FormValue("hold_endpoint") - // TODO: Update profile in PDS via ATProto client - // For now, just return success - _ = holdEndpoint + // Get access token and DPoP transport for the user + accessToken, _, dpopTransport, err := h.Refresher.GetAccessToken(r.Context(), user.DID) + if err != nil { + http.Error(w, "Failed to get access token: "+err.Error(), http.StatusInternalServerError) + return + } + + // Create ATProto client with DPoP transport + client := atproto.NewClientWithDPoP(user.PDSEndpoint, user.DID, accessToken, nil, dpopTransport) + + // Fetch existing profile or create new one + profile, err := atproto.GetProfile(r.Context(), client) + if err != nil || profile == nil { + // Profile doesn't exist, create new one + profile = atproto.NewSailorProfileRecord(holdEndpoint) + } else { + // Update existing profile + profile.DefaultHold = holdEndpoint + profile.UpdatedAt = time.Now() + } + + // Save profile + if err := atproto.UpdateProfile(r.Context(), client, profile); err != nil { + http.Error(w, "Failed to update profile: "+err.Error(), http.StatusInternalServerError) + return + } w.Header().Set("Content-Type", "text/html") w.Write([]byte(`
✓ Default hold updated successfully!
`)) diff --git a/pkg/appview/jetstream/worker.go b/pkg/appview/jetstream/worker.go new file mode 100644 index 0000000..6f9cea0 --- /dev/null +++ b/pkg/appview/jetstream/worker.go @@ -0,0 +1,383 @@ +package jetstream + +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + "net/url" + "strings" + "time" + + "atcr.io/pkg/appview/db" + "atcr.io/pkg/atproto" + "github.com/gorilla/websocket" + "github.com/klauspost/compress/zstd" +) + +// UserCache caches DID -> handle/PDS mappings to avoid repeated lookups +type UserCache struct { + cache map[string]*db.User +} + +// Worker consumes Jetstream events and populates the UI database +type Worker struct { + db *sql.DB + jetstreamURL string + startCursor int64 + wantedCollections []string + debugCollectionCount int + userCache *UserCache + resolver *atproto.Resolver +} + +// NewWorker creates a new Jetstream worker +// startCursor: Unix microseconds timestamp to start from (0 = start from now) +func NewWorker(database *sql.DB, jetstreamURL string, startCursor int64) *Worker { + if jetstreamURL == "" { + jetstreamURL = "wss://jetstream2.us-west.bsky.network/subscribe" + } + + return &Worker{ + db: database, + jetstreamURL: jetstreamURL, + startCursor: startCursor, + wantedCollections: []string{ + atproto.ManifestCollection, // io.atcr.manifest + atproto.TagCollection, // io.atcr.tag + }, + userCache: &UserCache{ + cache: make(map[string]*db.User), + }, + resolver: atproto.NewResolver(), + } +} + +// Start begins consuming Jetstream events +// This is a blocking function that runs until the context is cancelled +func (w *Worker) Start(ctx context.Context) error { + // Build connection URL with filters + u, err := url.Parse(w.jetstreamURL) + if err != nil { + return fmt.Errorf("invalid jetstream URL: %w", err) + } + + q := u.Query() + for _, collection := range w.wantedCollections { + q.Add("wantedCollections", collection) + } + + // Add cursor if specified (for backfilling historical data) + if w.startCursor > 0 { + q.Set("cursor", fmt.Sprintf("%d", w.startCursor)) + fmt.Printf("Starting from cursor: %d (replaying historical events)\n", w.startCursor) + } + + // Disable compression for now to debug + // q.Set("compress", "true") + u.RawQuery = q.Encode() + + fmt.Printf("Connecting to Jetstream: %s\n", u.String()) + + // Connect to Jetstream + conn, _, err := websocket.DefaultDialer.DialContext(ctx, u.String(), nil) + if err != nil { + return fmt.Errorf("failed to connect to jetstream: %w", err) + } + defer conn.Close() + + // Create zstd decoder for decompressing messages + decoder, err := zstd.NewReader(nil) + if err != nil { + return fmt.Errorf("failed to create zstd decoder: %w", err) + } + defer decoder.Close() + + fmt.Println("Connected to Jetstream, listening for events...") + + // Start heartbeat ticker to show Jetstream is alive + heartbeatTicker := time.NewTicker(30 * time.Second) + defer heartbeatTicker.Stop() + + eventCount := 0 + lastHeartbeat := time.Now() + + // Read messages + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-heartbeatTicker.C: + elapsed := time.Since(lastHeartbeat) + fmt.Printf("Jetstream: Alive (processed %d events in last %.0fs)\n", eventCount, elapsed.Seconds()) + eventCount = 0 + lastHeartbeat = time.Now() + default: + _, message, err := conn.ReadMessage() + if err != nil { + return fmt.Errorf("failed to read message: %w", err) + } + + // For now, process uncompressed messages + // TODO: Re-enable compression once debugging is complete + _ = decoder // Keep decoder to avoid unused variable error + + if err := w.processMessage(message); err != nil { + fmt.Printf("ERROR processing message: %v\n", err) + // Continue processing other messages + } else { + eventCount++ + } + } + } +} + +// processMessage processes a single Jetstream event +func (w *Worker) processMessage(message []byte) error { + var event JetstreamEvent + if err := json.Unmarshal(message, &event); err != nil { + return fmt.Errorf("failed to unmarshal event: %w", err) + } + + // Only process commit events + if event.Kind != "commit" { + return nil + } + + commit := event.Commit + if commit == nil { + return nil + } + + // Set DID on commit from parent event + commit.DID = event.DID + + // Debug: log first few collections we see to understand what's coming through + if w.debugCollectionCount < 5 { + fmt.Printf("Jetstream DEBUG: Received collection=%s, did=%s\n", commit.Collection, commit.DID) + w.debugCollectionCount++ + } + + // Process based on collection + switch commit.Collection { + case atproto.ManifestCollection: + fmt.Printf("Jetstream: Processing manifest event: did=%s, operation=%s, rkey=%s\n", + commit.DID, commit.Operation, commit.RKey) + return w.processManifest(commit) + case atproto.TagCollection: + fmt.Printf("Jetstream: Processing tag event: did=%s, operation=%s, rkey=%s\n", + commit.DID, commit.Operation, commit.RKey) + return w.processTag(commit) + default: + // Ignore other collections + return nil + } +} + +// ensureUser resolves and upserts a user by DID +func (w *Worker) ensureUser(ctx context.Context, did string) error { + // Check cache first + if user, ok := w.userCache.cache[did]; ok { + // Update last seen + user.LastSeen = time.Now() + return db.UpsertUser(w.db, user) + } + + // Resolve DID to get handle and PDS endpoint + resolvedDID, pdsEndpoint, err := w.resolver.ResolveIdentity(ctx, did) + if err != nil { + fmt.Printf("WARNING: Failed to resolve DID %s: %v (using DID as handle)\n", did, err) + // Fallback: use DID as handle + resolvedDID = did + pdsEndpoint = "https://bsky.social" // Default PDS endpoint as fallback + } + + // Get handle from DID document + handle, err := w.resolver.ResolveHandleFromDID(ctx, resolvedDID) + if err != nil { + fmt.Printf("WARNING: Failed to get handle for DID %s: %v (using DID as handle)\n", resolvedDID, err) + handle = resolvedDID // Fallback to DID + } + + // Cache the user + user := &db.User{ + DID: resolvedDID, + Handle: handle, + PDSEndpoint: pdsEndpoint, + LastSeen: time.Now(), + } + w.userCache.cache[did] = user + + // Upsert to database + return db.UpsertUser(w.db, user) +} + +// processManifest processes a manifest commit event +func (w *Worker) processManifest(commit *CommitEvent) error { + // Resolve and upsert user with handle/PDS endpoint + if err := w.ensureUser(context.Background(), commit.DID); err != nil { + return fmt.Errorf("failed to ensure user: %w", err) + } + + if commit.Operation == "delete" { + // Delete manifest + repo := extractRepoFromRKey(commit.RKey) + digest := commit.RKey + return db.DeleteManifest(w.db, commit.DID, repo, digest) + } + + // Parse manifest record + var manifestRecord atproto.ManifestRecord + if commit.Record != nil { + recordBytes, err := json.Marshal(commit.Record) + if err != nil { + return fmt.Errorf("failed to marshal record: %w", err) + } + if err := json.Unmarshal(recordBytes, &manifestRecord); err != nil { + return fmt.Errorf("failed to unmarshal manifest: %w", err) + } + } else { + // No record data, can't process + return nil + } + + // Serialize full manifest as JSON for storage + manifestJSON, err := json.Marshal(manifestRecord) + if err != nil { + return fmt.Errorf("failed to marshal manifest: %w", err) + } + + // Insert manifest + manifestID, err := db.InsertManifest(w.db, &db.Manifest{ + DID: commit.DID, + Repository: manifestRecord.Repository, + Digest: manifestRecord.Digest, + MediaType: manifestRecord.MediaType, + SchemaVersion: manifestRecord.SchemaVersion, + ConfigDigest: manifestRecord.Config.Digest, + ConfigSize: manifestRecord.Config.Size, + RawManifest: string(manifestJSON), + HoldEndpoint: manifestRecord.HoldEndpoint, + CreatedAt: manifestRecord.CreatedAt, + }) + if err != nil { + return fmt.Errorf("failed to insert manifest: %w", err) + } + + // Insert layers + for i, layer := range manifestRecord.Layers { + if err := db.InsertLayer(w.db, &db.Layer{ + ManifestID: manifestID, + Digest: layer.Digest, + MediaType: layer.MediaType, + Size: layer.Size, + LayerIndex: i, + }); err != nil { + // Continue on error - layer might already exist + continue + } + } + + return nil +} + +// processTag processes a tag commit event +func (w *Worker) processTag(commit *CommitEvent) error { + // Resolve and upsert user with handle/PDS endpoint + if err := w.ensureUser(context.Background(), commit.DID); err != nil { + return fmt.Errorf("failed to ensure user: %w", err) + } + + if commit.Operation == "delete" { + // Delete tag + parts := strings.Split(commit.RKey, "/") + if len(parts) < 2 { + return fmt.Errorf("invalid tag rkey: %s", commit.RKey) + } + repo := strings.Join(parts[:len(parts)-1], "/") + tag := parts[len(parts)-1] + return db.DeleteTag(w.db, commit.DID, repo, tag) + } + + // Parse tag record + var tagRecord atproto.TagRecord + if commit.Record != nil { + recordBytes, err := json.Marshal(commit.Record) + if err != nil { + return fmt.Errorf("failed to marshal record: %w", err) + } + if err := json.Unmarshal(recordBytes, &tagRecord); err != nil { + return fmt.Errorf("failed to unmarshal tag: %w", err) + } + } else { + return nil + } + + // Insert or update tag + return db.UpsertTag(w.db, &db.Tag{ + DID: commit.DID, + Repository: tagRecord.Repository, + Tag: tagRecord.Tag, + Digest: tagRecord.ManifestDigest, + CreatedAt: tagRecord.UpdatedAt, + }) +} + +// JetstreamEvent represents a Jetstream event +type JetstreamEvent struct { + DID string `json:"did"` + TimeUS int64 `json:"time_us"` + Kind string `json:"kind"` // "commit", "identity", "account" + Commit *CommitEvent `json:"commit,omitempty"` + Identity *IdentityInfo `json:"identity,omitempty"` + Account *AccountInfo `json:"account,omitempty"` +} + +// CommitEvent represents a commit event (create/update/delete) +type CommitEvent struct { + Rev string `json:"rev"` + Operation string `json:"operation"` // "create", "update", "delete" + Collection string `json:"collection"` + RKey string `json:"rkey"` + Record map[string]interface{} `json:"record,omitempty"` + CID string `json:"cid,omitempty"` + DID string `json:"-"` // Set from parent event +} + +// IdentityInfo represents an identity event +type IdentityInfo struct { + DID string `json:"did"` + Handle string `json:"handle"` + Seq int64 `json:"seq"` + Time string `json:"time"` +} + +// AccountInfo represents an account status event +type AccountInfo struct { + Active bool `json:"active"` + DID string `json:"did"` + Seq int64 `json:"seq"` + Time string `json:"time"` + Status string `json:"status,omitempty"` +} + +// Helper functions + +func extractRepoFromRKey(rkey string) string { + // RKey format: or / + // For manifest, it's just the digest + parts := strings.Split(rkey, "/") + if len(parts) > 1 { + return parts[0] + } + return "" +} + +func calculateManifestSize(manifest *atproto.ManifestRecord) int64 { + var total int64 + total += manifest.Config.Size + for _, layer := range manifest.Layers { + total += layer.Size + } + return total +} diff --git a/pkg/appview/middleware/auth.go b/pkg/appview/middleware/auth.go index 5d5198e..f886d2d 100644 --- a/pkg/appview/middleware/auth.go +++ b/pkg/appview/middleware/auth.go @@ -29,8 +29,9 @@ func RequireAuth(store *session.Store) func(http.Handler) http.Handler { } user := &db.User{ - DID: sess.DID, - Handle: sess.Handle, + DID: sess.DID, + Handle: sess.Handle, + PDSEndpoint: sess.PDSEndpoint, } ctx := context.WithValue(r.Context(), userKey, user) @@ -47,8 +48,9 @@ func OptionalAuth(store *session.Store) func(http.Handler) http.Handler { if ok { if sess, ok := store.Get(sessionID); ok { user := &db.User{ - DID: sess.DID, - Handle: sess.Handle, + DID: sess.DID, + Handle: sess.Handle, + PDSEndpoint: sess.PDSEndpoint, } ctx := context.WithValue(r.Context(), userKey, user) r = r.WithContext(ctx) diff --git a/pkg/appview/session/session.go b/pkg/appview/session/session.go index 5d8cfe3..08af7d7 100644 --- a/pkg/appview/session/session.go +++ b/pkg/appview/session/session.go @@ -3,34 +3,92 @@ package session import ( "crypto/rand" "encoding/base64" + "encoding/json" + "fmt" "net/http" + "os" "sync" "time" ) // Session represents a user session type Session struct { - ID string - DID string - Handle string - ExpiresAt time.Time + ID string + DID string + Handle string + PDSEndpoint string + ExpiresAt time.Time } // Store manages user sessions type Store struct { mu sync.RWMutex sessions map[string]*Session + filePath string } -// NewStore creates a new session store -func NewStore() *Store { - return &Store{ +// NewStore creates a new session store with file persistence +func NewStore(filePath string) *Store { + store := &Store{ sessions: make(map[string]*Session), + filePath: filePath, } + + // Load existing sessions from file + if err := store.load(); err != nil { + fmt.Printf("Warning: Failed to load sessions from %s: %v\n", filePath, err) + } + + return store } -// Create creates a new session and returns the full Session struct -func (s *Store) Create(did, handle string, duration time.Duration) (string, error) { +// load reads sessions from disk +func (s *Store) load() error { + if s.filePath == "" { + return nil + } + + data, err := os.ReadFile(s.filePath) + if err != nil { + if os.IsNotExist(err) { + return nil // File doesn't exist yet, that's fine + } + return err + } + + var sessions map[string]*Session + if err := json.Unmarshal(data, &sessions); err != nil { + return err + } + + // Filter out expired sessions + now := time.Now() + for id, sess := range sessions { + if now.Before(sess.ExpiresAt) { + s.sessions[id] = sess + } + } + + fmt.Printf("Loaded %d active sessions from disk\n", len(s.sessions)) + return nil +} + +// save writes sessions to disk +func (s *Store) save() error { + if s.filePath == "" { + return nil + } + + data, err := json.Marshal(s.sessions) + if err != nil { + return err + } + + return os.WriteFile(s.filePath, data, 0600) +} + +// Create creates a new session and returns the session ID +func (s *Store) Create(did, handle, pdsEndpoint string, duration time.Duration) (string, error) { s.mu.Lock() defer s.mu.Unlock() @@ -41,13 +99,20 @@ func (s *Store) Create(did, handle string, duration time.Duration) (string, erro } sess := &Session{ - ID: base64.URLEncoding.EncodeToString(b), - DID: did, - Handle: handle, - ExpiresAt: time.Now().Add(duration), + ID: base64.URLEncoding.EncodeToString(b), + DID: did, + Handle: handle, + PDSEndpoint: pdsEndpoint, + ExpiresAt: time.Now().Add(duration), } s.sessions[sess.ID] = sess + + // Save to disk + if err := s.save(); err != nil { + fmt.Printf("Warning: Failed to save sessions to disk: %v\n", err) + } + return sess.ID, nil } @@ -70,6 +135,11 @@ func (s *Store) Delete(id string) { defer s.mu.Unlock() delete(s.sessions, id) + + // Save to disk + if err := s.save(); err != nil { + fmt.Printf("Warning: Failed to save sessions to disk: %v\n", err) + } } // Cleanup removes expired sessions @@ -78,9 +148,18 @@ func (s *Store) Cleanup() { defer s.mu.Unlock() now := time.Now() + deleted := 0 for id, sess := range s.sessions { if now.After(sess.ExpiresAt) { delete(s.sessions, id) + deleted++ + } + } + + if deleted > 0 { + // Save to disk + if err := s.save(); err != nil { + fmt.Printf("Warning: Failed to save sessions to disk: %v\n", err) } } } diff --git a/pkg/appview/static/js/app.js b/pkg/appview/static/js/app.js index 58fde17..3f3a65b 100644 --- a/pkg/appview/static/js/app.js +++ b/pkg/appview/static/js/app.js @@ -58,3 +58,17 @@ document.addEventListener('htmx:afterSwap', updateTimestamps); // Update timestamps periodically setInterval(updateTimestamps, 60000); // Every minute + +// Toggle repository details (for images page) +function toggleRepo(name) { + const details = document.getElementById('repo-' + name); + const btn = document.getElementById('btn-' + name); + + if (details.style.display === 'none') { + details.style.display = 'block'; + btn.textContent = '▲'; + } else { + details.style.display = 'none'; + btn.textContent = '▼'; + } +} diff --git a/pkg/appview/templates/pages/images.html b/pkg/appview/templates/pages/images.html index 0faeb1f..3315427 100644 --- a/pkg/appview/templates/pages/images.html +++ b/pkg/appview/templates/pages/images.html @@ -7,6 +7,7 @@ Your Images - ATCR + {{ template "nav" . }} @@ -17,10 +18,11 @@ {{ if .Repositories }} {{ range .Repositories }} + {{ $repoName := .Name }}
-
+
-

{{ .Name }}

+

{{ $repoName }}

{{ .TagCount }} tags • @@ -31,16 +33,16 @@
- +
-