From 5df480c851c9315016d30c742efd75a79b5b1a89 Mon Sep 17 00:00:00 2001 From: Akshay Date: Sun, 02 Mar 2025 10:18:00 +0000 Subject: [PATCH] defer last event time in appview ingester --- knotserver/middleware.go | 2 -- appview/db/timeline.go | 2 -- appview/state/jetstream.go | 50 ++++++++++++++++++++++++++++++++++++++++++++++++++ appview/state/state.go | 28 +--------------------------- 4 file(s) changed, 51 insertion(s)(+), 31 deletion(s)(-) diff --git a/knotserver/middleware.go b/knotserver/middleware.go --- a/knotserver/middleware.go +++ b/knotserver/middleware.go @@ -4,7 +4,6 @@ "crypto/hmac" "crypto/sha256" "encoding/hex" - "log" "net/http" "time" ) @@ -15,7 +14,6 @@ } return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { signature := r.Header.Get("X-Signature") - log.Println(signature) if signature == "" || !h.verifyHMAC(signature, r) { writeError(w, "signature verification failed", http.StatusForbidden) return diff --git a/appview/db/timeline.go b/appview/db/timeline.go --- a/appview/db/timeline.go +++ b/appview/db/timeline.go @@ -1,7 +1,6 @@ package db import ( - "log" "sort" "time" ) @@ -26,7 +25,6 @@ } for _, repo := range repos { - log.Println(repo.Created) events = append(events, TimelineEvent{ Repo: &repo, Follow: nil, diff --git a/appview/state/jetstream.go b/appview/state/jetstream.go new file mode 100644 --- /dev/null +++ b/appview/state/jetstream.go @@ -0,0 +1,50 @@ +package state + +import ( + "context" + "encoding/json" + "fmt" + "log" + + "github.com/bluesky-social/jetstream/pkg/models" + tangled "github.com/sotangled/tangled/api/tangled" + "github.com/sotangled/tangled/appview/db" +) + +type Ingester func(ctx context.Context, e *models.Event) error + +func jetstreamIngester(db *db.DB) Ingester { + return func(ctx context.Context, e *models.Event) error { + var err error + defer func() { + eventTime := e.TimeUS + lastTimeUs := eventTime + 1 + if err := db.UpdateLastTimeUs(lastTimeUs); err != nil { + err = fmt.Errorf("(deferred) failed to save last time us: %w", err) + } + }() + + if e.Kind != models.EventKindCommit { + return nil + } + + did := e.Did + raw := json.RawMessage(e.Commit.Record) + + switch e.Commit.Collection { + case tangled.GraphFollowNSID: + record := tangled.GraphFollow{} + err := json.Unmarshal(raw, &record) + if err != nil { + log.Println("invalid record") + return err + } + err = db.AddFollow(did, record.Subject, e.Commit.RKey) + if err != nil { + return fmt.Errorf("failed to add follow to db: %w", err) + } + } + + return err + } +} diff --git a/appview/state/state.go b/appview/state/state.go --- a/appview/state/state.go +++ b/appview/state/state.go @@ -5,7 +5,6 @@ "crypto/hmac" "crypto/sha256" "encoding/hex" - "encoding/json" "fmt" "log" "log/slog" @@ -16,7 +15,6 @@ comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/atproto/syntax" lexutil "github.com/bluesky-social/indigo/lex/util" - "github.com/bluesky-social/jetstream/pkg/models" securejoin "github.com/cyphar/filepath-securejoin" "github.com/go-chi/chi/v5" tangled "github.com/sotangled/tangled/api/tangled" @@ -64,31 +62,7 @@ if err != nil { return nil, fmt.Errorf("failed to create jetstream client: %w", err) } - err = jc.StartJetstream(context.Background(), func(ctx context.Context, e *models.Event) error { - if e.Kind != models.EventKindCommit { - return nil - } - - did := e.Did - raw := json.RawMessage(e.Commit.Record) - - switch e.Commit.Collection { - case tangled.GraphFollowNSID: - record := tangled.GraphFollow{} - err := json.Unmarshal(raw, &record) - if err != nil { - log.Println("invalid record") - return err - } - err = db.AddFollow(did, record.Subject, e.Commit.RKey) - if err != nil { - return fmt.Errorf("failed to add follow to db: %w", err) - } - return db.UpdateLastTimeUs(e.TimeUS) - } - - return nil - }) + err = jc.StartJetstream(context.Background(), jetstreamIngester(db)) if err != nil { return nil, fmt.Errorf("failed to start jetstream watcher: %w", err) } -- tangled.sh