From 965c8c97bb8e19ac50523fdc144b9f2dc1f5b17d Mon Sep 17 00:00:00 2001 From: "Customer.io Open Source Bot" Date: Tue, 8 Sep 2026 11:06:15 +0530 Subject: [PATCH] Add Customer.io CLI source CioCliPublicExport-RevId: 982c65f90fa5722429f5f6ec5a5ea2d09d7dbce7 --- cmd/api.go | 47 ++++- cmd/api_preflight_cachekey_test.go | 145 ++++++++++++++++ cmd/api_preflight_coldstart_test.go | 198 ++++++++++++++++++++++ cmd/api_preflight_test.go | 29 +++- cmd/helpers.go | 31 +++- cmd/schema.go | 2 +- internal/client/client.go | 46 ++++- internal/client/client_test.go | 33 ++++ internal/filelock/filelock.go | 25 ++- internal/filelock/filelock_try_test.go | 101 +++++++++++ internal/filelock/filelock_unix.go | 11 ++ internal/filelock/filelock_windows.go | 17 +- internal/routes/cache.go | 72 ++++++-- internal/routes/cache_nonblocking_test.go | 173 +++++++++++++++++++ internal/routes/cache_test.go | 10 +- internal/routes/pathindex.go | 53 ++++-- 16 files changed, 928 insertions(+), 65 deletions(-) create mode 100644 cmd/api_preflight_cachekey_test.go create mode 100644 cmd/api_preflight_coldstart_test.go create mode 100644 internal/filelock/filelock_try_test.go create mode 100644 internal/routes/cache_nonblocking_test.go diff --git a/cmd/api.go b/cmd/api.go index 62ab533..54ea766 100644 --- a/cmd/api.go +++ b/cmd/api.go @@ -1,11 +1,13 @@ package cmd import ( + "context" "encoding/json" "fmt" "net/url" "regexp" "strings" + "time" "github.com/customerio/cli/internal/client" "github.com/customerio/cli/internal/output" @@ -215,6 +217,14 @@ func requestBody(jsonBody json.RawMessage, fileParts []client.FilePart) (*client return client.NewMultipartBody(fileParts, fields) } +// preflightSpecFetchBudget bounds everything the path check does over the +// network when it finds the cache cold: the token exchange, if the caller's +// credential needs one, and the one spec download. The spec is ~2.7MB and +// normally arrives well inside this; a branch that does not finish in time is +// a sign of a slow or stalled host, and the check gives up rather than hold +// the caller's request behind it. +const preflightSpecFetchBudget = 5 * time.Second + // preflightPath rejects a path the API spec does not describe, before the // request is sent. // @@ -234,15 +244,38 @@ func preflightPath(cmd *cobra.Command, c *client.Client, httpMethod, resolvedPat return nil } - // Whatever the spec cache already holds, read without a lock or a - // download: this sits in front of every request, so it must not be able to - // block on another process or wait on the network. A cold cache means the - // check has no opinion, and `cio schema` is what fills it. + // Read what the cache already holds first: no lock, no network, ~25ms. idx, err := routes.LoadPathIndexFromCache(specCacheOptions(c)) if err != nil { - // Fail open: an absent or unparseable spec must never stop a caller - // from reaching an endpoint that does exist. - return nil + // Cold cache. Fetch the spec once so the check can exist at all — a + // session that never runs `cio schema` would otherwise never be + // checked. The first version of this check refused to download here, + // because EnsureSpecs then meant an unbounded flock and two 30s + // fetches in front of every request; this fetch is bounded in each of + // those respects instead: the lock is tried, not waited for (another + // process downloading means we step aside), the download runs under a + // short deadline, no prose reaches stderr, and any of those failing + // means proceeding without an opinion. The identity's token is used so + // what lands in the cache is the same plan-filtered spec `cio schema` + // would fetch, never an anonymous one that would mislead it later. + // One budget covers the whole branch — the token exchange as well as + // the download — so a stalled token endpoint cannot hold the request + // any longer than a stalled spec host can. + ctx, cancel := context.WithTimeout(cmd.Context(), preflightSpecFetchBudget) + defer cancel() + opts := specLoadOptions(ctx, c) + if c.ServiceAccountToken() != "" && opts.AccessToken == "" { + // The token could not be exchanged, or not within the budget. + // `cio schema` falls back to an anonymous fetch here; this check + // must not: an anonymous spec is the wrong thing to judge a + // plan-filtered identity against, and it would land in a partition + // the identity's reads never look in. + return nil + } + idx, err = routes.EnsurePathIndex(ctx, opts) + if err != nil { + return nil + } } // The route exists: nothing to block, so let the request proceed. The diff --git a/cmd/api_preflight_cachekey_test.go b/cmd/api_preflight_cachekey_test.go new file mode 100644 index 0000000..bb59fdd --- /dev/null +++ b/cmd/api_preflight_cachekey_test.go @@ -0,0 +1,145 @@ +package cmd + +import ( + "encoding/base64" + "encoding/json" + "testing" +) + +// unsignedJWT builds a syntactically valid JWT with the given payload and a +// placeholder signature. The CLI never verifies signatures — it forwards the +// token and reads claims as hints — so the test server accepts it as-is. +func unsignedJWT(t *testing.T, claims map[string]any) string { + t.Helper() + header := base64.RawURLEncoding.EncodeToString([]byte(`{"alg":"EdDSA","typ":"JWT"}`)) + payload, err := json.Marshal(claims) + if err != nil { + t.Fatalf("marshal claims: %v", err) + } + return header + "." + base64.RawURLEncoding.EncodeToString(payload) + ".sig" +} + +// farFuture keeps the fixture tokens unexpired so the client never tries to +// refresh one. +const farFuture = 4102444800 // 2100-01-01T00:00:00Z + +// A caller authenticating with a pre-exchanged JWT gets a freshly signed token +// from its issuer on every command: same session, new iat, different bytes. +// The cache written by one command must be found by the next, or the path +// check never has a spec to judge against. +func TestAPIPreflight_CacheSurvivesTokenResign(t *testing.T) { + server := newPreflightServer(t) + t.Setenv("CIO_TOKEN", "") + + // First command: warm the cache under this session's token. + t.Setenv("CIO_ACCESS_TOKEN", unsignedJWT(t, map[string]any{ + "jti": "session-A", "sub": "sa:1", "iat": 1700000000, "exp": farFuture, + })) + if _, _, err := executeCommand("schema", "--api-url", server.URL); err != nil { + t.Fatalf("warming the cache via schema: %v", err) + } + server.forget() + + // Second command: same session, re-signed one second later. + t.Setenv("CIO_ACCESS_TOKEN", unsignedJWT(t, map[string]any{ + "jti": "session-A", "sub": "sa:1", "iat": 1700000001, "exp": farFuture, + })) + _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", + "--api-url", server.URL) + if err == nil { + t.Fatal("expected the re-signed token to find the cache and reject the invented path") + } + if got := server.requested(); len(got) != 0 { + t.Errorf("no API request should have been made, got %v", got) + } + if got := server.specFetches(); len(got) != 0 { + t.Errorf("the second command must reuse the cache, not fetch; got %v", got) + } +} + +// The same keying serves `cio schema`: a second call under a re-signed token +// must read the spec it cached moments ago, not download it again. +func TestSchema_ReusesCacheAcrossTokenResign(t *testing.T) { + server := newPreflightServer(t) + t.Setenv("CIO_TOKEN", "") + + t.Setenv("CIO_ACCESS_TOKEN", unsignedJWT(t, map[string]any{ + "jti": "session-A", "sub": "sa:1", "iat": 1700000000, "exp": farFuture, + })) + if _, _, err := executeCommand("schema", "--api-url", server.URL); err != nil { + t.Fatalf("first schema: %v", err) + } + if got := server.specFetches(); len(got) != 2 { + t.Fatalf("first schema should download both specs, got %v", got) + } + server.forget() + + t.Setenv("CIO_ACCESS_TOKEN", unsignedJWT(t, map[string]any{ + "jti": "session-A", "sub": "sa:1", "iat": 1700000001, "exp": farFuture, + })) + if _, _, err := executeCommand("schema", "--api-url", server.URL); err != nil { + t.Fatalf("second schema: %v", err) + } + if got := server.specFetches(); len(got) != 0 { + t.Errorf("second schema under a re-signed token must reuse the cache, got fetches %v", got) + } +} + +// Partitioning by session must still hold: a different session does not see +// another session's cache, because the served spec may differ per identity. +func TestAPIPreflight_CacheIsolatedPerSession(t *testing.T) { + server := newPreflightServer(t) + t.Setenv("CIO_TOKEN", "") + + t.Setenv("CIO_ACCESS_TOKEN", unsignedJWT(t, map[string]any{ + "jti": "session-A", "sub": "sa:1", "iat": 1700000000, "exp": farFuture, + })) + if _, _, err := executeCommand("schema", "--api-url", server.URL); err != nil { + t.Fatalf("warming the cache via schema: %v", err) + } + server.forget() + + // A different session: its cache is cold, so it must fetch its own spec — + // never read session A's, which may be filtered for a different plan. The + // tell is the download: A's spec is already on disk, and B fetches anyway. + t.Setenv("CIO_ACCESS_TOKEN", unsignedJWT(t, map[string]any{ + "jti": "session-B", "sub": "sa:2", "iat": 1700000001, "exp": farFuture, + })) + if _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", + "--api-url", server.URL); err == nil { + t.Fatal("expected session B's own freshly fetched spec to reject the invented path") + } + if got := server.specFetches(); len(got) != 2 { + t.Errorf("session B must download its own spec rather than read session A's, got fetches %v", got) + } + if got := server.requested(); len(got) != 0 { + t.Errorf("no API request should have been made, got %v", got) + } +} + +// A token that is not a JWT, or a JWT without jti, keys on its raw bytes — +// the behaviour every caller had before, unchanged. +func TestSpecCacheKey_FallsBackToRawToken(t *testing.T) { + for _, token := range []string{ + "sa_live_abc123", + "opaque-token", + unsignedJWT(t, map[string]any{"sub": "sa:1", "exp": farFuture}), // JWT, no jti + } { + if got := specCacheKey(token); got != token { + t.Errorf("specCacheKey(%q) = %q, want the raw token", token, got) + } + } +} + +func TestSpecCacheKey_UsesJTIAcrossResigns(t *testing.T) { + a := specCacheKey(unsignedJWT(t, map[string]any{"jti": "session-A", "iat": 1, "exp": farFuture})) + b := specCacheKey(unsignedJWT(t, map[string]any{"jti": "session-A", "iat": 2, "exp": farFuture})) + c := specCacheKey(unsignedJWT(t, map[string]any{"jti": "session-B", "iat": 2, "exp": farFuture})) + + if a != "session-A" || a != b { + t.Errorf("same session must yield the same key across re-signs; got %q and %q", a, b) + } + if c == a { + t.Errorf("different sessions must yield different keys; both %q", a) + } +} diff --git a/cmd/api_preflight_coldstart_test.go b/cmd/api_preflight_coldstart_test.go new file mode 100644 index 0000000..eeb1156 --- /dev/null +++ b/cmd/api_preflight_coldstart_test.go @@ -0,0 +1,198 @@ +package cmd + +import ( + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" +) + +// A spec host that never answers must cost the caller the fetch budget at most, +// and then the request must go out as if there were no check at all. +func TestAPIPreflight_StalledSpecHostFailsOpenWithinBudget(t *testing.T) { + t.Setenv("HOME", t.TempDir()) + t.Setenv("CIO_TOKEN", "sa_live_test123") + t.Setenv("CIO_ACCESS_TOKEN", "") + + release := make(chan struct{}) + + var mu sync.Mutex + var apiPaths []string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/v1/service_accounts/oauth/token": + _, _ = w.Write([]byte(`{"access_token":"jwt-test-session","token_type":"Bearer","expires_in":3600}`)) + case "/v1/openapi.json", "/cdp/api/openapi.json": + // Hold the response until the test ends; the client must give up first. + select { + case <-release: + case <-r.Context().Done(): + } + default: + mu.Lock() + apiPaths = append(apiPaths, r.URL.Path) + mu.Unlock() + _, _ = w.Write([]byte(`{"ok":true}`)) + } + })) + // Deferred LIFO: release the stalled handler before Close waits on it. + defer server.Close() + defer close(release) + + start := time.Now() + _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", "--api-url", server.URL) + elapsed := time.Since(start) + + if err != nil { + t.Fatalf("a stalled spec host must not fail the call, got: %v", err) + } + mu.Lock() + got := append([]string(nil), apiPaths...) + mu.Unlock() + if len(got) != 1 { + t.Errorf("expected the request to be sent after giving up on the spec, got %v", got) + } + // Anything near the 30s HTTP client timeout means the deadline was not + // honoured. Generous headroom over the budget keeps this off the flake list. + if elapsed > preflightSpecFetchBudget+5*time.Second { + t.Errorf("call took %v; the spec fetch must be bounded by the %v budget", elapsed, preflightSpecFetchBudget) + } +} + +// A service-account token that cannot be exchanged must not make the check +// fetch anonymously: that spec is the wrong one to judge the identity against, +// and it would be written to a partition the identity's own reads never open. +// The check steps aside and the request proceeds to fail — or not — on its own. +func TestAPIPreflight_UnexchangeableTokenSkipsCheckWithoutAnonymousFetch(t *testing.T) { + home := t.TempDir() + t.Setenv("HOME", home) + t.Setenv("CIO_TOKEN", "sa_live_revoked") + t.Setenv("CIO_ACCESS_TOKEN", "") + + var mu sync.Mutex + var specFetches, apiPaths []string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/v1/service_accounts/oauth/token": + w.WriteHeader(http.StatusUnauthorized) + _, _ = w.Write([]byte(`{"error":"invalid_client"}`)) + case "/v1/openapi.json", "/cdp/api/openapi.json": + mu.Lock() + specFetches = append(specFetches, r.URL.Path) + mu.Unlock() + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(preflightSpec())) + default: + mu.Lock() + apiPaths = append(apiPaths, r.URL.Path) + mu.Unlock() + w.WriteHeader(http.StatusUnauthorized) + _, _ = w.Write([]byte(`{"error":"unauthorized"}`)) + } + })) + defer server.Close() + + // The call fails on auth, as it should — at the token exchange, which the + // client performs before it can send anything — and not on the check. + _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", "--api-url", server.URL) + if err == nil { + t.Fatal("expected the call to fail authentication") + } + if strings.Contains(err.Error(), "not an endpoint in the API spec") { + t.Fatalf("the check must step aside when the token cannot be exchanged, got: %v", err) + } + if !strings.Contains(err.Error(), "token exchange failed") { + t.Fatalf("expected the failure to be the token exchange, got: %v", err) + } + + mu.Lock() + fetched := append([]string(nil), specFetches...) + sent := append([]string(nil), apiPaths...) + mu.Unlock() + if len(fetched) != 0 { + t.Errorf("must not fetch a spec anonymously, got %v", fetched) + } + // No Bearer could be minted, so nothing reaches the API; documented here so + // a future client that sends anyway shows up as a deliberate change. + if len(sent) != 0 { + t.Errorf("expected no API request without a token, got %v", sent) + } + if _, statErr := os.Stat(filepath.Join(home, ".cio", "cache", "specs", "openapi.json")); statErr == nil { + t.Error("an anonymous spec must not be written to the shared base partition") + } +} + +// The budget covers the token exchange too. A token endpoint that stalls on +// the check's exchange must cost at most the budget; the check then steps +// aside, and the request's own exchange (which is answered) goes through. +func TestAPIPreflight_StalledTokenExchangeIsBoundedByBudget(t *testing.T) { + t.Setenv("HOME", t.TempDir()) + t.Setenv("CIO_TOKEN", "sa_live_test123") + t.Setenv("CIO_ACCESS_TOKEN", "") + + // Deferred LIFO: the server must not be closed while its stalled handler + // still waits, so release is deferred after Close and therefore runs first. + release := make(chan struct{}) + + var mu sync.Mutex + var exchanges int + var specFetches, apiPaths []string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/v1/service_accounts/oauth/token": + mu.Lock() + exchanges++ + first := exchanges == 1 + mu.Unlock() + if first { + // The check's exchange: never answered. + select { + case <-release: + case <-r.Context().Done(): + } + return + } + _, _ = w.Write([]byte(`{"access_token":"jwt-test-session","token_type":"Bearer","expires_in":3600}`)) + case "/v1/openapi.json", "/cdp/api/openapi.json": + mu.Lock() + specFetches = append(specFetches, r.URL.Path) + mu.Unlock() + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(preflightSpec())) + default: + mu.Lock() + apiPaths = append(apiPaths, r.URL.Path) + mu.Unlock() + _, _ = w.Write([]byte(`{"ok":true}`)) + } + })) + defer server.Close() + defer close(release) + + start := time.Now() + _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", "--api-url", server.URL) + elapsed := time.Since(start) + + if err != nil { + t.Fatalf("a stalled exchange on the check must not fail the call, got: %v", err) + } + mu.Lock() + fetched := append([]string(nil), specFetches...) + sent := append([]string(nil), apiPaths...) + mu.Unlock() + if len(fetched) != 0 { + t.Errorf("without a token the check must not fetch, got %v", fetched) + } + if len(sent) != 1 { + t.Errorf("expected the request to be sent once the check stepped aside, got %v", sent) + } + // Without the budget on the exchange this would sit for the 30s client + // timeout before the check even gave up. + if elapsed > preflightSpecFetchBudget+5*time.Second { + t.Errorf("call took %v; the exchange must be bounded by the %v budget", elapsed, preflightSpecFetchBudget) + } +} diff --git a/cmd/api_preflight_test.go b/cmd/api_preflight_test.go index b5c2ffb..7fa7dae 100644 --- a/cmd/api_preflight_test.go +++ b/cmd/api_preflight_test.go @@ -289,20 +289,31 @@ func TestAPIPreflight_NeverFetchesTheSpecItself(t *testing.T) { } } -// A cold cache means no opinion: the call proceeds, and still nothing is -// fetched to form one. -func TestAPIPreflight_ColdCacheIsInertAndSilent(t *testing.T) { +// A cold cache is filled, once, so the check can judge — and the call it was +// asked about is judged against what was just fetched. +func TestAPIPreflight_ColdCacheFetchesOnceThenJudges(t *testing.T) { server := newPreflightServer(t) - if _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", - "--api-url", server.URL); err != nil { - t.Fatalf("a cold cache must not block the call, got: %v", err) + _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", + "--api-url", server.URL) + if err == nil { + t.Fatal("expected the freshly fetched spec to reject the invented path") } - if got := server.requested(); len(got) != 1 { - t.Errorf("expected the request to be sent, got %v", got) + if got := server.requested(); len(got) != 0 { + t.Errorf("no API request should have been made, got %v", got) + } + if got := server.specFetches(); len(got) != 2 { + t.Errorf("expected exactly one download of each spec, got %v", got) + } + + // Warm now: the next check reads the cache and fetches nothing. + server.forget() + if _, _, err := executeCommand("api", "/v1/environments/456/campaigns/48/actions", + "--api-url", server.URL); err == nil { + t.Fatal("expected the cached spec to reject the invented path") } if got := server.specFetches(); len(got) != 0 { - t.Errorf("the check must not download specs to populate itself, got %v", got) + t.Errorf("a warm cache must not fetch, got %v", got) } } diff --git a/cmd/helpers.go b/cmd/helpers.go index c1e51c2..0e9cf24 100644 --- a/cmd/helpers.go +++ b/cmd/helpers.go @@ -2,6 +2,7 @@ package cmd import ( "bytes" + "context" "encoding/json" "fmt" "io" @@ -13,10 +14,32 @@ import ( "github.com/spf13/cobra" ) +// specCacheKey picks the value that partitions the spec cache for a +// pre-exchanged access token. +// +// The cache is per identity because the served spec is: an authenticated +// fetch returns the plan-filtered spec for that account, and two identities +// must never read each other's. A service-account token is a stable secret and +// serves as its own key. A pre-exchanged JWT does not: its issuer re-signs it +// on every request, so the bytes change while the identity does not, and a key +// made of those bytes hands every call an empty cache — and a fresh download. +// The token's jti names the session and survives the re-signs, so it is the +// key, with the raw token as the fallback for a JWT that carries none. +func specCacheKey(accessToken string) string { + if id, ok := client.JWTID(accessToken); ok { + return id + } + return accessToken +} + // specLoadOptions builds the spec-download options shared by the commands that // read the route registry, so `cio api`'s pre-send check and `cio schema` // always resolve the same spec — and the same cache entry — for one identity. -func specLoadOptions(cmd *cobra.Command, c *client.Client) routes.LoadRegistryOptions { +// +// A service-account caller has its token exchanged here, under ctx: a caller +// working to a deadline passes it so the exchange is bounded like the download +// that follows it, rather than being an unbounded step in front of a bounded one. +func specLoadOptions(ctx context.Context, c *client.Client) routes.LoadRegistryOptions { var opts routes.LoadRegistryOptions if c == nil { return opts @@ -25,7 +48,7 @@ func specLoadOptions(cmd *cobra.Command, c *client.Client) routes.LoadRegistryOp opts.BaseURL = c.BaseURL() switch { case c.ServiceAccountToken() != "": - jwt, err := c.EnsureAccessToken(cmd.Context()) + jwt, err := c.EnsureAccessToken(ctx) if err == nil { opts.AccessToken = jwt opts.CacheKey = c.ServiceAccountToken() @@ -34,7 +57,7 @@ func specLoadOptions(cmd *cobra.Command, c *client.Client) routes.LoadRegistryOp // Pre-exchanged JWT (e.g. CIO_ACCESS_TOKEN, as the in-product agent uses) // — send it so the server returns the plan-filtered spec, not the full one. opts.AccessToken = c.AccessToken() - opts.CacheKey = opts.AccessToken + opts.CacheKey = specCacheKey(opts.AccessToken) } return opts } @@ -56,7 +79,7 @@ func specCacheOptions(c *client.Client) routes.LoadRegistryOptions { case c.ServiceAccountToken() != "": opts.CacheKey = c.ServiceAccountToken() case c.AccessToken() != "": - opts.CacheKey = c.AccessToken() + opts.CacheKey = specCacheKey(c.AccessToken()) } return opts } diff --git a/cmd/schema.go b/cmd/schema.go index 8cea6d0..b6e0c7d 100644 --- a/cmd/schema.go +++ b/cmd/schema.go @@ -39,7 +39,7 @@ func runSchema(cmd *cobra.Command, args []string) error { refresh, _ := cmd.Flags().GetBool("refresh") compact, _ := cmd.Flags().GetBool("compact") - opts := specLoadOptions(cmd, clientFromCmd(cmd)) + opts := specLoadOptions(cmd.Context(), clientFromCmd(cmd)) opts.ForceRefresh = refresh if GetDryRun(cmd) { diff --git a/internal/client/client.go b/internal/client/client.go index 6470e6c..ff587d4 100644 --- a/internal/client/client.go +++ b/internal/client/client.go @@ -730,26 +730,37 @@ func MintLoginCLILink(ctx context.Context, baseURL, saToken string, timeout time return &out, nil } -// parseJWTExpiry extracts the exp claim from a JWT without verifying the -// signature. Returns the expiry time or an error if the token is not a -// valid 3-part JWT or lacks an exp claim. -func parseJWTExpiry(token string) (time.Time, error) { +// decodeJWTClaims decodes a JWT's payload into claims without verifying the +// signature. The client reads claims only as hints — an expiry to schedule a +// refresh, an identifier to partition a cache — never to decide what a token +// may do; the server verifies every request it receives. +func decodeJWTClaims(token string, claims any) error { parts := strings.SplitN(token, ".", 4) if len(parts) != 3 { - return time.Time{}, fmt.Errorf("not a JWT") + return fmt.Errorf("not a JWT") } // The payload is base64url-encoded (no padding). payload, err := base64.RawURLEncoding.DecodeString(parts[1]) if err != nil { - return time.Time{}, fmt.Errorf("decode JWT payload: %w", err) + return fmt.Errorf("decode JWT payload: %w", err) } + if err := json.Unmarshal(payload, claims); err != nil { + return fmt.Errorf("parse JWT claims: %w", err) + } + return nil +} + +// parseJWTExpiry extracts the exp claim from a JWT without verifying the +// signature. Returns the expiry time or an error if the token is not a +// valid 3-part JWT or lacks an exp claim. +func parseJWTExpiry(token string) (time.Time, error) { var claims struct { Exp json.Number `json:"exp"` } - if err := json.Unmarshal(payload, &claims); err != nil { - return time.Time{}, fmt.Errorf("parse JWT claims: %w", err) + if err := decodeJWTClaims(token, &claims); err != nil { + return time.Time{}, err } expFloat, err := claims.Exp.Float64() @@ -759,3 +770,22 @@ func parseJWTExpiry(token string) (time.Time, error) { return time.Unix(int64(expFloat), 0), nil } + +// JWTID returns a JWT's jti claim without verifying the signature, and false +// when the token is not a JWT or carries no jti. +// +// For the session tokens this CLI is handed, jti names the session rather +// than the individual signing: the issuer signs a fresh token for the same +// session on each request, which changes iat and with it every byte of the +// token, while jti stays the same. That makes jti the handle for anything +// that must hold steady across those re-signs — a cache partition, say — +// where the token bytes do not. +func JWTID(token string) (string, bool) { + var claims struct { + ID string `json:"jti"` + } + if err := decodeJWTClaims(token, &claims); err != nil || claims.ID == "" { + return "", false + } + return claims.ID, true +} diff --git a/internal/client/client_test.go b/internal/client/client_test.go index c9ef9b6..6572b60 100644 --- a/internal/client/client_test.go +++ b/internal/client/client_test.go @@ -678,6 +678,39 @@ func TestParseJWTExpiry_OpaqueToken(t *testing.T) { } } +func TestJWTID_ValidJWT(t *testing.T) { + // Payload: {"jti":"sess-1","exp":1700000000} + token := "eyJhbGciOiJIUzI1NiJ9.eyJqdGkiOiJzZXNzLTEiLCJleHAiOjE3MDAwMDAwMDB9.signature" + id, ok := JWTID(token) + if !ok { + t.Fatal("expected a jti to be found") + } + if id != "sess-1" { + t.Errorf("expected sess-1, got %q", id) + } +} + +func TestJWTID_MissingJTI(t *testing.T) { + // Payload: {"sub":"test"} — a JWT with no jti reports false, not an empty key. + if _, ok := JWTID("eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiJ0ZXN0In0.signature"); ok { + t.Fatal("expected no jti for a JWT without one") + } +} + +func TestJWTID_EmptyJTI(t *testing.T) { + // Payload: {"jti":"","exp":1700000000} — present but empty is treated as absent, + // so it can never become a shared cache partition. + if _, ok := JWTID("eyJhbGciOiJIUzI1NiJ9.eyJqdGkiOiIiLCJleHAiOjE3MDAwMDAwMDB9.signature"); ok { + t.Fatal("expected an empty jti to be treated as absent") + } +} + +func TestJWTID_OpaqueToken(t *testing.T) { + if _, ok := JWTID("sa_live_abc123"); ok { + t.Fatal("expected no jti for an opaque token") + } +} + func TestClient_Do_ValidateHeader(t *testing.T) { // X-Validate: strict must be stamped on every request so the server // rejects unknown JSON fields with a 400 instead of silently dropping diff --git a/internal/filelock/filelock.go b/internal/filelock/filelock.go index 449b5ce..f775581 100644 --- a/internal/filelock/filelock.go +++ b/internal/filelock/filelock.go @@ -1,19 +1,42 @@ package filelock import ( + "errors" "fmt" "os" ) +// ErrLocked reports that another process holds the lock. TryLock returns it +// instead of waiting, so the caller decides what a held lock means for it. +var ErrLocked = errors.New("filelock: held by another process") + // Lock acquires an exclusive lock on path and returns an unlock function. +// +// It blocks until the lock is free. There is no timeout — the underlying +// primitive has none, and a portable one cannot be layered on top — so a +// caller that cannot afford an open-ended wait must use TryLock instead. func Lock(path string, mode os.FileMode) (func(), error) { + return acquire(path, mode, lockFile) +} + +// TryLock acquires an exclusive lock on path without waiting. When another +// process holds the lock it returns ErrLocked at once; any other failure is +// reported as it is for Lock. +func TryLock(path string, mode os.FileMode) (func(), error) { + return acquire(path, mode, tryLockFile) +} + +func acquire(path string, mode os.FileMode, lock func(*os.File) error) (func(), error) { f, err := os.OpenFile(path, os.O_CREATE|os.O_RDWR, mode) if err != nil { return nil, fmt.Errorf("open lock file: %w", err) } - if err := lockFile(f); err != nil { + if err := lock(f); err != nil { f.Close() + if errors.Is(err, ErrLocked) { + return nil, err + } return nil, fmt.Errorf("acquire lock: %w", err) } diff --git a/internal/filelock/filelock_try_test.go b/internal/filelock/filelock_try_test.go new file mode 100644 index 0000000..e521a2b --- /dev/null +++ b/internal/filelock/filelock_try_test.go @@ -0,0 +1,101 @@ +package filelock + +import ( + "errors" + "path/filepath" + "testing" + "time" +) + +func TestTryLock_AcquiresWhenFree(t *testing.T) { + path := filepath.Join(t.TempDir(), "test.lock") + + unlock, err := TryLock(path, 0600) + if err != nil { + t.Fatalf("TryLock() on a free lock: %v", err) + } + unlock() +} + +// The whole point of TryLock: a held lock is reported, not waited for. Both +// flock and LockFileEx conflict across separate opens of the same file within +// one process, so holding via Lock and probing via TryLock is a real contest. +func TestTryLock_HeldReturnsErrLockedWithoutWaiting(t *testing.T) { + path := filepath.Join(t.TempDir(), "test.lock") + + unlock, err := Lock(path, 0600) + if err != nil { + t.Fatalf("Lock(): %v", err) + } + defer unlock() + + start := time.Now() + _, err = TryLock(path, 0600) + elapsed := time.Since(start) + + if !errors.Is(err, ErrLocked) { + t.Fatalf("TryLock() on a held lock: want ErrLocked, got %v", err) + } + // Lock would sit here until the holder released; TryLock must not. The + // bound is generous so a slow CI runner cannot turn it into a flake — any + // blocking implementation would exceed it by orders of magnitude. + if elapsed > 2*time.Second { + t.Errorf("TryLock() took %v on a held lock; it must return at once", elapsed) + } +} + +func TestTryLock_SucceedsOnceReleased(t *testing.T) { + path := filepath.Join(t.TempDir(), "test.lock") + + unlock, err := Lock(path, 0600) + if err != nil { + t.Fatalf("Lock(): %v", err) + } + if _, err := TryLock(path, 0600); !errors.Is(err, ErrLocked) { + t.Fatalf("expected ErrLocked while held, got %v", err) + } + unlock() + + unlock, err = TryLock(path, 0600) + if err != nil { + t.Fatalf("TryLock() after release: %v", err) + } + unlock() +} + +// Lock's error contract is unchanged: it still blocks, and still succeeds once +// the holder lets go — TryLock is an addition, not a change to Lock. +func TestLock_StillBlocksThenAcquires(t *testing.T) { + path := filepath.Join(t.TempDir(), "test.lock") + + unlock, err := Lock(path, 0600) + if err != nil { + t.Fatalf("Lock(): %v", err) + } + + acquired := make(chan struct{}) + go func() { + u, err := Lock(path, 0600) + if err != nil { + t.Errorf("second Lock(): %v", err) + close(acquired) + return + } + u() + close(acquired) + }() + + select { + case <-acquired: + t.Fatal("second Lock() returned while the first was still held") + case <-time.After(200 * time.Millisecond): + // Still waiting, as it should be. + } + + unlock() + select { + case <-acquired: + case <-time.After(5 * time.Second): + t.Fatal("second Lock() did not acquire after the first was released") + } +} diff --git a/internal/filelock/filelock_unix.go b/internal/filelock/filelock_unix.go index cc67a54..c7dafab 100644 --- a/internal/filelock/filelock_unix.go +++ b/internal/filelock/filelock_unix.go @@ -3,6 +3,7 @@ package filelock import ( + "errors" "os" "golang.org/x/sys/unix" @@ -12,6 +13,16 @@ func lockFile(f *os.File) error { return unix.Flock(int(f.Fd()), unix.LOCK_EX) } +// tryLockFile is lockFile with LOCK_NB: flock returns EWOULDBLOCK instead of +// sleeping when the lock is held elsewhere. +func tryLockFile(f *os.File) error { + err := unix.Flock(int(f.Fd()), unix.LOCK_EX|unix.LOCK_NB) + if errors.Is(err, unix.EWOULDBLOCK) || errors.Is(err, unix.EAGAIN) { + return ErrLocked + } + return err +} + func unlockFile(f *os.File) error { return unix.Flock(int(f.Fd()), unix.LOCK_UN) } diff --git a/internal/filelock/filelock_windows.go b/internal/filelock/filelock_windows.go index d02a4d8..5bad4f2 100644 --- a/internal/filelock/filelock_windows.go +++ b/internal/filelock/filelock_windows.go @@ -3,16 +3,31 @@ package filelock import ( + "errors" "os" "golang.org/x/sys/windows" ) func lockFile(f *os.File) error { + return lockFileEx(f, windows.LOCKFILE_EXCLUSIVE_LOCK) +} + +// tryLockFile is lockFile with LOCKFILE_FAIL_IMMEDIATELY: LockFileEx returns +// ERROR_LOCK_VIOLATION instead of waiting when the lock is held elsewhere. +func tryLockFile(f *os.File) error { + err := lockFileEx(f, windows.LOCKFILE_EXCLUSIVE_LOCK|windows.LOCKFILE_FAIL_IMMEDIATELY) + if errors.Is(err, windows.ERROR_LOCK_VIOLATION) { + return ErrLocked + } + return err +} + +func lockFileEx(f *os.File, flags uint32) error { var overlapped windows.Overlapped return windows.LockFileEx( windows.Handle(f.Fd()), - windows.LOCKFILE_EXCLUSIVE_LOCK, + flags, 0, 1, 0, diff --git a/internal/routes/cache.go b/internal/routes/cache.go index d70c3a1..0b5d0b4 100644 --- a/internal/routes/cache.go +++ b/internal/routes/cache.go @@ -56,9 +56,12 @@ type LoadRegistryOptions struct { // when authenticated. Do NOT pass the raw sa_live_ token here — it must // be exchanged first via the client's EnsureAccessToken. AccessToken string - // CacheKey is an opaque string used to isolate cached specs per identity - // (typically derived from the sa_live_ token). When set, specs are cached - // in a key-specific subdirectory to avoid cross-account contamination. + // CacheKey is an opaque string that isolates cached specs per identity, so + // one account's plan-filtered spec is never read back for another. When + // set, specs live in a subdirectory named by its hash. It must be stable + // for as long as the identity is: a service-account token qualifies, and + // so does a session JWT's jti — but not the JWT itself, whose bytes change + // every time the issuer re-signs it. CacheKey string // ForceRefresh bypasses TTL and re-downloads with ETag validation. ForceRefresh bool @@ -68,6 +71,28 @@ type LoadRegistryOptions struct { TTL time.Duration // HTTPClient overrides the HTTP client (for testing). HTTPClient *http.Client + // NonBlocking makes EnsureSpecs give up with ErrCacheBusy when another + // process holds the cache lock, instead of waiting for it. A caller in + // front of a request wants this: if someone is already downloading, the + // right move is to step aside, not to queue behind them for however long + // their download takes. + NonBlocking bool + // Quiet keeps non-fatal cache warnings off stderr. Commands whose stderr + // is a JSON contract cannot have prose appear there. + Quiet bool +} + +// ErrCacheBusy reports that another process holds the spec cache lock. It is +// returned only when NonBlocking is set; otherwise EnsureSpecs waits. +var ErrCacheBusy = errors.New("spec cache: another process is refreshing it") + +// warnf reports a non-fatal cache problem on stderr unless the caller asked +// for quiet. +func (o *LoadRegistryOptions) warnf(format string, args ...any) { + if o.Quiet { + return + } + fmt.Fprintf(os.Stderr, format, args...) } // specMeta holds per-spec cache metadata. @@ -101,16 +126,16 @@ func (o *LoadRegistryOptions) resolveCacheDir() (string, error) { base = filepath.Join(home, ".cio", specCacheSubdir) } if o.CacheKey != "" { - return filepath.Join(base, tokenCacheKey(o.CacheKey)), nil + return filepath.Join(base, cacheKeyDir(o.CacheKey)), nil } return base, nil } -// tokenCacheKey returns a short, filesystem-safe hash of a token for use -// as a cache subdirectory name. Different tokens produce different keys -// so that personalized specs don't collide. -func tokenCacheKey(token string) string { - h := sha256.Sum256([]byte(token)) +// cacheKeyDir returns a short, filesystem-safe subdirectory name for a cache +// key: a hash, so a secret used as the key never lands on disk in the clear, +// and distinct keys never share a directory. +func cacheKeyDir(key string) string { + h := sha256.Sum256([]byte(key)) return "auth-" + hex.EncodeToString(h[:8]) } @@ -133,12 +158,24 @@ func (o *LoadRegistryOptions) resolveHTTPClient() *http.Client { return &http.Client{Timeout: 30 * time.Second} } -// lockCacheDir acquires an exclusive file lock on the cache directory. -// Returns an unlock function that must be called when done. -func lockCacheDir(cacheDir string) (unlock func(), err error) { +// lockCacheDir acquires an exclusive file lock on the cache directory and +// returns an unlock function that must be called when done. With nonBlocking +// set, a lock held elsewhere yields ErrCacheBusy at once instead of a wait. +// +// The lock file doubles as a download lease with no expiry to manage: the +// kernel releases a flock when its holder exits, so a downloader that dies +// mid-fetch never leaves the cache locked. +func lockCacheDir(cacheDir string, nonBlocking bool) (unlock func(), err error) { lockPath := filepath.Join(cacheDir, "specs.lock") - unlock, err = filelock.Lock(lockPath, 0600) + if nonBlocking { + unlock, err = filelock.TryLock(lockPath, 0600) + } else { + unlock, err = filelock.Lock(lockPath, 0600) + } if err != nil { + if errors.Is(err, filelock.ErrLocked) { + return nil, ErrCacheBusy + } return nil, fmt.Errorf("acquire lock: %w", err) } @@ -155,7 +192,7 @@ func EnsureSpecs(ctx context.Context, opts LoadRegistryOptions) (journeys, cdp [ return nil, nil, fmt.Errorf("create cache dir: %w", err) } - unlock, err := lockCacheDir(cacheDir) + unlock, err := lockCacheDir(cacheDir, opts.NonBlocking) if err != nil { return nil, nil, err } @@ -171,7 +208,7 @@ func EnsureSpecs(ctx context.Context, opts LoadRegistryOptions) (journeys, cdp [ var errs []error for i, src := range defaultSpecSources { - data, changed, specErr := ensureSpec(ctx, httpClient, cacheDir, baseURL, opts.AccessToken, src, meta, ttl, opts.ForceRefresh) + data, changed, specErr := ensureSpec(ctx, httpClient, cacheDir, baseURL, opts.AccessToken, src, meta, ttl, opts.ForceRefresh, opts.warnf) if specErr != nil { errs = append(errs, fmt.Errorf("%s: %w", src.Name, specErr)) continue @@ -184,7 +221,7 @@ func EnsureSpecs(ctx context.Context, opts LoadRegistryOptions) (journeys, cdp [ if updated { if err := writeMeta(cacheDir, meta); err != nil { - fmt.Fprintf(os.Stderr, "warning: failed to write spec cache metadata: %v\n", err) + opts.warnf("warning: failed to write spec cache metadata: %v\n", err) } } @@ -205,6 +242,7 @@ func ensureSpec( meta *cacheMeta, ttl time.Duration, forceRefresh bool, + warnf func(format string, args ...any), ) (data []byte, changed bool, err error) { filename := src.Name + ".json" cachedPath := filepath.Join(cacheDir, filename) @@ -236,7 +274,7 @@ func ensureSpec( // Try stale cache on download failure. data, readErr := os.ReadFile(cachedPath) if readErr == nil { - fmt.Fprintf(os.Stderr, "warning: using stale cached %s (download failed: %v)\n", src.Name, dlErr) + warnf("warning: using stale cached %s (download failed: %v)\n", src.Name, dlErr) return data, false, nil } return nil, false, dlErr diff --git a/internal/routes/cache_nonblocking_test.go b/internal/routes/cache_nonblocking_test.go new file mode 100644 index 0000000..25831c5 --- /dev/null +++ b/internal/routes/cache_nonblocking_test.go @@ -0,0 +1,173 @@ +package routes + +import ( + "context" + "errors" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync/atomic" + "testing" + "time" + + "github.com/customerio/cli/internal/filelock" +) + +const minimalSpec = `{"openapi":"3.1.0","info":{"title":"T","version":"1"},"paths":{"/v1/environments/{environment_id}/campaigns":{"get":{}}}}` + +// countingSpecServer serves a minimal spec and counts how many times it was asked. +func countingSpecServer(t *testing.T) (*httptest.Server, *atomic.Int32) { + t.Helper() + var hits atomic.Int32 + s := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + hits.Add(1) + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(minimalSpec)) + })) + t.Cleanup(s.Close) + return s, &hits +} + +// NonBlocking: a lock held by someone else means ErrCacheBusy now — not a wait, +// and not a download. +func TestEnsureSpecs_NonBlockingStepsAsideWhenLockHeld(t *testing.T) { + dir := t.TempDir() + server, hits := countingSpecServer(t) + + // Simulate another process mid-download. + unlock, err := filelock.Lock(filepath.Join(dir, "specs.lock"), 0600) + if err != nil { + t.Fatalf("holding lock: %v", err) + } + defer unlock() + + start := time.Now() + _, _, err = EnsureSpecs(context.Background(), LoadRegistryOptions{ + CacheDir: dir, BaseURL: server.URL, NonBlocking: true, + }) + elapsed := time.Since(start) + + if !errors.Is(err, ErrCacheBusy) { + t.Fatalf("want ErrCacheBusy, got %v", err) + } + if hits.Load() != 0 { + t.Errorf("must not download while another process holds the lock, got %d fetches", hits.Load()) + } + if elapsed > 2*time.Second { + t.Errorf("took %v; a held lock must be reported at once, not waited on", elapsed) + } +} + +// The default (blocking) behaviour is untouched: the same held lock is waited +// for, and released, and the download then proceeds. +func TestEnsureSpecs_BlockingStillWaitsForLock(t *testing.T) { + dir := t.TempDir() + server, hits := countingSpecServer(t) + + unlock, err := filelock.Lock(filepath.Join(dir, "specs.lock"), 0600) + if err != nil { + t.Fatalf("holding lock: %v", err) + } + go func() { + time.Sleep(150 * time.Millisecond) + unlock() + }() + + if _, _, err := EnsureSpecs(context.Background(), LoadRegistryOptions{CacheDir: dir, BaseURL: server.URL}); err != nil { + t.Fatalf("blocking EnsureSpecs after the holder released: %v", err) + } + if hits.Load() == 0 { + t.Error("expected the download to proceed once the lock was released") + } +} + +// The caller's deadline bounds the download. +func TestEnsureSpecs_HonoursContextDeadline(t *testing.T) { + dir := t.TempDir() + release := make(chan struct{}) + defer close(release) + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + select { + case <-release: + case <-r.Context().Done(): + } + })) + defer server.Close() + + ctx, cancel := context.WithTimeout(context.Background(), 200*time.Millisecond) + defer cancel() + + start := time.Now() + _, _, err := EnsureSpecs(ctx, LoadRegistryOptions{CacheDir: dir, BaseURL: server.URL, NonBlocking: true}) + elapsed := time.Since(start) + + if err == nil { + t.Fatal("expected an error from a stalled host under a 200ms deadline") + } + if elapsed > 3*time.Second { + t.Errorf("took %v; the 200ms deadline was not honoured", elapsed) + } +} + +// Quiet keeps the stale-cache warning off stderr; without it the warning is +// written as before. +func TestEnsureSpecs_QuietSuppressesWarnings(t *testing.T) { + for _, quiet := range []bool{false, true} { + t.Run(map[bool]string{false: "loud", true: "quiet"}[quiet], func(t *testing.T) { + dir := t.TempDir() + good, _ := countingSpecServer(t) + if _, _, err := EnsureSpecs(context.Background(), LoadRegistryOptions{CacheDir: dir, BaseURL: good.URL}); err != nil { + t.Fatalf("priming cache: %v", err) + } + + // A host that now fails, and a TTL that forces a revalidation + // attempt, drives the "using stale cached" warning path. + bad := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "down", http.StatusInternalServerError) + })) + defer bad.Close() + + stderr := captureStderr(t, func() { + _, _, err := EnsureSpecs(context.Background(), LoadRegistryOptions{ + CacheDir: dir, BaseURL: bad.URL, TTL: time.Nanosecond, Quiet: quiet, + }) + if err != nil { + t.Fatalf("stale fallback should succeed: %v", err) + } + }) + + if quiet && strings.Contains(stderr, "warning:") { + t.Errorf("Quiet must suppress cache warnings, got stderr: %q", stderr) + } + if !quiet && !strings.Contains(stderr, "warning: using stale cached") { + t.Errorf("without Quiet the warning must still be written, got stderr: %q", stderr) + } + }) + } +} + +// captureStderr runs fn with os.Stderr redirected to a pipe and returns what +// was written. +func captureStderr(t *testing.T, fn func()) string { + t.Helper() + r, w, err := os.Pipe() + if err != nil { + t.Fatalf("pipe: %v", err) + } + old := os.Stderr + os.Stderr = w + defer func() { os.Stderr = old }() + + done := make(chan string) + go func() { + b, _ := io.ReadAll(r) + done <- string(b) + }() + + fn() + _ = w.Close() + return <-done +} diff --git a/internal/routes/cache_test.go b/internal/routes/cache_test.go index 7afd5e9..f52a1ff 100644 --- a/internal/routes/cache_test.go +++ b/internal/routes/cache_test.go @@ -470,16 +470,16 @@ func TestEnsureSpecs_AuthenticatedCacheIsolation(t *testing.T) { t.Errorf("unauthenticated cache missing: %v", err) } - // Authenticated specs should be in a token-specific subdir. - authDir := filepath.Join(cacheDir, tokenCacheKey(token)) + // Authenticated specs should be in a key-specific subdir. + authDir := filepath.Join(cacheDir, cacheKeyDir(token)) if _, err := os.Stat(filepath.Join(authDir, "openapi.json")); err != nil { t.Errorf("authenticated cache missing: %v", err) } - // Different token should use a different subdir. + // Different keys should use different subdirs. token2 := "test-service-account-token-b" - if tokenCacheKey(token) == tokenCacheKey(token2) { - t.Error("different tokens should produce different cache keys") + if cacheKeyDir(token) == cacheKeyDir(token2) { + t.Error("different keys should produce different cache directories") } } diff --git a/internal/routes/pathindex.go b/internal/routes/pathindex.go index 9bd3974..c8334c9 100644 --- a/internal/routes/pathindex.go +++ b/internal/routes/pathindex.go @@ -1,6 +1,7 @@ package routes import ( + "context" "encoding/json" "fmt" "os" @@ -47,21 +48,14 @@ var pathItemVerbs = map[string]bool{ // LoadPathIndexFromCache decodes the path templates out of whatever specs the // cache already holds. It reads files and nothing else: no download, no cache -// lock, no metadata. -// -// EnsureSpecs, which LoadRegistry goes through, is the wrong tool in front of a -// request. It takes an exclusive lock on the cache directory before reading -// anything, and that lock is a plain flock — it cannot be bounded or -// cancelled — so concurrent callers serialize behind whichever one holds it, -// and the holder may sit through two sequential spec downloads first. A check -// worth ~25ms must not be able to cost a minute, nor make parallel calls queue. +// lock, no metadata — the common, ~25ms case for a check that runs before +// every request. // // Reading without the lock is safe because cached specs are replaced by rename, // so a reader sees one whole version or the previous one, never a torn file. -// The cost is that this reports what the cache knows rather than what the -// server currently serves: a caller that needs freshness (`cio schema`) keeps -// going through EnsureSpecs, and a caller checking a path treats a cold cache -// as "no opinion". +// It reports what the cache knows rather than what the server serves; a caller +// that needs freshness (`cio schema`) goes through EnsureSpecs, and a caller +// that finds the cache cold uses EnsurePathIndex to fill it, once. func LoadPathIndexFromCache(opts LoadRegistryOptions) (*PathIndex, error) { cacheDir, err := opts.resolveCacheDir() if err != nil { @@ -86,6 +80,41 @@ func LoadPathIndexFromCache(opts LoadRegistryOptions) (*PathIndex, error) { return idx, nil } +// EnsurePathIndex fills a cold cache and indexes it. It goes through +// EnsureSpecs — the one download path, so what it writes is exactly the spec +// `cio schema` would have fetched for this identity — but in the bounded form a +// caller in front of a request can afford: +// +// - NonBlocking: another process already downloading means step aside at +// once (ErrCacheBusy), never queue behind it. +// - The caller's context deadline bounds the download itself. +// - Quiet: cache warnings stay off stderr, which the calling command keeps +// for structured errors. +// +// Any error means the caller has no spec to judge with, and should proceed +// without an opinion rather than fail. +func EnsurePathIndex(ctx context.Context, opts LoadRegistryOptions) (*PathIndex, error) { + opts.NonBlocking = true + opts.Quiet = true + + journeysData, cdpData, err := EnsureSpecs(ctx, opts) + if err != nil { + return nil, err + } + + idx := &PathIndex{} + if err := idx.add(journeysData); err != nil { + return nil, err + } + if len(cdpData) > 0 { + _ = idx.add(cdpData) + } + if len(idx.entries) == 0 { + return nil, fmt.Errorf("API spec described no routes") + } + return idx, nil +} + // add indexes one spec document, accepting either an OpenAPI document or the // walked-routes JSON that LoadRegistryFromData reads. func (idx *PathIndex) add(data []byte) error {