From a24fc11caf0762c1dd4edbbb782c131dc66fffc4 Mon Sep 17 00:00:00 2001 From: Pasha Sviderski Date: Wed, 2 Sep 2026 12:52:41 +1000 Subject: [PATCH] feat(caddystorage): generic namespace-scoped key-value storage using cluster table in Corrosion --- internal/machine/store/keyspace.go | 167 +++++++++++++++++++++++++++++ internal/machine/store/store.go | 19 +++- 2 files changed, 183 insertions(+), 3 deletions(-) create mode 100644 internal/machine/store/keyspace.go diff --git a/internal/machine/store/keyspace.go b/internal/machine/store/keyspace.go new file mode 100644 index 00000000..5adaf169 --- /dev/null +++ b/internal/machine/store/keyspace.go @@ -0,0 +1,167 @@ +package store + +import ( + "context" + "fmt" + "regexp" + "strings" + "time" + "unicode/utf8" + + "github.com/psviderski/uncloud/internal/corrosion" +) + +const namespaceSeparator = ":" + +var namespacePattern = regexp.MustCompile(`^[a-z][a-z0-9]*(?:[_-][a-z0-9]+)*$`) + +// Record is a key-value entry and its metadata. +type Record struct { + Key string + Value []byte + UpdatedAt time.Time +} + +// KeyspaceListOptions controls the records returned by [Keyspace.List]. +type KeyspaceListOptions struct { + // KeysOnly omits record values from the query when true. + KeysOnly bool +} + +// KeyspaceDeleteOptions controls which records [Keyspace.Delete] deletes. +type KeyspaceDeleteOptions struct { + // Prefix deletes the key and every key that starts with it. + Prefix bool +} + +// Keyspace provides namespace-scoped key-value storage. +// Keys are relative to the namespace. Prefix operations use simple string-prefix matching. +type Keyspace struct { + corro *corrosion.APIClient + namespace string +} + +// Keyspace returns key-value storage scoped to namespace. +func (s *Store) Keyspace(namespace string) (*Keyspace, error) { + if !namespacePattern.MatchString(namespace) { + return nil, fmt.Errorf("invalid namespace %q", namespace) + } + return &Keyspace{ + corro: s.corro, + namespace: namespace, + }, nil +} + +// Get returns the entry for key. It returns [ErrKeyNotFound] when the key does not exist. +func (k *Keyspace) Get(ctx context.Context, key string) (Record, error) { + if key == "" { + return Record{}, fmt.Errorf("key is empty") + } + + rows, err := k.corro.QueryContext(ctx, "SELECT value, updated_at FROM cluster WHERE key = ?", k.namespacedKey(key)) + if err != nil { + return Record{}, err + } + defer rows.Close() + + if !rows.Next() { + if err = rows.Err(); err != nil { + return Record{}, err + } + return Record{}, ErrKeyNotFound + } + + record := Record{Key: key} + var updatedAt string + if err = rows.Scan(&record.Value, &updatedAt); err != nil { + return Record{}, err + } + if record.UpdatedAt, err = parseTimestamp(updatedAt); err != nil { + return Record{}, fmt.Errorf("parse updated_at for key %q: %w", key, err) + } + return record, nil +} + +// Put stores value at key. +func (k *Keyspace) Put(ctx context.Context, key string, value []byte) error { + if key == "" { + return fmt.Errorf("key is empty") + } + + _, err := k.corro.ExecContext(ctx, + "INSERT OR REPLACE INTO cluster (key, value, updated_at) VALUES (?, ?, datetime('now', 'subsec'))", + k.namespacedKey(key), value) + return err +} + +// Delete deletes key. When Prefix is true, it also deletes every key that starts with key. +func (k *Keyspace) Delete(ctx context.Context, key string, opts KeyspaceDeleteOptions) error { + if key == "" { + return fmt.Errorf("key is empty") + } + + namespacedKey := k.namespacedKey(key) + if !opts.Prefix { + _, err := k.corro.ExecContext(ctx, "DELETE FROM cluster WHERE key = ?", namespacedKey) + return err + } + + _, err := k.corro.ExecContext(ctx, `DELETE FROM cluster WHERE key >= ? AND key < ?`, + namespacedKey, prefixRangeEnd(namespacedKey)) + return err +} + +// List returns entries whose keys start with prefix, ordered by key. An empty prefix lists every entry in the keyspace. +func (k *Keyspace) List(ctx context.Context, prefix string, opts KeyspaceListOptions) ([]Record, error) { + keyspacePrefix := k.namespacedKey("") + namespacedPrefix := k.namespacedKey(prefix) + query := "SELECT key, value, updated_at FROM cluster WHERE key >= ? AND key < ? ORDER BY key" + if opts.KeysOnly { + query = "SELECT key, updated_at FROM cluster WHERE key >= ? AND key < ? ORDER BY key" + } + + rows, err := k.corro.QueryContext(ctx, query, namespacedPrefix, prefixRangeEnd(namespacedPrefix)) + if err != nil { + return nil, err + } + defer rows.Close() + + var records []Record + for rows.Next() { + var ( + record Record + updatedAt string + ) + if opts.KeysOnly { + err = rows.Scan(&record.Key, &updatedAt) + } else { + err = rows.Scan(&record.Key, &record.Value, &updatedAt) + } + if err != nil { + return nil, err + } + if record.UpdatedAt, err = parseTimestamp(updatedAt); err != nil { + return nil, fmt.Errorf("parse updated_at for key %q: %w", record.Key, err) + } + record.Key = strings.TrimPrefix(record.Key, keyspacePrefix) + records = append(records, record) + } + if err = rows.Err(); err != nil { + return nil, err + } + + return records, nil +} + +func (k *Keyspace) namespacedKey(key string) string { + return k.namespace + namespaceSeparator + key +} + +func parseTimestamp(value string) (time.Time, error) { + return time.Parse(time.DateTime, value) +} + +// prefixRangeEnd returns an exclusive upper bound for a key prefix. +func prefixRangeEnd(prefix string) string { + return prefix + string(utf8.MaxRune) +} diff --git a/internal/machine/store/store.go b/internal/machine/store/store.go index 3cd72474..d1d78868 100644 --- a/internal/machine/store/store.go +++ b/internal/machine/store/store.go @@ -31,14 +31,19 @@ func New(corro *corrosion.APIClient) *Store { return &Store{corro: corro} } +// Get retrieves an unnamespaced legacy value. +// +// Deprecated: Existing callers may continue to use Get for legacy records. New +// code should use [Store.Keyspace]. func (s *Store) Get(ctx context.Context, key string, value any) error { rows, err := s.corro.QueryContext(ctx, "SELECT value FROM cluster WHERE key = ?", key) if err != nil { return err } + defer rows.Close() if !rows.Next() { - if rows.Err() != nil { - return rows.Err() + if err = rows.Err(); err != nil { + return err } return ErrKeyNotFound } @@ -48,13 +53,21 @@ func (s *Store) Get(ctx context.Context, key string, value any) error { return nil } +// Put stores an unnamespaced legacy value. +// +// Deprecated: Existing callers may continue to use Put for legacy records. New +// code should use [Store.Keyspace]. func (s *Store) Put(ctx context.Context, key string, value any) error { _, err := s.corro.ExecContext(ctx, - "INSERT OR REPLACE INTO cluster (key, value, updated_at) VALUES (?, ?, datetime('now'))", + "INSERT OR REPLACE INTO cluster (key, value, updated_at) VALUES (?, ?, datetime('now', 'subsec'))", key, value) return err } +// Delete deletes an unnamespaced legacy value. +// +// Deprecated: Existing callers may continue to use Delete for legacy records. +// New code should use [Store.Keyspace]. func (s *Store) Delete(ctx context.Context, key string) error { _, err := s.corro.ExecContext(ctx, "DELETE FROM cluster WHERE key = ?", key) return err