From 02ee33db807f932a1533740b0b7edf8f0d7a8ed8 Mon Sep 17 00:00:00 2001 From: Anirudh Oppiliappan Date: Tue, 25 Aug 2026 15:44:02 +0300 Subject: [PATCH] deliberi/kv: add email-did namespace cache client --- deliberi/kv/kv.go | 103 ++++++++++++++++++++++ deliberi/kv/kv_test.go | 190 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 293 insertions(+) create mode 100644 deliberi/kv/kv.go create mode 100644 deliberi/kv/kv_test.go diff --git a/deliberi/kv/kv.go b/deliberi/kv/kv.go new file mode 100644 index 000000000..eaaf92dbb --- /dev/null +++ b/deliberi/kv/kv.go @@ -0,0 +1,103 @@ +// Package kv keeps the EMAIL_DID namespace in sync with deliberi's email mutations, fire-and-forget. +package kv + +import ( + "fmt" + "io" + "log/slog" + "net/http" + "net/url" + "strings" + "time" +) + +type EmailKV interface { + PutEmailDid(email, did string) + SetPrimaryEmail(did, email string) + DeleteEmailKey(email string) +} + +type Client struct { + apiToken string + baseURL string + client *http.Client + logger *slog.Logger +} + +func NewClient(token, accountID, namespaceID string, logger *slog.Logger) *Client { + if logger == nil { + logger = slog.Default() + } + return &Client{ + apiToken: token, + baseURL: fmt.Sprintf( + "https://api.cloudflare.com/client/v4/accounts/%s/storage/kv/namespaces/%s/", + accountID, namespaceID, + ), + client: &http.Client{Timeout: 10 * time.Second}, + logger: logger, + } +} + +func (c *Client) PutEmailDid(email, did string) { + c.write("put email→did", strings.ToLower(email), did) +} + +func (c *Client) SetPrimaryEmail(did, email string) { + c.write("put primary→email", "__primary:"+did, email) +} + +func (c *Client) DeleteEmailKey(email string) { + c.delete("delete email key", strings.ToLower(email)) +} + +func (c *Client) write(action, key, value string) { + if c == nil { + return + } + target := c.baseURL + "values/" + url.PathEscape(key) + req, err := http.NewRequest(http.MethodPut, target, strings.NewReader(value)) + if err != nil { + c.logger.Warn("kv "+action+" failed to build request", "err", err, "key", key) + return + } + req.Header.Set("Authorization", "Bearer "+c.apiToken) + req.Header.Set("Content-Type", "text/plain") + resp, err := c.client.Do(req) + if err != nil { + c.logger.Warn("kv "+action+" request failed", "err", err, "key", key) + return + } + defer resp.Body.Close() + body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + c.logger.Warn("kv "+action+" rejected", "status", resp.Status, "key", key, "body", string(body)) + return + } + c.logger.Info("kv "+action+" ok", "key", key) +} + +func (c *Client) delete(action, key string) { + if c == nil { + return + } + target := c.baseURL + "values/" + url.PathEscape(key) + req, err := http.NewRequest(http.MethodDelete, target, nil) + if err != nil { + c.logger.Warn("kv "+action+" failed to build request", "err", err, "key", key) + return + } + req.Header.Set("Authorization", "Bearer "+c.apiToken) + resp, err := c.client.Do(req) + if err != nil { + c.logger.Warn("kv "+action+" request failed", "err", err, "key", key) + return + } + defer resp.Body.Close() + body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + c.logger.Warn("kv "+action+" rejected", "status", resp.Status, "key", key, "body", string(body)) + return + } + c.logger.Info("kv "+action+" ok", "key", key) +} diff --git a/deliberi/kv/kv_test.go b/deliberi/kv/kv_test.go new file mode 100644 index 000000000..be4831a5c --- /dev/null +++ b/deliberi/kv/kv_test.go @@ -0,0 +1,190 @@ +package kv + +import ( + "bytes" + "io" + "log/slog" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + "time" +) + +type recordedRequest struct { + method string + path string + body string + auth string + ctype string +} + +func newRecordingClient(t *testing.T) (*Client, func() []recordedRequest) { + t.Helper() + + var mu sync.Mutex + var reqs []recordedRequest + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(r.Body) + mu.Lock() + reqs = append(reqs, recordedRequest{ + method: r.Method, + path: r.URL.Path, + body: string(body), + auth: r.Header.Get("Authorization"), + ctype: r.Header.Get("Content-Type"), + }) + mu.Unlock() + w.WriteHeader(http.StatusOK) + })) + t.Cleanup(srv.Close) + + c := &Client{ + apiToken: "test-token", + baseURL: srv.URL + "/accounts/acct/storage/kv/namespaces/ns/", + client: srv.Client(), + logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + } + + snapshot := func() []recordedRequest { + mu.Lock() + defer mu.Unlock() + out := make([]recordedRequest, len(reqs)) + copy(out, reqs) + return out + } + return c, snapshot +} + +func TestPutEmailDid(t *testing.T) { + c, snapshot := newRecordingClient(t) + + c.PutEmailDid("Bob@Example.com", "did:plc:bob") + + got := snapshot() + if len(got) != 1 { + t.Fatalf("got %d requests, want 1", len(got)) + } + r := got[0] + if r.method != http.MethodPut { + t.Errorf("method = %s, want PUT", r.method) + } + // url.PathEscape leaves @ unescaped; the worker looks up the literal key + if !strings.HasSuffix(r.path, "/values/bob@example.com") { + t.Errorf("path = %s, want suffix /values/bob@example.com", r.path) + } + if r.body != "did:plc:bob" { + t.Errorf("body = %q, want did:plc:bob", r.body) + } + if r.auth != "Bearer test-token" { + t.Errorf("auth = %q, want Bearer test-token", r.auth) + } + if r.ctype != "text/plain" { + t.Errorf("content-type = %q, want text/plain", r.ctype) + } + if !strings.Contains(r.path, "bob@example.com") { + t.Errorf("key was not lowercased: %s", r.path) + } +} + +func TestSetPrimaryEmail(t *testing.T) { + c, snapshot := newRecordingClient(t) + + c.SetPrimaryEmail("did:plc:bob", "bob@example.com") + + got := snapshot() + if len(got) != 1 { + t.Fatalf("got %d requests, want 1", len(got)) + } + r := got[0] + if r.method != http.MethodPut { + t.Errorf("method = %s, want PUT", r.method) + } + if !strings.HasSuffix(r.path, "/values/__primary:did:plc:bob") { + t.Errorf("path = %s, want suffix /values/__primary:did:plc:bob", r.path) + } + if r.body != "bob@example.com" { + t.Errorf("body = %q, want bob@example.com", r.body) + } +} + +func TestDeleteEmailKey(t *testing.T) { + c, snapshot := newRecordingClient(t) + + c.DeleteEmailKey("Bob@Example.com") + + got := snapshot() + if len(got) != 1 { + t.Fatalf("got %d requests, want 1", len(got)) + } + r := got[0] + if r.method != http.MethodDelete { + t.Errorf("method = %s, want DELETE", r.method) + } + if !strings.HasSuffix(r.path, "/values/bob@example.com") { + t.Errorf("path = %s, want suffix /values/bob@example.com", r.path) + } + if r.auth != "Bearer test-token" { + t.Errorf("auth = %q, want Bearer test-token", r.auth) + } + if r.ctype != "" { + t.Errorf("content-type = %q, want none on a bodyless DELETE", r.ctype) + } +} + +func TestFailureLoggedNotFatal(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + t.Cleanup(srv.Close) + + c := &Client{ + apiToken: "test-token", + baseURL: srv.URL + "/", + client: srv.Client(), + logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + } + + c.PutEmailDid("bob@example.com", "did:plc:bob") + c.DeleteEmailKey("bob@example.com") +} + +func TestSuccessLogged(t *testing.T) { + var buf bytes.Buffer + logger := slog.New(slog.NewTextHandler(&buf, nil)) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusOK) + })) + t.Cleanup(srv.Close) + + c := &Client{apiToken: "t", baseURL: srv.URL + "/", client: srv.Client(), logger: logger} + c.PutEmailDid("bob@example.com", "did:plc:bob") + c.DeleteEmailKey("bob@example.com") + + out := buf.String() + for _, want := range []string{"put email", "delete email", "bob@example.com"} { + if !strings.Contains(out, want) { + t.Errorf("missing %q in success log:\n%s", want, out) + } + } +} + +func TestNewClientSetsTimeout(t *testing.T) { + c := NewClient("t", "acct", "ns", slog.New(slog.NewTextHandler(io.Discard, nil))) + if c.client.Timeout != 10*time.Second { + t.Errorf("client timeout = %v, want 10s", c.client.Timeout) + } + if c.baseURL == "" || c.apiToken != "t" { + t.Errorf("client url/token not configured: %q %q", c.baseURL, c.apiToken) + } +} + +func TestNilClientSafe(t *testing.T) { + var c *Client + + c.PutEmailDid("bob@example.com", "did:plc:bob") + c.SetPrimaryEmail("did:plc:bob", "bob@example.com") + c.DeleteEmailKey("bob@example.com") +} -- 2.51.2