Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
496bf88
migrations: 0008 group scopes — billing_event.member_group_ids + grou…
hhuuggoo Oct 7, 2026
1a17764
identity: X-Saturn-Group-Scopes joins the pinned trusted set; members…
hhuuggoo Oct 7, 2026
911ff22
metering/drain: carry member_group_ids to billing_event (NULL when em…
hhuuggoo Oct 7, 2026
81e0ed5
rating: upsert the group_usage attribution rollup while rating a window
hhuuggoo Oct 7, 2026
c67ecf0
admission: group rate scopes + monthly spend caps (membership-aware q…
hhuuggoo Oct 7, 2026
48b9b62
proxy: parse the group scope envelope and enforce it at admission
hhuuggoo Oct 7, 2026
5452f4c
interceptor: wire the group spend store from DATABASE_URL
hhuuggoo Oct 7, 2026
6889d1c
admission: bound the group spend check; proxy: enforce the frozen env…
hhuuggoo Oct 7, 2026
1976681
admission: collapse the spend-check stampede after a store error
hhuuggoo Oct 7, 2026
2864330
admission: one spend-query budget per request; proxy: skip admission for
hhuuggoo Oct 7, 2026
f24050e
tests: prove the group-attribution gates and spend-check boundaries; …
hhuuggoo Oct 7, 2026
b3a7924
admission: a starved spend read is not a store outage; tests: e2e group
hhuuggoo Oct 7, 2026
f784bae
proxy: pin the group generated-token cap in the undeclared-max_tokens…
hhuuggoo Oct 7, 2026
accf4d6
interceptor: test the buildAdmission DATABASE_URL spend-store wiring
hhuuggoo Oct 7, 2026
3b0dbe8
admission: rename cap parameter to spendCap (revive redefines-builtin…
hhuuggoo Oct 7, 2026
354034c
admission: rename test cap locals to spendCap (revive redefines-built…
hhuuggoo Oct 7, 2026
edada95
admission: rename remaining cap identifiers in integration tests (rev…
hhuuggoo Oct 7, 2026
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
27 changes: 25 additions & 2 deletions cmd/interceptor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ func main() {

// R3 trusted-header registry: parse PHOEBE_TRUSTED_HEADERS (rendered
// from the phoebe chart's ConfigMap) into the active set the parser's
// envelope reads resolve through. The pinned 13 already govern from
// envelope reads resolve through. The pinned 14 already govern from
// package init; this engages the runtime config once, at startup, and
// warns loudly if the chart render was empty/malformed.
identity.LoadTrustedHeaders(log)
Expand Down Expand Up @@ -100,6 +100,26 @@ func buildAdmission(s *config.Settings, log *logging.Logger) (admission.Admitter
// docs/shared-tier-admission.md).
log.Info.Printf("admission: admission.valkeyAddr is empty; using the metering Valkey from emit.valkeyAddr (%s) as the admission store", cfg.ValkeyAddr)
}
// The monthly group spend cap reads group_usage in Postgres (the rater's
// attribution rollup), NOT the admission Valkey. Without a DATABASE_URL
// (a serving-only spoke install runs no Postgres at all) the cap cannot
// be checked: group rate limits still enforce, the spend cap is bypassed
// fail-open — said loudly here, never assumed silently.
closeSpend := func() {}
if dsn := os.Getenv("DATABASE_URL"); dsn != "" {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
store, err := admission.OpenPostgresSpendStore(ctx, dsn)
cancel()
if err != nil {
log.Error.Printf("admission: group spend caps DISABLED (spend-check store could not be built: %v); group rate limits still enforce", err)
} else {
admitter.WithGroupSpend(store, log)
closeSpend = func() { _ = store.Close() }
log.Info.Printf("admission: group spend caps enabled (spend check reads group_usage in Postgres)")
}
} else {
log.Error.Printf("admission: DATABASE_URL is unset; group spend caps are NOT enforced (group rate limits still enforce)")
}
mode := "contract limits only (admission.enabled=false: operator capacity tiers off, envelope-less shared routes allowed)"
if cfg.Enabled {
mode = "contract limits + operator capacity tiers (admission.enabled=true: every shared request must carry the trusted envelope)"
Expand All @@ -114,7 +134,10 @@ func buildAdmission(s *config.Settings, log *logging.Logger) (admission.Admitter
} else {
log.Info.Printf("admission: enforcing %s (valkey %s, lease ttl %s)", mode, cfg.ValkeyAddr, cfg.LeaseTTL)
}
return admitter, func() { _ = client.Close() }
return admitter, func() {
closeSpend()
_ = client.Close()
}
}

// buildGateway constructs the TF gateway (org, model) resolver. DEFAULT: the
Expand Down
137 changes: 137 additions & 0 deletions cmd/interceptor/main_integration_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
//go:build integration

// Package main integration test: runs buildAdmission's DATABASE_URL spend
// wiring against a LIVE Postgres loaded with the production migrations, so the
// "group spend caps enabled" startup line is proven to describe a working
// check — the unit test pins the two failure logs, but only this half
// exercises OpenPostgresSpendStore + WithGroupSpend together.
//
// Gated behind the `integration` build tag AND a non-empty
// PHOEBE_TEST_DATABASE_URL. Run with:
//
// PHOEBE_TEST_DATABASE_URL=postgres://... go test -tags=integration ./cmd/interceptor/...
package main

import (
"bytes"
"context"
"database/sql"
"errors"
"fmt"
"io"
"log"
"os"
"sort"
"strings"
"testing"
"time"

"github.com/alicebob/miniredis/v2"
_ "github.com/jackc/pgx/v5/stdlib"

"github.com/saturncloud/phoebe/internal/admission"
"github.com/saturncloud/phoebe/internal/logging"
"github.com/saturncloud/phoebe/migrations"
)

// newSpendWiringHarness creates an isolated schema — named per test PROCESS so
// concurrent runs against a shared database stop stomping each other — with
// ALL production migrations applied (discovered from the embedded
// migrations.FS, ordered by filename; a hand-maintained list could go stale
// while green), and returns a pool pinned to it plus the schema-pinned DSN for
// handing to buildAdmission as DATABASE_URL.
func newSpendWiringHarness(t *testing.T) (*sql.DB, string) {
t.Helper()
dsn := os.Getenv("PHOEBE_TEST_DATABASE_URL")
if dsn == "" {
t.Skip("PHOEBE_TEST_DATABASE_URL not set; skipping interceptor spend-wiring integration test")
}
schema := fmt.Sprintf("phoebe_interceptor_spend_it_%d", os.Getpid())

admin, err := sql.Open("pgx", dsn)
if err != nil {
t.Fatalf("open admin pool: %v", err)
}
t.Cleanup(func() {
_, _ = admin.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE")
_ = admin.Close()
})
execSpendWiring(t, admin, "DROP SCHEMA IF EXISTS "+schema+" CASCADE")
execSpendWiring(t, admin, "CREATE SCHEMA "+schema)

sep := "?"
if strings.Contains(dsn, "?") {
sep = "&"
}
pinned := dsn + sep + "search_path=" + schema
db, err := sql.Open("pgx", pinned)
if err != nil {
t.Fatalf("open schema pool: %v", err)
}
t.Cleanup(func() { _ = db.Close() })

entries, err := migrations.FS.ReadDir(".")
if err != nil {
t.Fatalf("read migrations FS: %v", err)
}
names := make([]string, 0, len(entries))
for _, e := range entries {
if !e.IsDir() && strings.HasSuffix(e.Name(), ".up.sql") {
names = append(names, e.Name())
}
}
sort.Strings(names)
for _, name := range names {
b, err := migrations.FS.ReadFile(name)
if err != nil {
t.Fatalf("read migration %s: %v", err, name)
}
execSpendWiring(t, db, string(b))
}
return db, pinned
}

func execSpendWiring(t *testing.T, db *sql.DB, stmt string) {
t.Helper()
if _, err := db.Exec(stmt); err != nil {
t.Fatalf("exec failed: %v\nstatement: %s", err, stmt)
}
}

// TestIntegration_BuildAdmissionGroupSpendCapsEnabled is the positive half of
// the DATABASE_URL wiring: against live Postgres, buildAdmission logs the
// "enabled" line AND the admitter it returns enforces the monthly spend cap —
// a spend-capped Admit whose group_usage spend reached the cap is the
// contractual monthly_spend rejection. A regression deleting the
// WithGroupSpend call admits this request, so the wiring's production switch
// is now mutation-covered, not just logged.
func TestIntegration_BuildAdmissionGroupSpendCapsEnabled(t *testing.T) {
db, pinned := newSpendWiringHarness(t)
t.Setenv("DATABASE_URL", pinned)
mr := miniredis.RunT(t)
var buf bytes.Buffer
logger := &logging.Logger{Debug: log.New(io.Discard, "", 0), Info: log.New(&buf, "", 0), Warn: log.New(&buf, "", 0), Error: log.New(&buf, "", 0)}
admitter, closeAdmission := buildAdmission(loadTestSettings(t, "emit:\n valkeyAddr: "+mr.Addr()+"\n"), logger)
defer closeAdmission()
if admitter == nil {
t.Fatal("no admitter built")
}
if out := buf.String(); !strings.Contains(out, "group spend caps enabled") {
t.Fatalf("startup log missing the enabled line:\n%s", out)
}

// The group's month-to-date attribution spend is 10 — over the cap of 1 —
// in the current UTC hour bucket, so the month predicate cannot skip it.
_, err := db.Exec(`INSERT INTO group_usage (group_id, window_start, cost, event_count)
VALUES ($1, $2, '10', 1)`, spendWiringGID, time.Now().UTC().Truncate(time.Hour))
if err != nil {
t.Fatalf("seed group_usage: %v", err)
}

req := spendWiringRequest(admission.GroupScope{GroupID: spendWiringGID, SpendCap: "1"})
_, err = admitter.Admit(context.Background(), req)
var rejected *admission.Rejected
if !errors.As(err, &rejected) || !rejected.Contractual || rejected.Dimension != "monthly_spend" {
t.Fatalf("spend-capped admit = %v, want a contractual monthly_spend rejection (the wired spend check reads group_usage)", err)
}
}
79 changes: 79 additions & 0 deletions cmd/interceptor/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package main
import (
"bytes"
"context"
"errors"
"io"
"log"
"os"
Expand Down Expand Up @@ -91,3 +92,81 @@ func TestBuildAdmissionLogsMeteringStoreFallback(t *testing.T) {
t.Fatalf("fallback logged although admission.valkeyAddr is set: %q", explicit)
}
}

const spendWiringGID = "a1b2c3d4e5f60718293a4b5c6d7e8f90"

// spendWiringRequest is a shared-shaped admit request carrying one group
// scope; scope carries either the rate limit or the spend cap under test.
func spendWiringRequest(scope admission.GroupScope) admission.Request {
return admission.Request{
Graph: "graph-a", Organization: "org-a", Owner: "owner-a", Model: "m",
PromptBytes: 10, EstimatedInputTokens: 10, ReservedOutputTokens: 20,
GroupScopes: []admission.GroupScope{scope},
}
}

// TestBuildAdmissionGroupSpendCapsDatabaseURL pins the production DATABASE_URL
// switch that turns the monthly group spend cap on. This wiring is the
// feature's only production call site, and every other spend test wires the
// store by hand — so a regression here (WithGroupSpend not called on a
// successful open, or the open error made fatal) disables group spend caps on
// every interceptor, or crash-loops it, with the rest of the suite green.
func TestBuildAdmissionGroupSpendCapsDatabaseURL(t *testing.T) {
build := func(t *testing.T) (admission.Admitter, *bytes.Buffer) {
t.Helper()
mr := miniredis.RunT(t)
var buf bytes.Buffer
logger := &logging.Logger{Debug: log.New(io.Discard, "", 0), Info: log.New(&buf, "", 0), Warn: log.New(&buf, "", 0), Error: log.New(&buf, "", 0)}
admitter, closeAdmission := buildAdmission(loadTestSettings(t, "emit:\n valkeyAddr: "+mr.Addr()+"\n"), logger)
t.Cleanup(closeAdmission)
return admitter, &buf
}

t.Run("unset logs NOT enforced and group rate limits still enforce", func(t *testing.T) {
t.Setenv("DATABASE_URL", "")
admitter, buf := build(t)
if admitter == nil {
t.Fatal("no admitter with DATABASE_URL unset: contract and group rate limits would not be enforced")
}
if out := buf.String(); !strings.Contains(out, "group spend caps are NOT enforced") {
t.Fatalf("startup log missing the loud NOT-enforced line:\n%s", out)
}

// The startup line's promise, pinned behaviorally: group rate limits
// still enforce — a capped window rejects the second request with the
// contractual 429 mapping.
one := int64(1)
rateScope := admission.GroupScope{GroupID: spendWiringGID, Limits: admission.RateLimits{Requests: &one}}
lease, err := admitter.Admit(context.Background(), spendWiringRequest(rateScope))
if err != nil {
t.Fatalf("first group-rate admit: %v, want admitted", err)
}
_ = lease.Complete(context.Background(), 0)
_, err = admitter.Admit(context.Background(), spendWiringRequest(rateScope))
var rejected *admission.Rejected
if !errors.As(err, &rejected) || !rejected.Contractual {
t.Fatalf("second group-rate admit: %v, want a contractual rejection (group rate limits still enforce without DATABASE_URL)", err)
}

// ... and the spend cap does not: a spend-capped group admits
// fail-open with no store behind it.
cappedScope := admission.GroupScope{GroupID: spendWiringGID, SpendCap: "0"}
if _, err := admitter.Admit(context.Background(), spendWiringRequest(cappedScope)); err != nil {
t.Fatalf("spend-capped admit with no store: %v, want admitted (spend caps are NOT enforced without DATABASE_URL)", err)
}
})

t.Run("unroutable DSN logs DISABLED and serving continues", func(t *testing.T) {
// 127.0.0.1:1 answers with connection refused, so the 5s open budget
// is not what bounds this test.
t.Setenv("DATABASE_URL", "postgres://postgres:pw@127.0.0.1:1/postgres?sslmode=disable")
admitter, buf := build(t)
if admitter == nil {
t.Fatal("no admitter with an unroutable DATABASE_URL: serving must continue without spend caps")
}
out := buf.String()
if !strings.Contains(out, "group spend caps DISABLED") || !strings.Contains(out, "group rate limits still enforce") {
t.Fatalf("startup log missing the loud DISABLED line naming the surviving rate limits:\n%s", out)
}
})
}
Loading
Loading