From 3f0650f6ce8ad85b0eaeca29f9d41bb9eae80198 Mon Sep 17 00:00:00 2001 From: Akshay Date: Mon, 24 Mar 2025 23:52:11 +0000 Subject: [PATCH] add in-memory jetstream did filter --- appview/db/jetstream.go | 16 ++- appview/pages/templates/knots.html | 3 +- appview/state/jetstream.go | 2 +- appview/state/state.go | 10 +- cmd/jstest/main.go | 150 ----------------------------- cmd/knotserver/main.go | 2 +- jetstream/jetstream.go | 110 ++++++++++++++++----- knotserver/db/jetstream.go | 16 ++- knotserver/db/pubkeys.go | 4 +- knotserver/handler.go | 4 +- knotserver/jetstream.go | 4 +- knotserver/routes.go | 6 +- 12 files changed, 120 insertions(+), 207 deletions(-) delete mode 100644 cmd/jstest/main.go diff --git a/appview/db/jetstream.go b/appview/db/jetstream.go index f4196387..cdfbcda5 100644 --- a/appview/db/jetstream.go +++ b/appview/db/jetstream.go @@ -5,21 +5,17 @@ type DbWrapper struct { } func (db DbWrapper) SaveLastTimeUs(lastTimeUs int64) error { - _, err := db.Exec(`insert into _jetstream (last_time_us) values (?)`, lastTimeUs) + _, err := db.Exec(` + insert into _jetstream (id, last_time_us) + values (1, ?) + on conflict(id) do update set last_time_us = excluded.last_time_us + `, lastTimeUs) return err } -func (db DbWrapper) UpdateLastTimeUs(lastTimeUs int64) error { - _, err := db.Exec(`update _jetstream set last_time_us = ? where rowid = 1`, lastTimeUs) - if err != nil { - return err - } - return nil -} - func (db DbWrapper) GetLastTimeUs() (int64, error) { var lastTimeUs int64 - row := db.QueryRow(`select last_time_us from _jetstream`) + row := db.QueryRow(`select last_time_us from _jetstream where id = 1;`) err := row.Scan(&lastTimeUs) return lastTimeUs, err } diff --git a/appview/pages/templates/knots.html b/appview/pages/templates/knots.html index 63c28dc4..a01ea011 100644 --- a/appview/pages/templates/knots.html +++ b/appview/pages/templates/knots.html @@ -8,8 +8,7 @@

Generate a key to initialize your knot server.

2*24*60*60*1000*1000 { lastTimeUs = time.Now().UnixMicro() l.Warn("last time us is older than 2 days; discarding that and starting from now") - err = j.db.UpdateLastTimeUs(lastTimeUs) + err = j.db.SaveLastTimeUs(lastTimeUs) if err != nil { l.Error("failed to save last time us", "error", err) } @@ -155,3 +179,41 @@ func (j *JetstreamClient) getLastTimeUs(ctx context.Context) *int64 { l.Info("found last time_us", "time_us", lastTimeUs) return &lastTimeUs } + +func (j *JetstreamClient) saveIfKilled(ctx context.Context) context.Context { + ctxWithCancel, cancel := context.WithCancel(ctx) + + sigChan := make(chan os.Signal, 1) + + signal.Notify(sigChan, + syscall.SIGINT, + syscall.SIGTERM, + syscall.SIGQUIT, + syscall.SIGHUP, + syscall.SIGKILL, + syscall.SIGSTOP, + ) + + go func() { + sig := <-sigChan + j.l.Info("Received signal, initiating graceful shutdown", "signal", sig) + + lastTimeUs := time.Now().UnixMicro() + if err := j.db.SaveLastTimeUs(lastTimeUs); err != nil { + j.l.Error("Failed to save last time during shutdown", "error", err) + } + j.l.Info("Saved lastTimeUs before shutdown", "lastTimeUs", lastTimeUs) + + j.cancelMu.Lock() + if j.cancel != nil { + j.cancel() + } + j.cancelMu.Unlock() + + cancel() + + os.Exit(0) + }() + + return ctxWithCancel +} diff --git a/knotserver/db/jetstream.go b/knotserver/db/jetstream.go index f2c4c9d6..f672fc15 100644 --- a/knotserver/db/jetstream.go +++ b/knotserver/db/jetstream.go @@ -1,21 +1,17 @@ package db func (d *DB) SaveLastTimeUs(lastTimeUs int64) error { - _, err := d.db.Exec(`insert into _jetstream (last_time_us) values (?)`, lastTimeUs) + _, err := d.db.Exec(` + insert into _jetstream (id, last_time_us) + values (1, ?) + on conflict(id) do update set last_time_us = excluded.last_time_us + `, lastTimeUs) return err } -func (d *DB) UpdateLastTimeUs(lastTimeUs int64) error { - _, err := d.db.Exec(`update _jetstream set last_time_us = ? where rowid = 1`, lastTimeUs) - if err != nil { - return err - } - return nil -} - func (d *DB) GetLastTimeUs() (int64, error) { var lastTimeUs int64 - row := d.db.QueryRow(`select last_time_us from _jetstream`) + row := d.db.QueryRow(`select last_time_us from _jetstream where id = 1;`) err := row.Scan(&lastTimeUs) return lastTimeUs, err } diff --git a/knotserver/db/pubkeys.go b/knotserver/db/pubkeys.go index 6c3a1b3d..2e447cd4 100644 --- a/knotserver/db/pubkeys.go +++ b/knotserver/db/pubkeys.go @@ -44,8 +44,8 @@ func (d *DB) RemovePublicKey(did string) error { return err } -func (pk *PublicKey) JSON() map[string]interface{} { - return map[string]interface{}{ +func (pk *PublicKey) JSON() map[string]any { + return map[string]any{ "did": pk.Did, "key": pk.Key, "created": pk.Created, diff --git a/knotserver/handler.go b/knotserver/handler.go index 0d372081..5207f40c 100644 --- a/knotserver/handler.go +++ b/knotserver/handler.go @@ -63,7 +63,9 @@ func Setup(ctx context.Context, c *config.Config, db *db.DB, e *rbac.Enforcer, j if len(dids) > 0 { h.knotInitialized = true close(h.init) - // h.jc.UpdateDids(dids) + for _, d := range dids { + h.jc.AddDid(d) + } } r.Get("/", h.Index) diff --git a/knotserver/jetstream.go b/knotserver/jetstream.go index 264c1a90..76a48cb5 100644 --- a/knotserver/jetstream.go +++ b/knotserver/jetstream.go @@ -53,6 +53,7 @@ func (h *Handle) processKnotMember(ctx context.Context, did string, record tangl l.Error("failed to add did", "error", err) return fmt.Errorf("failed to add did: %w", err) } + h.jc.AddDid(did) if err := h.fetchAndAddKeys(ctx, did); err != nil { return fmt.Errorf("failed to fetch and add keys: %w", err) @@ -115,10 +116,9 @@ func (h *Handle) processMessages(ctx context.Context, event *models.Event) error eventTime := event.TimeUS lastTimeUs := eventTime + 1 fmt.Println("lastTimeUs", lastTimeUs) - if err := h.db.UpdateLastTimeUs(lastTimeUs); err != nil { + if err := h.db.SaveLastTimeUs(lastTimeUs); err != nil { err = fmt.Errorf("(deferred) failed to save last time us: %w", err) } - // h.jc.UpdateDids([]string{did}) }() raw := json.RawMessage(event.Commit.Record) diff --git a/knotserver/routes.go b/knotserver/routes.go index 6ab61c19..a5ec37c8 100644 --- a/knotserver/routes.go +++ b/knotserver/routes.go @@ -448,7 +448,7 @@ func (h *Handle) Keys(w http.ResponseWriter, r *http.Request) { return } - data := make([]map[string]interface{}, 0) + data := make([]map[string]any, 0) for _, key := range keys { j := key.JSON() data = append(data, j) @@ -684,8 +684,8 @@ func (h *Handle) AddMember(w http.ResponseWriter, r *http.Request) { writeError(w, err.Error(), http.StatusInternalServerError) return } - h.jc.AddDid(did) + if err := h.e.AddMember(ThisServer, did); err != nil { l.Error("adding member", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) @@ -768,8 +768,8 @@ func (h *Handle) Init(w http.ResponseWriter, r *http.Request) { writeError(w, err.Error(), http.StatusInternalServerError) return } + h.jc.AddDid(data.Did) - // h.jc.UpdateDids([]string{data.Did}) if err := h.e.AddOwner(ThisServer, data.Did); err != nil { l.Error("adding owner", "error", err.Error()) writeError(w, err.Error(), http.StatusInternalServerError) -- 2.51.2