diff --git a/appview/state/knotstream.go b/appview/state/knotstream.go index 1797f12e..7472d0f4 100644 --- a/appview/state/knotstream.go +++ b/appview/state/knotstream.go @@ -12,6 +12,7 @@ import ( "tangled.sh/tangled.sh/core/appview/config" "tangled.sh/tangled.sh/core/appview/db" kc "tangled.sh/tangled.sh/core/knotclient" + "tangled.sh/tangled.sh/core/knotclient/cursor" "tangled.sh/tangled.sh/core/log" "tangled.sh/tangled.sh/core/rbac" @@ -32,7 +33,7 @@ func KnotstreamConsumer(ctx context.Context, c *config.Config, d *db.DB, enforce logger := log.New("knotstream") cache := cache.New(c.Redis.Addr) - cursorStore := kc.NewRedisCursorStore(cache) + cursorStore := cursor.NewRedisCursorStore(cache) cfg := kc.ConsumerConfig{ Sources: srcs, diff --git a/knotclient/cursor/memory.go b/knotclient/cursor/memory.go new file mode 100644 index 00000000..1acc49a9 --- /dev/null +++ b/knotclient/cursor/memory.go @@ -0,0 +1,23 @@ +package cursor + +import ( + "sync" +) + +type MemoryStore struct { + store sync.Map +} + +func (m *MemoryStore) Set(knot string, cursor int64) { + m.store.Store(knot, cursor) +} + +func (m *MemoryStore) Get(knot string) (cursor int64) { + if result, ok := m.store.Load(knot); ok { + if val, ok := result.(int64); ok { + return val + } + } + + return 0 +} diff --git a/knotclient/cursor/redis.go b/knotclient/cursor/redis.go new file mode 100644 index 00000000..88da382c --- /dev/null +++ b/knotclient/cursor/redis.go @@ -0,0 +1,43 @@ +package cursor + +import ( + "context" + "fmt" + "strconv" + + "tangled.sh/tangled.sh/core/appview/cache" +) + +const ( + cursorKey = "cursor:%s" +) + +type RedisStore struct { + rdb *cache.Cache +} + +func NewRedisCursorStore(cache *cache.Cache) RedisStore { + return RedisStore{ + rdb: cache, + } +} + +func (r *RedisStore) Set(knot string, cursor int64) { + key := fmt.Sprintf(cursorKey, knot) + r.rdb.Set(context.Background(), key, cursor, 0) +} + +func (r *RedisStore) Get(knot string) (cursor int64) { + key := fmt.Sprintf(cursorKey, knot) + val, err := r.rdb.Get(context.Background(), key).Result() + if err != nil { + return 0 + } + cursor, err = strconv.ParseInt(val, 10, 64) + if err != nil { + // TODO: log here + return 0 + } + + return cursor +} diff --git a/knotclient/cursor/store.go b/knotclient/cursor/store.go new file mode 100644 index 00000000..5b0938ef --- /dev/null +++ b/knotclient/cursor/store.go @@ -0,0 +1,6 @@ +package cursor + +type Store interface { + Set(knot string, cursor int64) + Get(knot string) (cursor int64) +} diff --git a/knotclient/events.go b/knotclient/events.go index ea2b593f..c13a65d5 100644 --- a/knotclient/events.go +++ b/knotclient/events.go @@ -7,11 +7,10 @@ import ( "log/slog" "math/rand" "net/url" - "strconv" "sync" "time" - "tangled.sh/tangled.sh/core/appview/cache" + "tangled.sh/tangled.sh/core/knotclient/cursor" "tangled.sh/tangled.sh/core/log" "github.com/gorilla/websocket" @@ -36,7 +35,7 @@ type ConsumerConfig struct { QueueSize int Logger *slog.Logger Dev bool - CursorStore CursorStore + CursorStore cursor.Store } func NewConsumerConfig() *ConsumerConfig { @@ -72,63 +71,6 @@ type EventConsumer struct { cfg ConsumerConfig } -type CursorStore interface { - Set(knot string, cursor int64) - Get(knot string) (cursor int64) -} - -type RedisCursorStore struct { - rdb *cache.Cache -} - -func NewRedisCursorStore(cache *cache.Cache) RedisCursorStore { - return RedisCursorStore{ - rdb: cache, - } -} - -const ( - cursorKey = "cursor:%s" -) - -func (r *RedisCursorStore) Set(knot string, cursor int64) { - key := fmt.Sprintf(cursorKey, knot) - r.rdb.Set(context.Background(), key, cursor, 0) -} - -func (r *RedisCursorStore) Get(knot string) (cursor int64) { - key := fmt.Sprintf(cursorKey, knot) - val, err := r.rdb.Get(context.Background(), key).Result() - if err != nil { - return 0 - } - - cursor, err = strconv.ParseInt(val, 10, 64) - if err != nil { - return 0 // optionally log parsing error - } - - return cursor -} - -type MemoryCursorStore struct { - store sync.Map -} - -func (m *MemoryCursorStore) Set(knot string, cursor int64) { - m.store.Store(knot, cursor) -} - -func (m *MemoryCursorStore) Get(knot string) (cursor int64) { - if result, ok := m.store.Load(knot); ok { - if val, ok := result.(int64); ok { - return val - } - } - - return 0 -} - func (e *EventConsumer) buildUrl(s EventSource, cursor int64) (*url.URL, error) { scheme := "wss" if e.cfg.Dev { @@ -173,7 +115,7 @@ func NewEventConsumer(cfg ConsumerConfig) *EventConsumer { cfg.QueueSize = 100 } if cfg.CursorStore == nil { - cfg.CursorStore = &MemoryCursorStore{} + cfg.CursorStore = &cursor.MemoryStore{} } return &EventConsumer{ cfg: cfg, diff --git a/nix/vm.nix b/nix/vm.nix index a808ae2d..f2512ba2 100644 --- a/nix/vm.nix +++ b/nix/vm.nix @@ -21,7 +21,7 @@ nixpkgs.lib.nixosSystem { g = config.services.tangled-knot.gitUser; in [ "d /var/lib/knot 0770 ${u} ${g} - -" # Create the directory first - "f+ /var/lib/knot/secret 0660 ${u} ${g} - KNOT_SERVER_SECRET=16154910ef55fe48121082c0b51fc0e360a8b15eb7bda7991d88dc9f7684427a" + "f+ /var/lib/knot/secret 0660 ${u} ${g} - KNOT_SERVER_SECRET=2650ecafdce279b09865fb1923051156eb773ee7485061b2e766086f07dbd85a" ]; services.tangled-knot = { enable = true;