Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 40 additions & 7 deletions cmd/api.go
Original file line number Diff line number Diff line change
@@ -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"
Expand Down Expand Up @@ -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.
//
Expand All @@ -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
Expand Down
145 changes: 145 additions & 0 deletions cmd/api_preflight_cachekey_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading