From 3e3f7f72aab6f1043717f7ad224cf590afc71637 Mon Sep 17 00:00:00 2001 From: mintaka Date: Fri, 2 Oct 2026 23:30:07 -0400 Subject: [PATCH 01/10] feat(usage): record compute-usage intervals with the session binding (RIG-2872) Each session binding is one billable compute interval. Binding writes now emit start and end rows into an append-only compute_usage_events table, in the same transaction as the binding change. ### Changes - Migration 0004_compute_usage.sql adds compute_usage_events with tenant RLS and grants, and backfills an estimated start for every binding already open. - RecordSessionBinding, DeleteSessionBinding and DeleteSessionBindingsForRunner write interval events atomically. A same-session runner change keeps its interval. - A system-role hourly sweep closes an interval whose binding vanished without an end event, marked estimated. - docs/concepts/tokens-and-billing.md describes the event log. Stop, a lost-session drop, and the Runner re-enroll reap all delete the binding, so each ends the interval. There is no pause path in the protocol today. ### Verification - pgtest race: ./internal/store, ./internal/usage, ./internal/runnerhub, ./server all ok. - sql-migration-gate (squawk + sqruff) and sqlc drift pass; golangci-lint 0 issues. Refs RIG-2872 Co-authored-by: Matt Wilkinson --- docs/concepts/tokens-and-billing.md | 11 +- .../backfill_empty_tenant_pgtest_test.go | 2 +- go/internal/store/compute_usage.go | 16 ++ .../store/compute_usage_pgtest_test.go | 236 ++++++++++++++++++ go/internal/store/db/compute_usage.sql.go | 51 ++++ go/internal/store/db/models.go | 26 +- go/internal/store/db/querier.go | 43 ++-- go/internal/store/db/session_bindings.sql.go | 170 ++++++++++--- .../store/migrate_upgrade_pgtest_test.go | 58 ++++- .../store/migrations/0004_compute_usage.sql | 60 +++++ go/internal/store/queries/compute_usage.sql | 31 +++ .../store/queries/session_bindings.sql | 106 +++++--- go/internal/store/rls_pgtest_test.go | 1 + go/internal/store/session_bindings.go | 50 ++-- go/internal/store/tenant_tx.go | 17 +- go/internal/usage/compute_sweep.go | 54 ++++ go/internal/usage/compute_sweep_test.go | 32 +++ go/server/serve.go | 3 + go/server/sinks.go | 12 + 19 files changed, 846 insertions(+), 133 deletions(-) create mode 100644 go/internal/store/compute_usage.go create mode 100644 go/internal/store/compute_usage_pgtest_test.go create mode 100644 go/internal/store/db/compute_usage.sql.go create mode 100644 go/internal/store/migrations/0004_compute_usage.sql create mode 100644 go/internal/store/queries/compute_usage.sql create mode 100644 go/internal/usage/compute_sweep.go create mode 100644 go/internal/usage/compute_sweep_test.go diff --git a/docs/concepts/tokens-and-billing.md b/docs/concepts/tokens-and-billing.md index a6ae952e5..482acff20 100644 --- a/docs/concepts/tokens-and-billing.md +++ b/docs/concepts/tokens-and-billing.md @@ -52,11 +52,12 @@ Three distinct things are recorded off the OMP gateway and the runtime. Keeping them separate matters: one is billing-grade, one is display, one is a health/quota signal. -- **Manager/agent compute usage — billing-grade on managed.** The metered - activity the managed service caps and charges overage on. Because it backs - billing, it must be exact, auditable, and reconstructable — an append-only - event record, not a lossy counter (the durable-event-log rationale is the obs - record's Decision D5). +- **Manager/agent compute usage — billing-grade on managed.** Each session + binding creates one append-only interval, billed as agent-active seconds from + its start to end. The server emits both events with the binding transition; + an hourly sweep marks an inferred end as estimated if a binding disappears + without its end event. Intervals already open when the log was introduced get + an estimated start. The event log is exact, auditable, and reconstructable. - **LLM token usage and spend — recorded for display, not billed day-1.** Every model call's tokens-in/out and cost, captured at the gateway. It powers the in-product usage/spend charts the user sees, and it is *recorded* even though diff --git a/go/internal/store/backfill_empty_tenant_pgtest_test.go b/go/internal/store/backfill_empty_tenant_pgtest_test.go index 3dd279ebe..be01dd547 100644 --- a/go/internal/store/backfill_empty_tenant_pgtest_test.go +++ b/go/internal/store/backfill_empty_tenant_pgtest_test.go @@ -24,7 +24,7 @@ func TestBackfillEmptyTenantMigration(t *testing.T) { } execAsSystem(t, s, "INSERT INTO agent_sessions (session_id, agent_account_id, tenant_id) VALUES ('sess-lone', $1, ''), ('sess-twice-old', $2, '')", string(lone.ID), string(twice.ID)) - execAsSystem(t, s, "INSERT INTO session_bindings (agent_account_id, session_id, runner_id, tenant_id) VALUES ($1, 'sess-lone', 'runner-1', ''), ($2, 'sess-twice-old', 'runner-1', '')", + execAsSystem(t, s, "INSERT INTO session_bindings (agent_account_id, session_id, runner_id, tenant_id, usage_interval_id) VALUES ($1, 'sess-lone', 'runner-1', '', 'iv-lone'), ($2, 'sess-twice-old', 'runner-1', '', 'iv-twice-old')", string(lone.ID), string(twice.ID)) execAsSystem(t, s, migrationSQL(t, 2)) diff --git a/go/internal/store/compute_usage.go b/go/internal/store/compute_usage.go new file mode 100644 index 000000000..98ca95f80 --- /dev/null +++ b/go/internal/store/compute_usage.go @@ -0,0 +1,16 @@ +package store + +import ( + "context" + "fmt" +) + +// CloseOrphanedComputeIntervals closes intervals whose binding disappeared +// without its matching end event. The system role lets one pass cover all tenants. +func (s *Store) CloseOrphanedComputeIntervals(ctx context.Context) (int64, error) { + closed, err := s.q.CloseOrphanedComputeIntervals(ctx) + if err != nil { + return 0, fmt.Errorf("store: close orphaned compute intervals: %w", err) + } + return closed, nil +} diff --git a/go/internal/store/compute_usage_pgtest_test.go b/go/internal/store/compute_usage_pgtest_test.go new file mode 100644 index 000000000..bda2d99d1 --- /dev/null +++ b/go/internal/store/compute_usage_pgtest_test.go @@ -0,0 +1,236 @@ +//go:build pgtest + +package store + +import ( + "context" + "testing" + + "github.com/jackc/pgx/v5" +) + +type computeUsageEvent struct { + TenantID string + ID string + IntervalID string + Kind string + AgentAccountID string + OwnerUserID string + SessionID string + RunnerID string + Estimated bool +} + +func computeEvents(t *testing.T, s *Store, tenant TenantID, accountID AccountID) []computeUsageEvent { + t.Helper() + rows, err := s.pool.Query(t.Context(), ` + SELECT tenant_id, id, interval_id, kind, agent_account_id, owner_user_id, + session_id, runner_id, estimated + FROM compute_usage_events + WHERE tenant_id = $1 AND agent_account_id = $2 + ORDER BY session_id, CASE kind WHEN 'start' THEN 0 ELSE 1 END`, string(tenant), string(accountID)) + if err != nil { + t.Fatalf("read compute usage events: %v", err) + } + defer rows.Close() + events := []computeUsageEvent{} + for rows.Next() { + var event computeUsageEvent + if err := rows.Scan(&event.TenantID, &event.ID, &event.IntervalID, &event.Kind, + &event.AgentAccountID, &event.OwnerUserID, &event.SessionID, &event.RunnerID, + &event.Estimated); err != nil { + t.Fatalf("scan compute usage event: %v", err) + } + events = append(events, event) + } + if err := rows.Err(); err != nil { + t.Fatalf("iterate compute usage events: %v", err) + } + return events +} + +func TestComputeUsageBindingLifecycle(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-owner") + agent := mustAgent(t, s, owner.ID, "compute-agent") + tenant := s.EffectiveTenant(ctx) + + mustBind(t, ctx, s, "compute-session", agent.ID, "compute-runner") + events := computeEvents(t, s, tenant, agent.ID) + if len(events) != 1 { + t.Fatalf("events after first bind = %+v, want one start", events) + } + start := events[0] + if start.Kind != "start" || start.TenantID != string(tenant) || start.OwnerUserID != string(owner.ID) || + start.AgentAccountID != string(agent.ID) || start.SessionID != "compute-session" || + start.RunnerID != "compute-runner" || start.Estimated { + t.Fatalf("start event = %+v, want the tenant, owner, agent, session, runner, and non-estimated start", start) + } + + if err := s.DeleteSessionBinding(ctx, "compute-session"); err != nil { + t.Fatalf("DeleteSessionBinding: %v", err) + } + events = computeEvents(t, s, tenant, agent.ID) + if len(events) != 2 { + t.Fatalf("events after delete = %+v, want start and end", events) + } + var end computeUsageEvent + for _, event := range events { + if event.Kind == "end" { + end = event + } + } + if end.ID == "" || end.IntervalID != start.IntervalID || end.TenantID != string(tenant) || + end.OwnerUserID != string(owner.ID) || end.SessionID != start.SessionID || end.RunnerID != start.RunnerID || end.Estimated { + t.Fatalf("end event = %+v, want same interval metadata and a non-estimated end", end) + } + if err := s.DeleteSessionBinding(ctx, "compute-session"); err != nil { + t.Fatalf("second DeleteSessionBinding: %v", err) + } + if got := computeEvents(t, s, tenant, agent.ID); len(got) != 2 { + t.Fatalf("events after duplicate delete = %+v, want no additional end", got) + } +} + +func TestComputeUsageRepointClosesPriorInterval(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-repoint-owner") + agent := mustAgent(t, s, owner.ID, "compute-repoint-agent") + tenant := s.EffectiveTenant(ctx) + mustBind(t, ctx, s, "session-before", agent.ID, "runner-before") + mustBind(t, ctx, s, "session-after", agent.ID, "runner-after") + + events := computeEvents(t, s, tenant, agent.ID) + if len(events) != 3 { + t.Fatalf("events after re-point = %+v, want prior start/end and replacement start", events) + } + starts := map[string]computeUsageEvent{} + ends := map[string]computeUsageEvent{} + for _, event := range events { + switch event.Kind { + case "start": + starts[event.SessionID] = event + case "end": + ends[event.SessionID] = event + } + } + before, hasBefore := starts["session-before"] + after, hasAfter := starts["session-after"] + closed, hasEnd := ends["session-before"] + if !hasBefore || !hasAfter || !hasEnd || before.IntervalID == after.IntervalID || closed.IntervalID != before.IntervalID { + t.Fatalf("re-point events = %+v, want distinct starts and end of displaced interval", events) + } + if after.RunnerID != "runner-after" || closed.RunnerID != "runner-before" || closed.OwnerUserID != string(owner.ID) { + t.Fatalf("re-point event metadata = %+v, want runner and owner preserved per interval", events) + } +} + +func TestComputeUsageSameSessionRunnerRebindKeepsInterval(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-runner-owner") + agent := mustAgent(t, s, owner.ID, "compute-runner-agent") + tenant := s.EffectiveTenant(ctx) + mustBind(t, ctx, s, "stable-session", agent.ID, "runner-old") + before := computeEvents(t, s, tenant, agent.ID) + mustBind(t, ctx, s, "stable-session", agent.ID, "runner-new") + after := computeEvents(t, s, tenant, agent.ID) + if len(before) != 1 || len(after) != 1 || before[0].Kind != "start" || after[0].IntervalID != before[0].IntervalID { + t.Fatalf("events before=%+v after=%+v, want one unchanged interval start", before, after) + } + if after[0].RunnerID != "runner-old" { + t.Fatalf("interval start runner = %q, want original runner runner-old", after[0].RunnerID) + } +} + +func TestComputeUsageRunnerSweepClosesIntervalsAndReturnsBindings(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-sweep-owner") + agentA := mustAgent(t, s, owner.ID, "compute-sweep-a") + agentB := mustAgent(t, s, owner.ID, "compute-sweep-b") + tenant := s.EffectiveTenant(ctx) + mustBind(t, ctx, s, "sweep-session-a", agentA.ID, "sweep-runner") + mustBind(t, ctx, s, "sweep-session-b", agentB.ID, "sweep-runner") + + swept, err := s.DeleteSessionBindingsForRunner(ctx, "sweep-runner") + if err != nil { + t.Fatalf("DeleteSessionBindingsForRunner: %v", err) + } + if len(swept) != 2 || swept[0].SessionID != "sweep-session-a" || swept[0].AccountID != agentA.ID || + swept[1].SessionID != "sweep-session-b" || swept[1].AccountID != agentB.ID { + t.Fatalf("swept = %+v, want both runner bindings in session order", swept) + } + for _, agent := range []Account{agentA, agentB} { + events := computeEvents(t, s, tenant, agent.ID) + if len(events) != 2 { + t.Fatalf("events for %s after sweep = %+v, want start and end", agent.ID, events) + } + if events[0].IntervalID != events[1].IntervalID || events[0].Kind == events[1].Kind || events[0].Estimated || events[1].Estimated { + t.Fatalf("events for %s = %+v, want matching explicit start/end", agent.ID, events) + } + } +} + +func TestCloseOrphanedComputeIntervals(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-orphan-owner") + orphan := mustAgent(t, s, owner.ID, "compute-orphan-agent") + live := mustAgent(t, s, owner.ID, "compute-live-agent") + tenant := s.EffectiveTenant(ctx) + mustBind(t, ctx, s, "orphan-session", orphan.ID, "orphan-runner") + mustBind(t, ctx, s, "live-session", live.ID, "live-runner") + if _, err := s.pool.Exec(ctx, + "DELETE FROM session_bindings WHERE tenant_id = $1 AND agent_account_id = $2", + string(tenant), string(orphan.ID)); err != nil { + t.Fatalf("delete orphan binding out of band: %v", err) + } + + closed, err := s.CloseOrphanedComputeIntervals(WithSystemRole(ctx)) + if err != nil { + t.Fatalf("CloseOrphanedComputeIntervals: %v", err) + } + if closed != 1 { + t.Fatalf("closed = %d, want one orphan interval", closed) + } + orphanEvents := computeEvents(t, s, tenant, orphan.ID) + if len(orphanEvents) != 2 || orphanEvents[1].Kind != "end" || + orphanEvents[1].IntervalID != orphanEvents[0].IntervalID || !orphanEvents[1].Estimated { + t.Fatalf("orphan events = %+v, want an estimated end for its start", orphanEvents) + } + liveEvents := computeEvents(t, s, tenant, live.ID) + if len(liveEvents) != 1 || liveEvents[0].Kind != "start" { + t.Fatalf("live events = %+v, want open start only", liveEvents) + } + closed, err = s.CloseOrphanedComputeIntervals(WithSystemRole(ctx)) + if err != nil { + t.Fatalf("second CloseOrphanedComputeIntervals: %v", err) + } + if closed != 0 { + t.Fatalf("second close count = %d, want 0", closed) + } +} + +func TestComputeUsageEventsAreTenantIsolated(t *testing.T) { + s := newTestStore(t) + tenantB := seedTenant(t, s, "compute-tenant-b") + ctxA := context.Background() + ctxB := WithTenant(context.Background(), tenantB) + owner := mustUser(t, s, "compute-isolation-owner") + agent := mustAgent(t, s, owner.ID, "compute-isolation-agent") + mustBind(t, ctxA, s, "isolated-session", agent.ID, "isolated-runner") + + var count int64 + err := s.WithTx(ctxB, func(tx pgx.Tx) error { + return tx.QueryRow(ctxB, "SELECT count(*) FROM compute_usage_events").Scan(&count) + }) + if err != nil { + t.Fatalf("read events under tenant B: %v", err) + } + if count != 0 { + t.Fatalf("tenant B sees %d events, want none from tenant A", count) + } +} diff --git a/go/internal/store/db/compute_usage.sql.go b/go/internal/store/db/compute_usage.sql.go new file mode 100644 index 000000000..f50d460b9 --- /dev/null +++ b/go/internal/store/db/compute_usage.sql.go @@ -0,0 +1,51 @@ +// Code generated by sqlc. DO NOT EDIT. +// versions: +// sqlc v1.31.1 +// source: compute_usage.sql + +package db + +import ( + "context" +) + +const closeOrphanedComputeIntervals = `-- name: CloseOrphanedComputeIntervals :execrows +WITH orphaned AS ( + SELECT starts.tenant_id, starts.interval_id, starts.agent_account_id, + starts.owner_user_id, starts.session_id, starts.runner_id + FROM compute_usage_events AS starts + WHERE starts.kind = 'start' + AND NOT EXISTS ( + SELECT 1 + FROM compute_usage_events AS ends + WHERE ends.tenant_id = starts.tenant_id + AND ends.interval_id = starts.interval_id + AND ends.kind = 'end' + ) + AND NOT EXISTS ( + SELECT 1 + FROM session_bindings AS bindings + WHERE bindings.tenant_id = starts.tenant_id + AND bindings.usage_interval_id = starts.interval_id + ) +) +INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id, estimated +) +SELECT orphaned.tenant_id, gen_random_uuid()::text AS id, orphaned.interval_id, 'end' AS kind, now() AS occurred_at, + orphaned.agent_account_id, orphaned.owner_user_id, orphaned.session_id, + orphaned.runner_id, true AS estimated + FROM orphaned +ON CONFLICT DO NOTHING +` + +// Close intervals only after both binding deletion and its end event are absent. +// The tenant id comes from each start because the system role has no tenant GUC. +func (q *Queries) CloseOrphanedComputeIntervals(ctx context.Context) (int64, error) { + result, err := q.db.Exec(ctx, closeOrphanedComputeIntervals) + if err != nil { + return 0, err + } + return result.RowsAffected(), nil +} diff --git a/go/internal/store/db/models.go b/go/internal/store/db/models.go index decf4e3c6..e54fc3043 100644 --- a/go/internal/store/db/models.go +++ b/go/internal/store/db/models.go @@ -151,6 +151,19 @@ type ChannelPin struct { TenantID string } +type ComputeUsageEvent struct { + TenantID string + ID string + IntervalID string + Kind string + OccurredAt pgtype.Timestamptz + AgentAccountID string + OwnerUserID string + SessionID string + RunnerID string + Estimated bool +} + type ForgeArtifactCursor struct { ForgeProvider int16 ForgeHost string @@ -303,12 +316,13 @@ type ServerSecret struct { } type SessionBinding struct { - TenantID string - AgentAccountID string - SessionID string - RunnerID string - CreatedAt pgtype.Timestamptz - UpdatedAt pgtype.Timestamptz + TenantID string + AgentAccountID string + SessionID string + RunnerID string + CreatedAt pgtype.Timestamptz + UpdatedAt pgtype.Timestamptz + UsageIntervalID string } type SystemAccount struct { diff --git a/go/internal/store/db/querier.go b/go/internal/store/db/querier.go index 2c595b936..35982368f 100644 --- a/go/internal/store/db/querier.go +++ b/go/internal/store/db/querier.go @@ -81,6 +81,9 @@ type Querier interface { ChannelVisibleTo(ctx context.Context, arg ChannelVisibleToParams) (bool, error) ChannelsByNameForViewer(ctx context.Context, arg ChannelsByNameForViewerParams) ([]ChannelsByNameForViewerRow, error) ClearOwedMention(ctx context.Context, arg ClearOwedMentionParams) (int64, error) + // Close intervals only after both binding deletion and its end event are absent. + // The tenant id comes from each start because the system role has no tenant GUC. + CloseOrphanedComputeIntervals(ctx context.Context) (int64, error) CollectSegment(ctx context.Context, arg CollectSegmentParams) ([]CollectSegmentRow, error) // Single-statement clear-and-return: the row is claimed and its actor returned // in ONE UPDATE, so a memo attributes at most one event and a concurrent second @@ -144,6 +147,9 @@ type Querier interface { // rollups can hold pruned events, so they stay. DeleteTokenUsageRollupsFrom(ctx context.Context, horizon pgtype.Timestamptz) error DeleteTopic(ctx context.Context, id string) error + // The account owner is read in SQL so callers continue to supply only the + // session binding identity. + EndComputeUsageInterval(ctx context.Context, arg EndComputeUsageIntervalParams) error // Agent-forge-subscription / artifact-cursor queries (sqlc adoption T6, // RIG-3034). These replace the inline SQL literals in // internal/store/forge_subscriptions.go; the hand-written Store methods keep @@ -382,11 +388,9 @@ type Querier interface { // map these rows into the SessionBinding domain struct (the AccountID newtype is // done inline in the Go, as agent_placements does). // - // No query here names tenant_id, and on the REQUEST path that is complete: the - // store arms every statement with SET LOCAL ROLE compass_app + the - // compass.tenant_id GUC (tenant_tx.go), so the RLS policy (0001_init.sql) scopes - // reads to the acting tenant and the tenant_id column DEFAULTs to that GUC on - // insert. That is how every other query file here is written. + // These request-path queries rely on SET LOCAL ROLE compass_app and the + // compass.tenant_id GUC for RLS and tenant defaults. Delete queries carry the + // deleted binding's tenant_id into their end-event insert explicitly. // // It is NOT complete under WithSystemRole (tenant_tx.go), which arms the // BYPASSRLS compass_system role and NO tenant GUC. Every query here then runs @@ -500,19 +504,8 @@ type Querier interface { // there is the NORMAL replay case, not a drop), so asserting rows-affected here // would wrongly fail an idempotent re-fire. RecordOwedMention(ctx context.Context, arg RecordOwedMentionParams) error - // The bind. Keyed on the ACCOUNT (see the table comment): the hub's 1:1 - // accountSessions map this replaces treats re-pointing an account at a newer - // session as an assignment, not a collision, so this is an upsert on - // (tenant_id, agent_account_id) and never refuses a re-point. - // // What it DISPLACED comes from SessionBindingForUpdate above, not from a - // RETURNING here: ON CONFLICT DO UPDATE's RETURNING sees the POST-update row, and - // the pre-update one is unreachable from this statement (`OLD`-aliased RETURNING - // is Postgres 18+; this targets 16). The two statements are nonetheless one - // operation, because they share a transaction and the row lock the read took — - // which is the property the caller needs. It reaps the displaced session from the - // delivery held-deliver registry, and reaping a session that is once again live - // would strand a live agent's deliveries. + // RETURNING here. The binding update and event writes share the Store tx. RecordSessionBinding(ctx context.Context, arg RecordSessionBindingParams) error // Forge state-transition memo queries (compass-forge-state-transition §Actor // attribution). The write chokepoint upserts one memo per forge coordinate @@ -571,12 +564,9 @@ type Querier interface { SessionBase(ctx context.Context, sessionID string) (int64, error) SessionBinding(ctx context.Context, sessionID string) (SessionBindingRow, error) SessionBindingForAccount(ctx context.Context, agentAccountID string) (string, error) - // The prior-value read of the bind, and the second of three statements the Store - // runs in ONE explicit transaction (beginTenantTx): the advisory lock above, this - // read, then the upsert below. FOR UPDATE takes a row lock on the binding this - // bind is about to overwrite, so the read and the write cannot be separated by a - // concurrent bind — that caller is already parked on the advisory lock, and would - // park here too. + // The prior-value read follows the account advisory lock. It shares a tx with + // event writes and the binding upsert. FOR UPDATE also protects against a writer + // that reaches the row without taking the advisory lock. // // The single-statement form this replaces (a `prev` CTE beside the upsert) could // not do that: under READ COMMITTED the CTE reads the statement-start snapshot @@ -587,7 +577,9 @@ type Querier interface { // // A MISS is not an error: a first-ever bind returns pgx.ErrNoRows and the Store // maps that to the empty displaced id. - SessionBindingForUpdate(ctx context.Context, agentAccountID string) (string, error) + // The prior binding, including the original Runner and interval identity. The + // Store closes a displaced interval before it writes a replacement binding. + SessionBindingForUpdate(ctx context.Context, agentAccountID string) (SessionBindingForUpdateRow, error) // The one query here meant for the system role: a Runner-originated call carries // no tenant, so the hub reads the session's tenant cross-tenant, then acts under it. // :many so a session id minted in two tenants is refused, not resolved arbitrarily. @@ -604,6 +596,9 @@ type Querier interface { SetIssueState(ctx context.Context, arg SetIssueStateParams) (int64, error) SetTopicArchived(ctx context.Context, arg SetTopicArchivedParams) error SharesVisibleChannel(ctx context.Context, arg SharesVisibleChannelParams) (bool, error) + // Event writes share RecordSessionBinding's transaction, so neither half of an + // interval can commit without its binding transition. + StartComputeUsageInterval(ctx context.Context, arg StartComputeUsageIntervalParams) error StoreForgeRepoWatermark(ctx context.Context, arg StoreForgeRepoWatermarkParams) (int64, error) SubscribeConvertedDMParties(ctx context.Context, channelID string) error // Delivery-consumer read queries (sqlc adoption T4, RIG-3034). These replace the diff --git a/go/internal/store/db/session_bindings.sql.go b/go/internal/store/db/session_bindings.sql.go index 8cf7d54c4..e5ee3b54c 100644 --- a/go/internal/store/db/session_bindings.sql.go +++ b/go/internal/store/db/session_bindings.sql.go @@ -12,7 +12,21 @@ import ( ) const deleteSessionBinding = `-- name: DeleteSessionBinding :exec -DELETE FROM session_bindings WHERE session_id = $1 +WITH d AS ( + DELETE FROM session_bindings AS b + WHERE b.session_id = $1 + RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, + b.session_id, b.runner_id +) +INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id +) +SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id +ON CONFLICT DO NOTHING ` func (q *Queries) DeleteSessionBinding(ctx context.Context, sessionID string) error { @@ -21,9 +35,24 @@ func (q *Queries) DeleteSessionBinding(ctx context.Context, sessionID string) er } const deleteSessionBindingsForRunner = `-- name: DeleteSessionBindingsForRunner :many -DELETE FROM session_bindings - WHERE runner_id = $1 -RETURNING session_id, agent_account_id +WITH d AS ( + DELETE FROM session_bindings AS b + WHERE b.runner_id = $1 + RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, + b.session_id, b.runner_id +), ins AS ( + INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id + ) + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id + ON CONFLICT DO NOTHING + RETURNING 1 +) +SELECT d.session_id, d.agent_account_id FROM d ` type DeleteSessionBindingsForRunnerRow struct { @@ -66,6 +95,37 @@ func (q *Queries) DeleteSessionBindingsForRunner(ctx context.Context, runnerID s return items, nil } +const endComputeUsageInterval = `-- name: EndComputeUsageInterval :exec +INSERT INTO compute_usage_events ( + id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, + session_id, runner_id +) +SELECT gen_random_uuid()::text, $1::text, 'end', now(), + a.account_id, a.owner_user_id, $2::text, $3::text + FROM agent_accounts AS a + WHERE a.account_id = $4::text +ON CONFLICT (tenant_id, interval_id, kind) DO NOTHING +` + +type EndComputeUsageIntervalParams struct { + IntervalID string + SessionID string + RunnerID string + AgentAccountID string +} + +// The account owner is read in SQL so callers continue to supply only the +// session binding identity. +func (q *Queries) EndComputeUsageInterval(ctx context.Context, arg EndComputeUsageIntervalParams) error { + _, err := q.db.Exec(ctx, endComputeUsageInterval, + arg.IntervalID, + arg.SessionID, + arg.RunnerID, + arg.AgentAccountID, + ) + return err +} + const lockSessionBindingAccount = `-- name: LockSessionBindingAccount :exec SELECT pg_advisory_xact_lock(hashtext('binding:' || $1 || ':' || $2)) @@ -82,11 +142,9 @@ type LockSessionBindingAccountParams struct { // map these rows into the SessionBinding domain struct (the AccountID newtype is // done inline in the Go, as agent_placements does). // -// No query here names tenant_id, and on the REQUEST path that is complete: the -// store arms every statement with SET LOCAL ROLE compass_app + the -// compass.tenant_id GUC (tenant_tx.go), so the RLS policy (0001_init.sql) scopes -// reads to the acting tenant and the tenant_id column DEFAULTs to that GUC on -// insert. That is how every other query file here is written. +// These request-path queries rely on SET LOCAL ROLE compass_app and the +// compass.tenant_id GUC for RLS and tenant defaults. Delete queries carry the +// deleted binding's tenant_id into their end-event insert explicitly. // // It is NOT complete under WithSystemRole (tenant_tx.go), which arms the // BYPASSRLS compass_system role and NO tenant GUC. Every query here then runs @@ -150,34 +208,30 @@ func (q *Queries) LockSessionBindingAccount(ctx context.Context, arg LockSession } const recordSessionBinding = `-- name: RecordSessionBinding :exec -INSERT INTO session_bindings (agent_account_id, session_id, runner_id) -VALUES ($1, $2, $3) +INSERT INTO session_bindings (agent_account_id, session_id, runner_id, usage_interval_id) +VALUES ($1, $2, $3, $4) ON CONFLICT (tenant_id, agent_account_id) DO UPDATE SET session_id = EXCLUDED.session_id, - runner_id = EXCLUDED.runner_id + runner_id = EXCLUDED.runner_id, + usage_interval_id = EXCLUDED.usage_interval_id ` type RecordSessionBindingParams struct { - AgentAccountID string - SessionID string - RunnerID string + AgentAccountID string + SessionID string + RunnerID string + UsageIntervalID string } -// The bind. Keyed on the ACCOUNT (see the table comment): the hub's 1:1 -// accountSessions map this replaces treats re-pointing an account at a newer -// session as an assignment, not a collision, so this is an upsert on -// (tenant_id, agent_account_id) and never refuses a re-point. -// // What it DISPLACED comes from SessionBindingForUpdate above, not from a -// RETURNING here: ON CONFLICT DO UPDATE's RETURNING sees the POST-update row, and -// the pre-update one is unreachable from this statement (`OLD`-aliased RETURNING -// is Postgres 18+; this targets 16). The two statements are nonetheless one -// operation, because they share a transaction and the row lock the read took — -// which is the property the caller needs. It reaps the displaced session from the -// delivery held-deliver registry, and reaping a session that is once again live -// would strand a live agent's deliveries. +// RETURNING here. The binding update and event writes share the Store tx. func (q *Queries) RecordSessionBinding(ctx context.Context, arg RecordSessionBindingParams) error { - _, err := q.db.Exec(ctx, recordSessionBinding, arg.AgentAccountID, arg.SessionID, arg.RunnerID) + _, err := q.db.Exec(ctx, recordSessionBinding, + arg.AgentAccountID, + arg.SessionID, + arg.RunnerID, + arg.UsageIntervalID, + ) return err } @@ -209,18 +263,21 @@ func (q *Queries) SessionBindingForAccount(ctx context.Context, agentAccountID s } const sessionBindingForUpdate = `-- name: SessionBindingForUpdate :one -SELECT b.session_id - FROM session_bindings b +SELECT b.session_id, b.usage_interval_id, b.runner_id + FROM session_bindings AS b WHERE b.agent_account_id = $1 FOR UPDATE ` -// The prior-value read of the bind, and the second of three statements the Store -// runs in ONE explicit transaction (beginTenantTx): the advisory lock above, this -// read, then the upsert below. FOR UPDATE takes a row lock on the binding this -// bind is about to overwrite, so the read and the write cannot be separated by a -// concurrent bind — that caller is already parked on the advisory lock, and would -// park here too. +type SessionBindingForUpdateRow struct { + SessionID string + UsageIntervalID string + RunnerID string +} + +// The prior-value read follows the account advisory lock. It shares a tx with +// event writes and the binding upsert. FOR UPDATE also protects against a writer +// that reaches the row without taking the advisory lock. // // The single-statement form this replaces (a `prev` CTE beside the upsert) could // not do that: under READ COMMITTED the CTE reads the statement-start snapshot @@ -231,11 +288,13 @@ SELECT b.session_id // // A MISS is not an error: a first-ever bind returns pgx.ErrNoRows and the Store // maps that to the empty displaced id. -func (q *Queries) SessionBindingForUpdate(ctx context.Context, agentAccountID string) (string, error) { +// The prior binding, including the original Runner and interval identity. The +// Store closes a displaced interval before it writes a replacement binding. +func (q *Queries) SessionBindingForUpdate(ctx context.Context, agentAccountID string) (SessionBindingForUpdateRow, error) { row := q.db.QueryRow(ctx, sessionBindingForUpdate, agentAccountID) - var session_id string - err := row.Scan(&session_id) - return session_id, err + var i SessionBindingForUpdateRow + err := row.Scan(&i.SessionID, &i.UsageIntervalID, &i.RunnerID) + return i, err } const sessionBindingTenants = `-- name: SessionBindingTenants :many @@ -269,3 +328,34 @@ func (q *Queries) SessionBindingTenants(ctx context.Context, arg SessionBindingT } return items, nil } + +const startComputeUsageInterval = `-- name: StartComputeUsageInterval :exec +INSERT INTO compute_usage_events ( + id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, + session_id, runner_id +) +SELECT gen_random_uuid()::text, $1::text, 'start', now(), + a.account_id, a.owner_user_id, $2::text, $3::text + FROM agent_accounts AS a + WHERE a.account_id = $4::text +ON CONFLICT (tenant_id, interval_id, kind) DO NOTHING +` + +type StartComputeUsageIntervalParams struct { + IntervalID string + SessionID string + RunnerID string + AgentAccountID string +} + +// Event writes share RecordSessionBinding's transaction, so neither half of an +// interval can commit without its binding transition. +func (q *Queries) StartComputeUsageInterval(ctx context.Context, arg StartComputeUsageIntervalParams) error { + _, err := q.db.Exec(ctx, startComputeUsageInterval, + arg.IntervalID, + arg.SessionID, + arg.RunnerID, + arg.AgentAccountID, + ) + return err +} diff --git a/go/internal/store/migrate_upgrade_pgtest_test.go b/go/internal/store/migrate_upgrade_pgtest_test.go index 4688115fe..b86b76934 100644 --- a/go/internal/store/migrate_upgrade_pgtest_test.go +++ b/go/internal/store/migrate_upgrade_pgtest_test.go @@ -22,6 +22,8 @@ func TestOpenUpgradesV1DatabaseToTokenUsage(t *testing.T) { ctx := t.Context() dsn := pgtest.RequireDSN(t) const tenant TenantID = "upgrade-tenant" + const userID = "upgrade-user" + const agentID = "upgrade-agent" applyV1Only(t, dsn) // The tenant predates the upgrade; the token-usage foreign keys must accept it. @@ -37,6 +39,39 @@ func TestOpenUpgradesV1DatabaseToTokenUsage(t *testing.T) { if err != nil { t.Fatalf("seed tenant at v1: %v", err) } + seed, err = pgxpool.New(ctx, dsn) + if err != nil { + t.Fatalf("connect to seed binding: %v", err) + } + _, err = seed.Exec(ctx, + "INSERT INTO accounts (id, handle, display_name, tenant_id) VALUES ($1, 'upgrade-user', 'Upgrade User', $2)", + userID, tenant, + ) + if err == nil { + _, err = seed.Exec(ctx, "INSERT INTO user_accounts (account_id, tenant_id) VALUES ($1, $2)", userID, tenant) + } + if err == nil { + _, err = seed.Exec(ctx, + "INSERT INTO accounts (id, handle, display_name, tenant_id) VALUES ($1, 'upgrade-agent', 'Upgrade Agent', $2)", + agentID, tenant, + ) + } + if err == nil { + _, err = seed.Exec(ctx, + "INSERT INTO agent_accounts (account_id, owner_user_id, tenant_id) VALUES ($1, $2, $3)", + agentID, userID, tenant, + ) + } + if err == nil { + _, err = seed.Exec(ctx, + "INSERT INTO session_bindings (tenant_id, agent_account_id, session_id, runner_id) VALUES ($1, $2, 'upgrade-session', 'upgrade-runner')", + tenant, agentID, + ) + } + seed.Close() + if err != nil { + t.Fatalf("seed v1 session binding: %v", err) + } s := openStore(t, dsn) @@ -52,7 +87,10 @@ func TestOpenUpgradesV1DatabaseToTokenUsage(t *testing.T) { t.Fatalf("schema version after upgrade = %d, want %d", version, want) } - for _, tbl := range []string{"token_usage_events", "token_usage_rollups_hourly", "token_usage_rollups_daily"} { + for _, tbl := range []string{ + "token_usage_events", "token_usage_rollups_hourly", "token_usage_rollups_daily", + "compute_usage_events", + } { var enabled, forced bool err := s.pool.QueryRow(ctx, `SELECT relrowsecurity, relforcerowsecurity FROM pg_class @@ -66,6 +104,24 @@ func TestOpenUpgradesV1DatabaseToTokenUsage(t *testing.T) { } } + var backfilledStartCount int + var backfilledInterval, boundInterval string + var allEstimated bool + if err := s.pool.QueryRow(ctx, ` + SELECT count(*), min(e.interval_id), bool_and(e.estimated), min(b.usage_interval_id) + FROM compute_usage_events AS e + JOIN session_bindings AS b + ON b.tenant_id = e.tenant_id AND b.agent_account_id = e.agent_account_id + WHERE e.tenant_id = $1 AND e.agent_account_id = $2 + AND e.kind = 'start' AND e.session_id = 'upgrade-session'`, string(tenant), agentID). + Scan(&backfilledStartCount, &backfilledInterval, &allEstimated, &boundInterval); err != nil { + t.Fatalf("read backfilled compute start: %v", err) + } + if backfilledStartCount != 1 || backfilledInterval == "" || backfilledInterval != boundInterval || !allEstimated { + t.Fatalf("backfilled starts = %d, interval = %q, bound = %q, estimated = %t; want one estimated start on the bound interval", + backfilledStartCount, backfilledInterval, boundInterval, allEstimated) + } + var horizon pgtype.Timestamptz if err := s.pool.QueryRow(ctx, "SELECT horizon FROM token_usage_prune_horizon").Scan(&horizon); err != nil { t.Fatalf("read prune horizon row: %v", err) diff --git a/go/internal/store/migrations/0004_compute_usage.sql b/go/internal/store/migrations/0004_compute_usage.sql new file mode 100644 index 000000000..1d7b7b60c --- /dev/null +++ b/go/internal/store/migrations/0004_compute_usage.sql @@ -0,0 +1,60 @@ +-- 0004_compute_usage: the append-only compute interval event log. Each +-- session binding is one billable interval; starts and ends commit with its row. + +CREATE TABLE compute_usage_events ( + tenant_id TEXT NOT NULL DEFAULT current_setting('compass.tenant_id', TRUE) REFERENCES tenants (id) ON DELETE RESTRICT, + id TEXT NOT NULL, + interval_id TEXT NOT NULL, + kind TEXT NOT NULL CHECK (kind IN ('start', 'end')), + occurred_at TIMESTAMPTZ NOT NULL, + agent_account_id TEXT NOT NULL, + owner_user_id TEXT NOT NULL, + session_id TEXT NOT NULL, + runner_id TEXT NOT NULL, + estimated BOOLEAN NOT NULL DEFAULT FALSE, + PRIMARY KEY (tenant_id, id), + UNIQUE (tenant_id, interval_id, kind) +); + +CREATE INDEX compute_usage_events_occurred_at_idx ON compute_usage_events (tenant_id, occurred_at); + +ALTER TABLE session_bindings ADD COLUMN usage_interval_id TEXT; + +UPDATE session_bindings AS b + SET usage_interval_id = gen_random_uuid()::TEXT; + +INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id, estimated +) +-- updated_at marks the last re-point, the best start this table still holds. +SELECT b.tenant_id, gen_random_uuid()::TEXT AS id, b.usage_interval_id, 'start' AS kind, b.updated_at AS occurred_at, + b.agent_account_id, a.owner_user_id, b.session_id, b.runner_id, TRUE AS estimated + FROM session_bindings AS b + JOIN agent_accounts AS a + ON a.account_id = b.agent_account_id + AND a.tenant_id = b.tenant_id; + +-- Every existing row got an id above, so the cutover cannot fail. +-- squawk-ignore adding-not-nullable-field +ALTER TABLE session_bindings ALTER COLUMN usage_interval_id SET NOT NULL; + +GRANT SELECT, INSERT ON compute_usage_events TO compass_app, compass_system; + +DO $$ +DECLARE + t text; + tenant_tables text[] := ARRAY['compute_usage_events']; +BEGIN + FOREACH t IN ARRAY tenant_tables LOOP + EXECUTE format('ALTER TABLE %I ENABLE ROW LEVEL SECURITY', t); + EXECUTE format('ALTER TABLE %I FORCE ROW LEVEL SECURITY', t); + EXECUTE format($f$ + CREATE POLICY tenant_isolation ON %I + USING ((SELECT current_setting('compass.tenant_id', TRUE)) <> '' + AND tenant_id = (SELECT current_setting('compass.tenant_id', TRUE))) + WITH CHECK ((SELECT current_setting('compass.tenant_id', TRUE)) <> '' + AND tenant_id = (SELECT current_setting('compass.tenant_id', TRUE))) + $f$, t); + END LOOP; +END $$; diff --git a/go/internal/store/queries/compute_usage.sql b/go/internal/store/queries/compute_usage.sql new file mode 100644 index 000000000..1d9096a06 --- /dev/null +++ b/go/internal/store/queries/compute_usage.sql @@ -0,0 +1,31 @@ +-- Close intervals only after both binding deletion and its end event are absent. +-- The tenant id comes from each start because the system role has no tenant GUC. +-- name: CloseOrphanedComputeIntervals :execrows +WITH orphaned AS ( + SELECT starts.tenant_id, starts.interval_id, starts.agent_account_id, + starts.owner_user_id, starts.session_id, starts.runner_id + FROM compute_usage_events AS starts + WHERE starts.kind = 'start' + AND NOT EXISTS ( + SELECT 1 + FROM compute_usage_events AS ends + WHERE ends.tenant_id = starts.tenant_id + AND ends.interval_id = starts.interval_id + AND ends.kind = 'end' + ) + AND NOT EXISTS ( + SELECT 1 + FROM session_bindings AS bindings + WHERE bindings.tenant_id = starts.tenant_id + AND bindings.usage_interval_id = starts.interval_id + ) +) +INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id, estimated +) +SELECT orphaned.tenant_id, gen_random_uuid()::text AS id, orphaned.interval_id, 'end' AS kind, now() AS occurred_at, + orphaned.agent_account_id, orphaned.owner_user_id, orphaned.session_id, + orphaned.runner_id, true AS estimated + FROM orphaned +ON CONFLICT DO NOTHING; diff --git a/go/internal/store/queries/session_bindings.sql b/go/internal/store/queries/session_bindings.sql index 2c4a2558e..dac0a8c02 100644 --- a/go/internal/store/queries/session_bindings.sql +++ b/go/internal/store/queries/session_bindings.sql @@ -4,11 +4,9 @@ -- map these rows into the SessionBinding domain struct (the AccountID newtype is -- done inline in the Go, as agent_placements does). -- --- No query here names tenant_id, and on the REQUEST path that is complete: the --- store arms every statement with SET LOCAL ROLE compass_app + the --- compass.tenant_id GUC (tenant_tx.go), so the RLS policy (0001_init.sql) scopes --- reads to the acting tenant and the tenant_id column DEFAULTs to that GUC on --- insert. That is how every other query file here is written. +-- These request-path queries rely on SET LOCAL ROLE compass_app and the +-- compass.tenant_id GUC for RLS and tenant defaults. Delete queries carry the +-- deleted binding's tenant_id into their end-event insert explicitly. -- -- It is NOT complete under WithSystemRole (tenant_tx.go), which arms the -- BYPASSRLS compass_system role and NO tenant GUC. Every query here then runs @@ -70,12 +68,9 @@ -- name: LockSessionBindingAccount :exec SELECT pg_advisory_xact_lock(hashtext('binding:' || $1 || ':' || $2)); --- The prior-value read of the bind, and the second of three statements the Store --- runs in ONE explicit transaction (beginTenantTx): the advisory lock above, this --- read, then the upsert below. FOR UPDATE takes a row lock on the binding this --- bind is about to overwrite, so the read and the write cannot be separated by a --- concurrent bind — that caller is already parked on the advisory lock, and would --- park here too. +-- The prior-value read follows the account advisory lock. It shares a tx with +-- event writes and the binding upsert. FOR UPDATE also protects against a writer +-- that reaches the row without taking the advisory lock. -- -- The single-statement form this replaces (a `prev` CTE beside the upsert) could -- not do that: under READ COMMITTED the CTE reads the statement-start snapshot @@ -86,31 +81,49 @@ SELECT pg_advisory_xact_lock(hashtext('binding:' || $1 || ':' || $2)); -- -- A MISS is not an error: a first-ever bind returns pgx.ErrNoRows and the Store -- maps that to the empty displaced id. +-- The prior binding, including the original Runner and interval identity. The +-- Store closes a displaced interval before it writes a replacement binding. -- name: SessionBindingForUpdate :one -SELECT b.session_id - FROM session_bindings b +SELECT b.session_id, b.usage_interval_id, b.runner_id + FROM session_bindings AS b WHERE b.agent_account_id = $1 FOR UPDATE; --- The bind. Keyed on the ACCOUNT (see the table comment): the hub's 1:1 --- accountSessions map this replaces treats re-pointing an account at a newer --- session as an assignment, not a collision, so this is an upsert on --- (tenant_id, agent_account_id) and never refuses a re-point. --- +-- Event writes share RecordSessionBinding's transaction, so neither half of an +-- interval can commit without its binding transition. +-- name: StartComputeUsageInterval :exec +INSERT INTO compute_usage_events ( + id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, + session_id, runner_id +) +SELECT gen_random_uuid()::text, @interval_id::text, 'start', now(), + a.account_id, a.owner_user_id, @session_id::text, @runner_id::text + FROM agent_accounts AS a + WHERE a.account_id = @agent_account_id::text +ON CONFLICT (tenant_id, interval_id, kind) DO NOTHING; + +-- The account owner is read in SQL so callers continue to supply only the +-- session binding identity. +-- name: EndComputeUsageInterval :exec +INSERT INTO compute_usage_events ( + id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, + session_id, runner_id +) +SELECT gen_random_uuid()::text, @interval_id::text, 'end', now(), + a.account_id, a.owner_user_id, @session_id::text, @runner_id::text + FROM agent_accounts AS a + WHERE a.account_id = @agent_account_id::text +ON CONFLICT (tenant_id, interval_id, kind) DO NOTHING; + -- What it DISPLACED comes from SessionBindingForUpdate above, not from a --- RETURNING here: ON CONFLICT DO UPDATE's RETURNING sees the POST-update row, and --- the pre-update one is unreachable from this statement (`OLD`-aliased RETURNING --- is Postgres 18+; this targets 16). The two statements are nonetheless one --- operation, because they share a transaction and the row lock the read took — --- which is the property the caller needs. It reaps the displaced session from the --- delivery held-deliver registry, and reaping a session that is once again live --- would strand a live agent's deliveries. +-- RETURNING here. The binding update and event writes share the Store tx. -- name: RecordSessionBinding :exec -INSERT INTO session_bindings (agent_account_id, session_id, runner_id) -VALUES ($1, $2, $3) +INSERT INTO session_bindings (agent_account_id, session_id, runner_id, usage_interval_id) +VALUES ($1, $2, $3, $4) ON CONFLICT (tenant_id, agent_account_id) DO UPDATE SET session_id = EXCLUDED.session_id, - runner_id = EXCLUDED.runner_id; + runner_id = EXCLUDED.runner_id, + usage_interval_id = EXCLUDED.usage_interval_id; -- name: SessionBinding :one SELECT agent_account_id, runner_id FROM session_bindings WHERE session_id = $1; @@ -119,7 +132,21 @@ SELECT agent_account_id, runner_id FROM session_bindings WHERE session_id = $1; SELECT session_id FROM session_bindings WHERE agent_account_id = $1; -- name: DeleteSessionBinding :exec -DELETE FROM session_bindings WHERE session_id = $1; +WITH d AS ( + DELETE FROM session_bindings AS b + WHERE b.session_id = $1 + RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, + b.session_id, b.runner_id +) +INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id +) +SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id +ON CONFLICT DO NOTHING; -- The reconnect sweep. Hub.enroll (internal/runnerhub/hub.go:905-957) clears -- every binding when a Runner (re-)enrolls: a reconnecting Runner has no live @@ -137,9 +164,24 @@ DELETE FROM session_bindings WHERE session_id = $1; -- A DELETE ... RETURNING takes no ORDER BY, so the Store method sorts the -- returned slice by session id to keep a sweep pass deterministic and diffable. -- name: DeleteSessionBindingsForRunner :many -DELETE FROM session_bindings - WHERE runner_id = $1 -RETURNING session_id, agent_account_id; +WITH d AS ( + DELETE FROM session_bindings AS b + WHERE b.runner_id = $1 + RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, + b.session_id, b.runner_id +), ins AS ( + INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id + ) + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id + ON CONFLICT DO NOTHING + RETURNING 1 +) +SELECT d.session_id, d.agent_account_id FROM d; -- The one query here meant for the system role: a Runner-originated call carries -- no tenant, so the hub reads the session's tenant cross-tenant, then acts under it. diff --git a/go/internal/store/rls_pgtest_test.go b/go/internal/store/rls_pgtest_test.go index 01f2a5ab9..54023ce65 100644 --- a/go/internal/store/rls_pgtest_test.go +++ b/go/internal/store/rls_pgtest_test.go @@ -606,6 +606,7 @@ func TestRLSCatalogEnabledAndForced(t *testing.T) { "issues", "forge_repo_subscriptions", "forge_artifact_cursors", "forge_state_transitions", "token_usage_events", "token_usage_rollups_hourly", "token_usage_rollups_daily", + "compute_usage_events", } for _, tbl := range tenantOwned { if !enumerated[tbl] { diff --git a/go/internal/store/session_bindings.go b/go/internal/store/session_bindings.go index 36ee5c564..8d60ba860 100644 --- a/go/internal/store/session_bindings.go +++ b/go/internal/store/session_bindings.go @@ -6,6 +6,7 @@ import ( "fmt" "slices" + "github.com/google/uuid" "github.com/jackc/pgx/v5/pgtype" "github.com/RigelBuild/compass/go/internal/store/db" @@ -50,12 +51,10 @@ type SessionBinding struct { // DeleteSessionBindingsForRunner's returned rows exist for. A displaced session // left unreaped holds deliveries for an account that has already moved on. // -// So the bind is THREE statements in one explicit transaction, and that is the -// whole reason this method opens a tx at all: a per-account advisory lock, then -// the prior-value read under FOR UPDATE, then the write. A concurrent bind for -// the same account parks on the advisory lock until this transaction commits, so -// it cannot interleave between this read and this write — the two binds -// serialize, and each caller is told the session IT actually displaced. +// The bind uses one explicit transaction: a per-account advisory lock, a prior +// row read, any interval close/start events, and the new binding write. Concurrent +// binds serialize at the advisory lock, so each caller gets the correct displaced +// session and no event can commit without its binding transition. // // A single statement could not give that. The `prev` CTE this replaces read the // statement-start snapshot while the upsert's ON CONFLICT DO UPDATE blocked on @@ -134,21 +133,43 @@ func (s *Store) RecordSessionBinding(ctx context.Context, sessionID string, acco return "", fmt.Errorf("store: lock session binding account: %w", err) } - displaced, err := qtx.SessionBindingForUpdate(ctx, string(accountID)) + prior, err := qtx.SessionBindingForUpdate(ctx, string(accountID)) if err != nil && !noRows(err) { return "", fmt.Errorf("store: read prior session binding: %w", err) } if noRows(err) { - // No prior binding: nothing to displace. Scan left displaced at its zero - // value, which is already the empty session id this method reports for - // "the account held none". - displaced = "" + prior = db.SessionBindingForUpdateRow{} + } + displaced := prior.SessionID + + intervalID := prior.UsageIntervalID + if prior.SessionID != sessionID { + if intervalID != "" { + if err := qtx.EndComputeUsageInterval(ctx, db.EndComputeUsageIntervalParams{ + IntervalID: intervalID, + AgentAccountID: string(accountID), + SessionID: prior.SessionID, + RunnerID: prior.RunnerID, + }); err != nil { + return "", fmt.Errorf("store: end displaced compute usage interval: %w", err) + } + } + intervalID = uuid.NewString() + if err := qtx.StartComputeUsageInterval(ctx, db.StartComputeUsageIntervalParams{ + IntervalID: intervalID, + AgentAccountID: string(accountID), + SessionID: sessionID, + RunnerID: runnerID, + }); err != nil { + return "", fmt.Errorf("store: start compute usage interval: %w", err) + } } if err := qtx.RecordSessionBinding(ctx, db.RecordSessionBindingParams{ - SessionID: sessionID, - AgentAccountID: string(accountID), - RunnerID: runnerID, + SessionID: sessionID, + AgentAccountID: string(accountID), + RunnerID: runnerID, + UsageIntervalID: intervalID, }); err != nil { if pgErrIs(err, pgForeignKeyViolation) { return "", fmt.Errorf("%w: agent account %q does not exist", ErrInvalidArgument, accountID) @@ -158,7 +179,6 @@ func (s *Store) RecordSessionBinding(ctx context.Context, sessionID string, acco } return "", fmt.Errorf("store: record session binding: %w", err) } - if err := tx.Commit(ctx); err != nil { return "", fmt.Errorf("store: commit record session binding: %w", err) } diff --git a/go/internal/store/tenant_tx.go b/go/internal/store/tenant_tx.go index 5888d4603..70e2e2f44 100644 --- a/go/internal/store/tenant_tx.go +++ b/go/internal/store/tenant_tx.go @@ -18,8 +18,8 @@ import ( // per-transaction compass.tenant_id GUC. // - systemRole is the narrowly-scoped BYPASSRLS role the cross-tenant // background loops (N5/OQ-4: delivery-cursor sweep, deliver-ack advance, -// reattach recovery, lag-resync) run under, and ONLY those. It carries no -// tenant GUC — it is cross-tenant by design. +// reattach recovery, lag-resync, compute-usage orphan sweep) run under, and +// ONLY those. It carries no tenant GUC — it is cross-tenant by design. const ( appRole = "compass_app" systemRole = "compass_system" @@ -33,18 +33,17 @@ const tenantGUC = "compass.tenant_id" // systemRoleKey marks a context as running the cross-tenant system path. When // present, the store arms statements with SET LOCAL ROLE compass_system -// (BYPASSRLS) and no tenant GUC, instead of the tenant-scoped app role. Only the -// four N5 background loops set it (WithSystemRole); every other path is -// tenant-scoped and fail-closed. +// (BYPASSRLS) and no tenant GUC, instead of the tenant-scoped app role. Only +// named background loops use it; every other path is tenant-scoped and fail-closed. type systemRoleKey struct{} // WithSystemRole marks ctx as the cross-tenant background/system path: store // calls made under it run as the BYPASSRLS compass_system role and see every // tenant's rows. It is the OQ-4 (Matt-ruled option 1) exemption, applied ONLY at -// the four named background-loop entrypoints (the delivery consumer's Run, the -// hub's deliver-ack / forge-notification-ack arms, and reattach recovery) — a -// request-path call NEVER sets it, so the request path stays tenant-scoped and -// fail-closed under RLS. +// named background-loop entrypoints (the delivery consumer's Run, the hub's +// deliver-ack / forge-notification-ack arms, reattach recovery, and the +// compute-usage orphan sweep) — a request-path call NEVER sets it, so the +// request path stays tenant-scoped and fail-closed under RLS. func WithSystemRole(ctx context.Context) context.Context { return context.WithValue(ctx, systemRoleKey{}, true) } diff --git a/go/internal/usage/compute_sweep.go b/go/internal/usage/compute_sweep.go new file mode 100644 index 000000000..3ccfdc827 --- /dev/null +++ b/go/internal/usage/compute_sweep.go @@ -0,0 +1,54 @@ +package usage + +import ( + "context" + "log/slog" + "time" +) + +const computeSweepInterval = time.Hour + +type computeIntervalCloser interface { + CloseOrphanedComputeIntervals(ctx context.Context) (int64, error) +} + +// ComputeUsageSweeper closes intervals left without a binding or end event. +type ComputeUsageSweeper struct { + store computeIntervalCloser + log *slog.Logger +} + +// NewComputeUsageSweeper returns a sweeper that closes orphaned intervals. +func NewComputeUsageSweeper(s computeIntervalCloser, log *slog.Logger) *ComputeUsageSweeper { + if log == nil { + log = slog.Default() + } + return &ComputeUsageSweeper{store: s, log: log} +} + +// Run sweeps once at startup and then hourly until ctx ends. A failed sweep is +// logged and retried on the next tick so it cannot stop the serve group. +func (w *ComputeUsageSweeper) Run(ctx context.Context) error { + w.sweep(ctx) + ticker := time.NewTicker(computeSweepInterval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return nil + case <-ticker.C: + w.sweep(ctx) + } + } +} + +func (w *ComputeUsageSweeper) sweep(ctx context.Context) { + closed, err := w.store.CloseOrphanedComputeIntervals(ctx) + switch { + case ctx.Err() != nil: + case err != nil: + w.log.ErrorContext(ctx, "compute usage: close orphaned intervals", "error", err) + case closed > 0: + w.log.InfoContext(ctx, "compute usage: closed orphaned intervals", "closed", closed) + } +} diff --git a/go/internal/usage/compute_sweep_test.go b/go/internal/usage/compute_sweep_test.go new file mode 100644 index 000000000..270aad28a --- /dev/null +++ b/go/internal/usage/compute_sweep_test.go @@ -0,0 +1,32 @@ +package usage_test + +import ( + "context" + "log/slog" + "testing" + + "github.com/RigelBuild/compass/go/internal/usage" +) + +type computeCloser struct { + calls int +} + +func (c *computeCloser) CloseOrphanedComputeIntervals(context.Context) (int64, error) { + c.calls++ + return 0, nil +} + +func TestComputeUsageSweeperRunsImmediatelyAndStopsOnCancellation(t *testing.T) { + closer := &computeCloser{} + sweeper := usage.NewComputeUsageSweeper(closer, slog.New(slog.DiscardHandler)) + ctx, cancel := context.WithCancel(t.Context()) + cancel() + + if err := sweeper.Run(ctx); err != nil { + t.Fatalf("Run = %v, want nil after cancellation", err) + } + if closer.calls != 1 { + t.Fatalf("CloseOrphanedComputeIntervals calls = %d, want one immediate pass", closer.calls) + } +} diff --git a/go/server/serve.go b/go/server/serve.go index 524ae404d..bd872a411 100644 --- a/go/server/serve.go +++ b/go/server/serve.go @@ -909,6 +909,9 @@ func Serve(ctx context.Context, cfg ServeConfig) error { startCommsConsumers(gctx, g, commsBus, fab, st, hub, hubLog) // A daily prune bounds the raw usage log; the rollups keep their sums. startUsageRetention(gctx, g, st, cfg.UsageEventRetention, hubLog) + // Close orphaned compute intervals on startup and once per hour; binding + // transitions normally close intervals in their own transaction. + startComputeUsageSweeper(gctx, g, st, hubLog) // Drain member of the same group: wake on gctx cancellation, then hand off to // drainDoors. A drain that overruns (a handler still wedged mid-replay) // surfaces as the error rather than a false clean shutdown; a real serve diff --git a/go/server/sinks.go b/go/server/sinks.go index 2279a47df..3fa823d8f 100644 --- a/go/server/sinks.go +++ b/go/server/sinks.go @@ -191,6 +191,18 @@ func startUsageRetention(gctx context.Context, g *errgroup.Group, st *store.Stor g.Go(func() error { return w.Run(gctx) }) } +// startComputeUsageSweeper closes orphaned compute intervals on the serve group. +func startComputeUsageSweeper(gctx context.Context, g *errgroup.Group, st *store.Store, log *slog.Logger) { + w := usage.NewComputeUsageSweeper(computeUsageCloser{st: st}, log) + g.Go(func() error { return w.Run(gctx) }) +} + +type computeUsageCloser struct{ st *store.Store } + +func (c computeUsageCloser) CloseOrphanedComputeIntervals(ctx context.Context) (int64, error) { + return c.st.CloseOrphanedComputeIntervals(store.WithSystemRole(ctx)) +} + // logFrameDiagnostics emits the hub's frame-loss snapshot as one line. Serve // calls it on shutdown, so every run states plainly how many relayed frames // never reached their surface. From 1ade0110cd147f95fe9e26a1fdafb92389a5f3b3 Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 03:18:07 -0400 Subject: [PATCH 02/10] fix(usage): keep backfilled starts, reap bindings across tenants, stamp events at write time (RIG-2872) - 0004 adds usage_interval_id with a volatile default instead of UPDATE, so the set_updated_at trigger no longer overwrites the estimated start, and older servers can still insert during a rolling deploy. - The backfill runs as compass_system, as 0002 does, so FORCE RLS cannot hide rows from a non-superuser owner. - Interval events use clock_timestamp(), so an event written after a lock wait is stamped when it is written. - Hub.enroll reaps a reconnecting Runner's bindings under the system role, so every tenant's interval ends; each archive runs under the row's tenant. - Tests: conflict rollback, re-point ordering, upgrade start timestamp, cross-tenant re-enroll reap. Open decisions: RIG-4228 (joint restart leaves bindings) and RIG-4229 (cross-tenant account binding). Refs RIG-2872 Co-authored-by: Matt Wilkinson --- go/internal/runnerhub/binding_cache_test.go | 6 +- go/internal/runnerhub/hub.go | 27 ++++--- .../runnerhub/runner_tenant_pgtest_test.go | 73 +++++++++++++------ .../store/compute_usage_pgtest_test.go | 34 ++++++++- go/internal/store/db/querier.go | 60 ++++----------- go/internal/store/db/session_bindings.sql.go | 73 ++++++------------- .../store/migrate_upgrade_pgtest_test.go | 15 ++-- .../store/migrations/0004_compute_usage.sql | 15 ++-- .../store/queries/session_bindings.sql | 70 ++++++------------ go/internal/store/session_bindings.go | 8 +- .../store/session_bindings_pgtest_test.go | 15 ++-- 11 files changed, 191 insertions(+), 205 deletions(-) diff --git a/go/internal/runnerhub/binding_cache_test.go b/go/internal/runnerhub/binding_cache_test.go index e8660321f..fec4fb581 100644 --- a/go/internal/runnerhub/binding_cache_test.go +++ b/go/internal/runnerhub/binding_cache_test.go @@ -87,7 +87,7 @@ func (f *fakeBindingStore) RecordSessionBinding(_ context.Context, sessionID str break } } - f.bindings[sessionID] = store.SessionBinding{SessionID: sessionID, AccountID: accountID, RunnerID: runnerID} + f.bindings[sessionID] = store.SessionBinding{TenantID: f.tenant, SessionID: sessionID, AccountID: accountID, RunnerID: runnerID} return displaced, nil } @@ -165,7 +165,7 @@ func (f *fakeBindingStore) SessionBindingTenant(_ context.Context, sessionID, ru func (f *fakeBindingStore) seed(sessionID string) { f.mu.Lock() defer f.mu.Unlock() - f.bindings[sessionID] = store.SessionBinding{SessionID: sessionID, AccountID: testAgentAccount, RunnerID: testRunnerID} + f.bindings[sessionID] = store.SessionBinding{TenantID: f.tenant, SessionID: sessionID, AccountID: testAgentAccount, RunnerID: testRunnerID} } // fakeRoutingFabric is an in-memory RoutingFabric double: PublishBindingChange @@ -610,7 +610,7 @@ func TestReusedSessionIDConflictIsSwallowed(t *testing.T) { // A row from before the joint restart, under a DIFFERENT account. bindings.mu.Lock() - bindings.bindings["sess-1"] = store.SessionBinding{SessionID: "sess-1", AccountID: "acct-stale", RunnerID: "runner-1"} + bindings.bindings["sess-1"] = store.SessionBinding{TenantID: bindings.tenant, SessionID: "sess-1", AccountID: "acct-stale", RunnerID: "runner-1"} bindings.mu.Unlock() hub.bindContainer("cont-1", testAgentAccount, "runner-1") diff --git a/go/internal/runnerhub/hub.go b/go/internal/runnerhub/hub.go index 4c8954019..9b15e60b1 100644 --- a/go/internal/runnerhub/hub.go +++ b/go/internal/runnerhub/hub.go @@ -981,9 +981,9 @@ type promotedPair struct { // (OQ-2): a restarted Runner could re-mint a still-bound id, so clearing forces a // re-minted id to CodeNotFound until bound anew. RIG-3108: the maps are a read-through // cache over session_bindings whose rows survive process death. Every enroll reaps -// this Runner's rows and drives OFFLINE + reap edges from them: a Runner enrolls once -// per process and sweeps its stale containers, so none of its pre-enroll sessions live. -// The enroll ctx carries no tenant, so the reap reaches only the bootstrap tenant's rows. +// this Runner's rows across tenants under the system role and drives OFFLINE + reap +// edges from them: a Runner enrolls once per process and sweeps its stale containers, +// so none of its pre-enroll sessions live. func (h *Hub) enroll(ctx context.Context, id string, subject store.Subject, tier compassv1.RuntimeTier, egressPosture compassv1.EgressPosture) (reattached bool) { h.mu.Lock() reattached = h.runner != nil @@ -1025,15 +1025,18 @@ func (h *Hub) enroll(ctx context.Context, id string, subject store.Subject, tier // snapshot taken above. offline := ramOffline reapedSessions := ramReaped + var durableReaped []store.SessionBinding + durableReapSucceeded := false if bindings != nil { - rows, err := bindings.DeleteSessionBindingsForRunner(ctx, id) + rows, err := bindings.DeleteSessionBindingsForRunner(store.WithSystemRole(ctx), id) if err != nil { - // A durable-reap fault must not wedge the reconnect: log and fall back to the - // in-RAM snapshot (still cleared). The rows SURVIVE and name dead sessions, so - // reapStale STAYS raised — a read-through would resurrect one. + // A durable-reap fault must not wedge reconnect; fall back to the in-RAM + // snapshot while read-through stays disabled until a reap succeeds. h.log.Error("durable session-binding reap failed on enroll; using in-RAM snapshot, read-through disabled until a reap succeeds", "runner_id", id, "error", err) } else { + durableReaped = rows + durableReapSucceeded = true offline = make([]promotedPair, 0, len(rows)) reapedSessions = make([]string, 0, len(rows)) for _, b := range rows { @@ -1061,8 +1064,14 @@ func (h *Hub) enroll(ctx context.Context, id string, subject store.Subject, tier if reap != nil { reap.OnSessionsReaped(reapedSessions) } - for _, sessionID := range reapedSessions { - h.archiveEnded(ctx, sessionID) + if durableReapSucceeded { + for _, b := range durableReaped { + h.archiveEnded(store.WithTenant(ctx, b.TenantID), b.SessionID) + } + } else { + for _, sessionID := range reapedSessions { + h.archiveEnded(ctx, sessionID) + } } return reattached } diff --git a/go/internal/runnerhub/runner_tenant_pgtest_test.go b/go/internal/runnerhub/runner_tenant_pgtest_test.go index 8faeecb21..49dd88698 100644 --- a/go/internal/runnerhub/runner_tenant_pgtest_test.go +++ b/go/internal/runnerhub/runner_tenant_pgtest_test.go @@ -19,6 +19,56 @@ import ( // woken: the hub resolves the session's tenant before it reads or deletes. func TestDropLostSessionScopesToTheSessionTenant(t *testing.T) { ctx := context.Background() + st, agent, ctxB := openTenantBSession(t, ctx) + + // A fresh hub: its cache is cold, as after a Server restart. + hub := newHubOnly() + hub.SetSessionBindingStore(st) + sink := &recordingLostSink{} + hub.SetSessionLostSink(sink) + hub.enroll(ctx, "runner-1", runnerSubject(), compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED) + + hub.dropLostSession(ctx, "runner-1", "sess-b") + + if len(sink.lost) != 1 || sink.lost[0] != agent.ID { + t.Fatalf("lost = %v, want [%s]: the tenant-B session was not resolved, so no wake", sink.lost, agent.ID) + } + if _, _, err := st.ResolveSessionBinding(ctxB, "sess-b"); err == nil { + t.Fatal("tenant B's durable binding survived the drop; a cache miss would resurrect it") + } +} + +// A re-enrolling Runner spans tenants, so its system-role sweep must reap tenant B's +// binding and record its end event with tenant B's id. +func TestEnrollReapsSessionBindingAcrossTenants(t *testing.T) { + ctx := context.Background() + st, agent, ctxB := openTenantBSession(t, ctx) + + hub := newHubOnly() + hub.SetSessionBindingStore(st) + // The first enroll keeps durable rows; the second is the reconnect that reaps. + hub.enroll(ctx, "runner-1", runnerSubject(), compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED) + hub.enroll(ctx, "runner-1", runnerSubject(), compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED) + + var remaining, ends int + if err := st.WithTx(ctxB, func(tx pgx.Tx) error { + if err := tx.QueryRow(ctxB, "SELECT count(*) FROM session_bindings WHERE session_id = 'sess-b'").Scan(&remaining); err != nil { + return err + } + return tx.QueryRow(ctxB, `SELECT count(*) FROM compute_usage_events + WHERE agent_account_id = $1 AND session_id = 'sess-b' AND kind = 'end'`, string(agent.ID)).Scan(&ends) + }); err != nil { + t.Fatalf("read tenant B after enroll reap: %v", err) + } + if remaining != 0 || ends != 1 { + t.Fatalf("tenant B binding remaining = %d, end events = %d; want reaped binding and one end", remaining, ends) + } +} + +// openTenantBSession opens a store with a non-bootstrap tenant B whose agent holds +// sess-b on runner-1. +func openTenantBSession(t *testing.T, ctx context.Context) (*store.Store, store.Account, context.Context) { + t.Helper() dsn := pgtest.RequireDSN(t) st, err := store.Open(ctx, dsn) if err != nil { @@ -28,7 +78,6 @@ func TestDropLostSessionScopesToTheSessionTenant(t *testing.T) { if _, err := st.BootstrapAdmin(ctx, store.NewUser{Handle: "admin", DisplayName: "admin"}); err != nil { t.Fatalf("BootstrapAdmin: %v", err) } - tenantB := seedTenantRow(t, ctx, dsn, "tenant-b") ctxB := store.WithTenant(ctx, tenantB) owner, err := st.CreateUser(ctxB, store.NewUser{Handle: "owner-b", DisplayName: "owner-b"}) @@ -42,27 +91,7 @@ func TestDropLostSessionScopesToTheSessionTenant(t *testing.T) { if _, err := st.RecordSessionBinding(ctxB, "sess-b", agent.ID, "runner-1"); err != nil { t.Fatalf("RecordSessionBinding: %v", err) } - - // A fresh hub: its cache is cold, as after a Server restart. - hub := newHubOnly() - hub.SetSessionBindingStore(st) - sink := &recordingLostSink{} - hub.SetSessionLostSink(sink) - hub.enroll(ctx, "runner-1", runnerSubject(), compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED) - // enroll's reap ran unscoped and cannot see tenant B's row; re-record it so the - // test does not depend on that. - if _, err := st.RecordSessionBinding(ctxB, "sess-b", agent.ID, "runner-1"); err != nil { - t.Fatalf("RecordSessionBinding after enroll: %v", err) - } - - hub.dropLostSession(ctx, "runner-1", "sess-b") - - if len(sink.lost) != 1 || sink.lost[0] != agent.ID { - t.Fatalf("lost = %v, want [%s]: the tenant-B session was not resolved, so no wake", sink.lost, agent.ID) - } - if _, _, err := st.ResolveSessionBinding(ctxB, "sess-b"); err == nil { - t.Fatal("tenant B's durable binding survived the drop; a cache miss would resurrect it") - } + return st, agent, ctxB } // seedTenantRow inserts a non-bootstrap tenant. tenants is RLS-exempt, so a plain diff --git a/go/internal/store/compute_usage_pgtest_test.go b/go/internal/store/compute_usage_pgtest_test.go index bda2d99d1..e2da565c5 100644 --- a/go/internal/store/compute_usage_pgtest_test.go +++ b/go/internal/store/compute_usage_pgtest_test.go @@ -4,7 +4,9 @@ package store import ( "context" + "errors" "testing" + "time" "github.com/jackc/pgx/v5" ) @@ -18,6 +20,7 @@ type computeUsageEvent struct { OwnerUserID string SessionID string RunnerID string + OccurredAt time.Time Estimated bool } @@ -25,7 +28,7 @@ func computeEvents(t *testing.T, s *Store, tenant TenantID, accountID AccountID) t.Helper() rows, err := s.pool.Query(t.Context(), ` SELECT tenant_id, id, interval_id, kind, agent_account_id, owner_user_id, - session_id, runner_id, estimated + session_id, runner_id, occurred_at, estimated FROM compute_usage_events WHERE tenant_id = $1 AND agent_account_id = $2 ORDER BY session_id, CASE kind WHEN 'start' THEN 0 ELSE 1 END`, string(tenant), string(accountID)) @@ -38,7 +41,7 @@ func computeEvents(t *testing.T, s *Store, tenant TenantID, accountID AccountID) var event computeUsageEvent if err := rows.Scan(&event.TenantID, &event.ID, &event.IntervalID, &event.Kind, &event.AgentAccountID, &event.OwnerUserID, &event.SessionID, &event.RunnerID, - &event.Estimated); err != nil { + &event.OccurredAt, &event.Estimated); err != nil { t.Fatalf("scan compute usage event: %v", err) } events = append(events, event) @@ -122,11 +125,38 @@ func TestComputeUsageRepointClosesPriorInterval(t *testing.T) { if !hasBefore || !hasAfter || !hasEnd || before.IntervalID == after.IntervalID || closed.IntervalID != before.IntervalID { t.Fatalf("re-point events = %+v, want distinct starts and end of displaced interval", events) } + if closed.OccurredAt.After(after.OccurredAt) { + t.Fatalf("re-point end occurred at %s, after replacement start at %s", closed.OccurredAt, after.OccurredAt) + } if after.RunnerID != "runner-after" || closed.RunnerID != "runner-before" || closed.OwnerUserID != string(owner.ID) { t.Fatalf("re-point event metadata = %+v, want runner and owner preserved per interval", events) } } +func TestComputeUsageConflictRollsBackIntervalEvents(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-conflict-owner") + accountX := mustAgent(t, s, owner.ID, "compute-conflict-x") + accountY := mustAgent(t, s, owner.ID, "compute-conflict-y") + tenant := s.EffectiveTenant(ctx) + mustBind(t, ctx, s, "sess-1", accountX.ID, "runner-1") + mustBind(t, ctx, s, "sess-2", accountY.ID, "runner-1") + + before := computeEvents(t, s, tenant, accountX.ID) + if _, err := s.RecordSessionBinding(ctx, "sess-2", accountX.ID, "runner-1"); !errors.Is(err, ErrConflict) { + t.Fatalf("RecordSessionBinding(X, sess-2) error = %v, want ErrConflict", err) + } + after := computeEvents(t, s, tenant, accountX.ID) + if len(before) != 1 || len(after) != 1 || after[0].Kind != "start" || after[0].SessionID != "sess-1" || + after[0].IntervalID != before[0].IntervalID || !after[0].OccurredAt.Equal(before[0].OccurredAt) { + t.Fatalf("account X events before=%+v after=%+v, want its unchanged open start only", before, after) + } + if sessionID, err := s.SessionForAccount(ctx, accountX.ID); err != nil || sessionID != "sess-1" { + t.Fatalf("account X resolves to (%q, %v), want sess-1", sessionID, err) + } +} + func TestComputeUsageSameSessionRunnerRebindKeepsInterval(t *testing.T) { ctx := context.Background() s := newTestStore(t) diff --git a/go/internal/store/db/querier.go b/go/internal/store/db/querier.go index 35982368f..d31c84216 100644 --- a/go/internal/store/db/querier.go +++ b/go/internal/store/db/querier.go @@ -124,21 +124,10 @@ type Querier interface { DeleteSecret(ctx context.Context, arg DeleteSecretParams) (int64, error) DeleteServerSecret(ctx context.Context, name string) (int64, error) DeleteSessionBinding(ctx context.Context, sessionID string) error - // The reconnect sweep. Hub.enroll (internal/runnerhub/hub.go:905-957) clears - // every binding when a Runner (re-)enrolls: a reconnecting Runner has no live - // sessions, so a surviving binding would resolve a re-minted session id to a - // stale account. A durable table does not forget on reconnect, so the sweep must - // be explicit. - // - // :many with RETURNING, deliberately NOT :exec. enroll snapshots the bindings - // BEFORE clearing them (hub.go:912-928) because each cleared binding drives a - // presence DISCONNECTED edge (RIG-1569 T8) and each cleared session id must be - // reaped from the delivery held-deliver registry (RIG-1569 T3). A bare DELETE - // would satisfy the invariant while silently dropping both side-effects, leaving - // a long-WORKING agent stuck WORKING in the projection forever. RETURNING is - // what preserves them, so the returned rows are load-bearing, not diagnostic. - // A DELETE ... RETURNING takes no ORDER BY, so the Store method sorts the - // returned slice by session id to keep a sweep pass deterministic and diffable. + // The reconnect sweep, run by Hub.enroll under the system role because a Runner + // is shared across tenants. :many with RETURNING: each removed row drives a + // presence DISCONNECTED edge, a held-deliver reap, and a tenant-scoped archive. + // DELETE ... RETURNING takes no ORDER BY, so the Store method sorts by session id. DeleteSessionBindingsForRunner(ctx context.Context, runnerID string) ([]DeleteSessionBindingsForRunnerRow, error) // DeleteTokenUsageEventsBefore deletes old events of the tx's tenant only, // because row-level security scopes it. The prune runs it once per tenant. @@ -392,39 +381,21 @@ type Querier interface { // compass.tenant_id GUC for RLS and tenant defaults. Delete queries carry the // deleted binding's tenant_id into their end-event insert explicitly. // - // It is NOT complete under WithSystemRole (tenant_tx.go), which arms the - // BYPASSRLS compass_system role and NO tenant GUC. Every query here then runs - // cross-tenant and unscoped, and each one changes meaning: + // Under WithSystemRole (BYPASSRLS, no tenant GUC) these queries run cross-tenant: // // * SessionBindingForAccount / SessionBindingAccount / SessionBindingForUpdate - // are :one, but with RLS gone their single predicate can match rows in - // SEVERAL tenants. pgx's QueryRow takes the first and discards the rest - // without error, so the caller gets a plausible answer from an arbitrary - // tenant. - // * DeleteSessionBindingsForRunner sweeps EVERY tenant's bindings for that - // runner id — runner ids are not tenant-unique either. - // * RecordSessionBinding does not fail closed. tenant_id DEFAULTs to - // current_setting('compass.tenant_id', TRUE); on a pooled connection that - // previously served an ARMED statement, the ended SET LOCAL leaves that - // custom GUC defined-and-EMPTY rather than undefined, so the DEFAULT - // resolves to '' and the NOT NULL is satisfied. The row lands stamped with a - // tenant that does not exist — and tenant_id here has no FK to tenants - // (accounts.tenant_id does), so nothing catches it. No RLS policy matches - // '', so that row is then invisible to every tenant and releasable by - // nothing on the request path. On a connection that never carried an armed - // statement the GUC is genuinely undefined, the DEFAULT is NULL, and the - // insert fails not-null instead — so which of the two a caller gets depends - // on the pooled connection it draws. + // are :one, so a predicate matching rows in several tenants returns an + // arbitrary tenant's row without error. + // * RecordSessionBinding does not fail closed: tenant_id DEFAULTs from the GUC, + // which a pooled connection may leave empty, landing a row no tenant can see. + // * DeleteSessionBindingsForRunner sweeps every tenant's bindings for that + // runner id. Hub.enroll calls it this way on purpose, since a Runner is shared + // across tenants; the returned tenant_id scopes each archive and end event. // - // Nothing calls these under the system role today (WithSystemRole is set at - // delivery/consumer.go and runnerhub/hub.go); a PR3 caller that wants to must - // scope them deliberately rather than inherit scoping from here. - // session_bindings_pgtest_test.go pins the observed behaviour so it is recorded - // rather than latent. + // SessionBindingTenants is the other deliberate system-role read. + // session_bindings_pgtest_test.go pins the :one behaviour. // - // updated_at is NEVER assigned here: the set_updated_at() BEFORE UPDATE trigger - // (0001_init.sql, RIG-3495) is the one mechanism, and a hand-written - // `updated_at = now()` is the exact defect that convention removes. + // updated_at is maintained by the set_updated_at() trigger, never here. // // The per-account serialization the bind takes FIRST, before it reads anything. // Auto-released at transaction end. It mirrors LockOwnerDM / LockOwnerCoordination @@ -598,6 +569,7 @@ type Querier interface { SharesVisibleChannel(ctx context.Context, arg SharesVisibleChannelParams) (bool, error) // Event writes share RecordSessionBinding's transaction, so neither half of an // interval can commit without its binding transition. + // clock_timestamp records after lock waits, unlike now() which uses tx start time. StartComputeUsageInterval(ctx context.Context, arg StartComputeUsageIntervalParams) error StoreForgeRepoWatermark(ctx context.Context, arg StoreForgeRepoWatermarkParams) (int64, error) SubscribeConvertedDMParties(ctx context.Context, channelID string) error diff --git a/go/internal/store/db/session_bindings.sql.go b/go/internal/store/db/session_bindings.sql.go index e5ee3b54c..6aa17ef50 100644 --- a/go/internal/store/db/session_bindings.sql.go +++ b/go/internal/store/db/session_bindings.sql.go @@ -22,7 +22,7 @@ INSERT INTO compute_usage_events ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) -SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), +SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', clock_timestamp(), d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id FROM d JOIN agent_accounts AS a ON a.account_id = d.agent_account_id @@ -45,36 +45,26 @@ WITH d AS ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) - SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', clock_timestamp(), d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id FROM d JOIN agent_accounts AS a ON a.account_id = d.agent_account_id ON CONFLICT DO NOTHING RETURNING 1 ) -SELECT d.session_id, d.agent_account_id FROM d +SELECT d.tenant_id, d.session_id, d.agent_account_id FROM d ` type DeleteSessionBindingsForRunnerRow struct { + TenantID string SessionID string AgentAccountID string } -// The reconnect sweep. Hub.enroll (internal/runnerhub/hub.go:905-957) clears -// every binding when a Runner (re-)enrolls: a reconnecting Runner has no live -// sessions, so a surviving binding would resolve a re-minted session id to a -// stale account. A durable table does not forget on reconnect, so the sweep must -// be explicit. -// -// :many with RETURNING, deliberately NOT :exec. enroll snapshots the bindings -// BEFORE clearing them (hub.go:912-928) because each cleared binding drives a -// presence DISCONNECTED edge (RIG-1569 T8) and each cleared session id must be -// reaped from the delivery held-deliver registry (RIG-1569 T3). A bare DELETE -// would satisfy the invariant while silently dropping both side-effects, leaving -// a long-WORKING agent stuck WORKING in the projection forever. RETURNING is -// what preserves them, so the returned rows are load-bearing, not diagnostic. -// A DELETE ... RETURNING takes no ORDER BY, so the Store method sorts the -// returned slice by session id to keep a sweep pass deterministic and diffable. +// The reconnect sweep, run by Hub.enroll under the system role because a Runner +// is shared across tenants. :many with RETURNING: each removed row drives a +// presence DISCONNECTED edge, a held-deliver reap, and a tenant-scoped archive. +// DELETE ... RETURNING takes no ORDER BY, so the Store method sorts by session id. func (q *Queries) DeleteSessionBindingsForRunner(ctx context.Context, runnerID string) ([]DeleteSessionBindingsForRunnerRow, error) { rows, err := q.db.Query(ctx, deleteSessionBindingsForRunner, runnerID) if err != nil { @@ -84,7 +74,7 @@ func (q *Queries) DeleteSessionBindingsForRunner(ctx context.Context, runnerID s var items []DeleteSessionBindingsForRunnerRow for rows.Next() { var i DeleteSessionBindingsForRunnerRow - if err := rows.Scan(&i.SessionID, &i.AgentAccountID); err != nil { + if err := rows.Scan(&i.TenantID, &i.SessionID, &i.AgentAccountID); err != nil { return nil, err } items = append(items, i) @@ -100,7 +90,7 @@ INSERT INTO compute_usage_events ( id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) -SELECT gen_random_uuid()::text, $1::text, 'end', now(), +SELECT gen_random_uuid()::text, $1::text, 'end', clock_timestamp(), a.account_id, a.owner_user_id, $2::text, $3::text FROM agent_accounts AS a WHERE a.account_id = $4::text @@ -146,39 +136,21 @@ type LockSessionBindingAccountParams struct { // compass.tenant_id GUC for RLS and tenant defaults. Delete queries carry the // deleted binding's tenant_id into their end-event insert explicitly. // -// It is NOT complete under WithSystemRole (tenant_tx.go), which arms the -// BYPASSRLS compass_system role and NO tenant GUC. Every query here then runs -// cross-tenant and unscoped, and each one changes meaning: +// Under WithSystemRole (BYPASSRLS, no tenant GUC) these queries run cross-tenant: // // - SessionBindingForAccount / SessionBindingAccount / SessionBindingForUpdate -// are :one, but with RLS gone their single predicate can match rows in -// SEVERAL tenants. pgx's QueryRow takes the first and discards the rest -// without error, so the caller gets a plausible answer from an arbitrary -// tenant. -// - DeleteSessionBindingsForRunner sweeps EVERY tenant's bindings for that -// runner id — runner ids are not tenant-unique either. -// - RecordSessionBinding does not fail closed. tenant_id DEFAULTs to -// current_setting('compass.tenant_id', TRUE); on a pooled connection that -// previously served an ARMED statement, the ended SET LOCAL leaves that -// custom GUC defined-and-EMPTY rather than undefined, so the DEFAULT -// resolves to ” and the NOT NULL is satisfied. The row lands stamped with a -// tenant that does not exist — and tenant_id here has no FK to tenants -// (accounts.tenant_id does), so nothing catches it. No RLS policy matches -// ”, so that row is then invisible to every tenant and releasable by -// nothing on the request path. On a connection that never carried an armed -// statement the GUC is genuinely undefined, the DEFAULT is NULL, and the -// insert fails not-null instead — so which of the two a caller gets depends -// on the pooled connection it draws. +// are :one, so a predicate matching rows in several tenants returns an +// arbitrary tenant's row without error. +// - RecordSessionBinding does not fail closed: tenant_id DEFAULTs from the GUC, +// which a pooled connection may leave empty, landing a row no tenant can see. +// - DeleteSessionBindingsForRunner sweeps every tenant's bindings for that +// runner id. Hub.enroll calls it this way on purpose, since a Runner is shared +// across tenants; the returned tenant_id scopes each archive and end event. // -// Nothing calls these under the system role today (WithSystemRole is set at -// delivery/consumer.go and runnerhub/hub.go); a PR3 caller that wants to must -// scope them deliberately rather than inherit scoping from here. -// session_bindings_pgtest_test.go pins the observed behaviour so it is recorded -// rather than latent. +// SessionBindingTenants is the other deliberate system-role read. +// session_bindings_pgtest_test.go pins the :one behaviour. // -// updated_at is NEVER assigned here: the set_updated_at() BEFORE UPDATE trigger -// (0001_init.sql, RIG-3495) is the one mechanism, and a hand-written -// `updated_at = now()` is the exact defect that convention removes. +// updated_at is maintained by the set_updated_at() trigger, never here. // // The per-account serialization the bind takes FIRST, before it reads anything. // Auto-released at transaction end. It mirrors LockOwnerDM / LockOwnerCoordination @@ -334,7 +306,7 @@ INSERT INTO compute_usage_events ( id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) -SELECT gen_random_uuid()::text, $1::text, 'start', now(), +SELECT gen_random_uuid()::text, $1::text, 'start', clock_timestamp(), a.account_id, a.owner_user_id, $2::text, $3::text FROM agent_accounts AS a WHERE a.account_id = $4::text @@ -350,6 +322,7 @@ type StartComputeUsageIntervalParams struct { // Event writes share RecordSessionBinding's transaction, so neither half of an // interval can commit without its binding transition. +// clock_timestamp records after lock waits, unlike now() which uses tx start time. func (q *Queries) StartComputeUsageInterval(ctx context.Context, arg StartComputeUsageIntervalParams) error { _, err := q.db.Exec(ctx, startComputeUsageInterval, arg.IntervalID, diff --git a/go/internal/store/migrate_upgrade_pgtest_test.go b/go/internal/store/migrate_upgrade_pgtest_test.go index b86b76934..9aa48f72f 100644 --- a/go/internal/store/migrate_upgrade_pgtest_test.go +++ b/go/internal/store/migrate_upgrade_pgtest_test.go @@ -64,7 +64,7 @@ func TestOpenUpgradesV1DatabaseToTokenUsage(t *testing.T) { } if err == nil { _, err = seed.Exec(ctx, - "INSERT INTO session_bindings (tenant_id, agent_account_id, session_id, runner_id) VALUES ($1, $2, 'upgrade-session', 'upgrade-runner')", + "INSERT INTO session_bindings (tenant_id, agent_account_id, session_id, runner_id, updated_at) VALUES ($1, $2, 'upgrade-session', 'upgrade-runner', now() - interval '3 days')", tenant, agentID, ) } @@ -106,20 +106,21 @@ func TestOpenUpgradesV1DatabaseToTokenUsage(t *testing.T) { var backfilledStartCount int var backfilledInterval, boundInterval string - var allEstimated bool + var allEstimated, timestampMatches bool if err := s.pool.QueryRow(ctx, ` - SELECT count(*), min(e.interval_id), bool_and(e.estimated), min(b.usage_interval_id) + SELECT count(*), min(e.interval_id), bool_and(e.estimated), min(b.usage_interval_id), + bool_and(e.occurred_at = b.updated_at) FROM compute_usage_events AS e JOIN session_bindings AS b ON b.tenant_id = e.tenant_id AND b.agent_account_id = e.agent_account_id WHERE e.tenant_id = $1 AND e.agent_account_id = $2 AND e.kind = 'start' AND e.session_id = 'upgrade-session'`, string(tenant), agentID). - Scan(&backfilledStartCount, &backfilledInterval, &allEstimated, &boundInterval); err != nil { + Scan(&backfilledStartCount, &backfilledInterval, &allEstimated, &boundInterval, ×tampMatches); err != nil { t.Fatalf("read backfilled compute start: %v", err) } - if backfilledStartCount != 1 || backfilledInterval == "" || backfilledInterval != boundInterval || !allEstimated { - t.Fatalf("backfilled starts = %d, interval = %q, bound = %q, estimated = %t; want one estimated start on the bound interval", - backfilledStartCount, backfilledInterval, boundInterval, allEstimated) + if backfilledStartCount != 1 || backfilledInterval == "" || backfilledInterval != boundInterval || !allEstimated || !timestampMatches { + t.Fatalf("backfilled starts = %d, interval = %q, bound = %q, estimated = %t, timestamp matches = %t; want one estimated start on the bound interval at its seeded updated_at", + backfilledStartCount, backfilledInterval, boundInterval, allEstimated, timestampMatches) } var horizon pgtype.Timestamptz diff --git a/go/internal/store/migrations/0004_compute_usage.sql b/go/internal/store/migrations/0004_compute_usage.sql index 1d7b7b60c..c61af3db9 100644 --- a/go/internal/store/migrations/0004_compute_usage.sql +++ b/go/internal/store/migrations/0004_compute_usage.sql @@ -18,10 +18,13 @@ CREATE TABLE compute_usage_events ( CREATE INDEX compute_usage_events_occurred_at_idx ON compute_usage_events (tenant_id, occurred_at); -ALTER TABLE session_bindings ADD COLUMN usage_interval_id TEXT; +-- Keep the default so older servers can insert bindings during a rolling deploy. +-- squawk-ignore adding-field-with-default +ALTER TABLE session_bindings ADD COLUMN usage_interval_id TEXT NOT NULL DEFAULT gen_random_uuid()::TEXT; -UPDATE session_bindings AS b - SET usage_interval_id = gen_random_uuid()::TEXT; +GRANT SELECT, INSERT ON compute_usage_events TO compass_app, compass_system; + +SET LOCAL ROLE compass_system; INSERT INTO compute_usage_events ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, @@ -35,11 +38,7 @@ SELECT b.tenant_id, gen_random_uuid()::TEXT AS id, b.usage_interval_id, 'start' ON a.account_id = b.agent_account_id AND a.tenant_id = b.tenant_id; --- Every existing row got an id above, so the cutover cannot fail. --- squawk-ignore adding-not-nullable-field -ALTER TABLE session_bindings ALTER COLUMN usage_interval_id SET NOT NULL; - -GRANT SELECT, INSERT ON compute_usage_events TO compass_app, compass_system; +RESET ROLE; DO $$ DECLARE diff --git a/go/internal/store/queries/session_bindings.sql b/go/internal/store/queries/session_bindings.sql index dac0a8c02..7c394648f 100644 --- a/go/internal/store/queries/session_bindings.sql +++ b/go/internal/store/queries/session_bindings.sql @@ -8,39 +8,21 @@ -- compass.tenant_id GUC for RLS and tenant defaults. Delete queries carry the -- deleted binding's tenant_id into their end-event insert explicitly. -- --- It is NOT complete under WithSystemRole (tenant_tx.go), which arms the --- BYPASSRLS compass_system role and NO tenant GUC. Every query here then runs --- cross-tenant and unscoped, and each one changes meaning: +-- Under WithSystemRole (BYPASSRLS, no tenant GUC) these queries run cross-tenant: -- -- * SessionBindingForAccount / SessionBindingAccount / SessionBindingForUpdate --- are :one, but with RLS gone their single predicate can match rows in --- SEVERAL tenants. pgx's QueryRow takes the first and discards the rest --- without error, so the caller gets a plausible answer from an arbitrary --- tenant. --- * DeleteSessionBindingsForRunner sweeps EVERY tenant's bindings for that --- runner id — runner ids are not tenant-unique either. --- * RecordSessionBinding does not fail closed. tenant_id DEFAULTs to --- current_setting('compass.tenant_id', TRUE); on a pooled connection that --- previously served an ARMED statement, the ended SET LOCAL leaves that --- custom GUC defined-and-EMPTY rather than undefined, so the DEFAULT --- resolves to '' and the NOT NULL is satisfied. The row lands stamped with a --- tenant that does not exist — and tenant_id here has no FK to tenants --- (accounts.tenant_id does), so nothing catches it. No RLS policy matches --- '', so that row is then invisible to every tenant and releasable by --- nothing on the request path. On a connection that never carried an armed --- statement the GUC is genuinely undefined, the DEFAULT is NULL, and the --- insert fails not-null instead — so which of the two a caller gets depends --- on the pooled connection it draws. +-- are :one, so a predicate matching rows in several tenants returns an +-- arbitrary tenant's row without error. +-- * RecordSessionBinding does not fail closed: tenant_id DEFAULTs from the GUC, +-- which a pooled connection may leave empty, landing a row no tenant can see. +-- * DeleteSessionBindingsForRunner sweeps every tenant's bindings for that +-- runner id. Hub.enroll calls it this way on purpose, since a Runner is shared +-- across tenants; the returned tenant_id scopes each archive and end event. -- --- Nothing calls these under the system role today (WithSystemRole is set at --- delivery/consumer.go and runnerhub/hub.go); a PR3 caller that wants to must --- scope them deliberately rather than inherit scoping from here. --- session_bindings_pgtest_test.go pins the observed behaviour so it is recorded --- rather than latent. +-- SessionBindingTenants is the other deliberate system-role read. +-- session_bindings_pgtest_test.go pins the :one behaviour. -- --- updated_at is NEVER assigned here: the set_updated_at() BEFORE UPDATE trigger --- (0001_init.sql, RIG-3495) is the one mechanism, and a hand-written --- `updated_at = now()` is the exact defect that convention removes. +-- updated_at is maintained by the set_updated_at() trigger, never here. -- -- The per-account serialization the bind takes FIRST, before it reads anything. @@ -91,12 +73,13 @@ SELECT b.session_id, b.usage_interval_id, b.runner_id -- Event writes share RecordSessionBinding's transaction, so neither half of an -- interval can commit without its binding transition. +-- clock_timestamp records after lock waits, unlike now() which uses tx start time. -- name: StartComputeUsageInterval :exec INSERT INTO compute_usage_events ( id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) -SELECT gen_random_uuid()::text, @interval_id::text, 'start', now(), +SELECT gen_random_uuid()::text, @interval_id::text, 'start', clock_timestamp(), a.account_id, a.owner_user_id, @session_id::text, @runner_id::text FROM agent_accounts AS a WHERE a.account_id = @agent_account_id::text @@ -109,7 +92,7 @@ INSERT INTO compute_usage_events ( id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) -SELECT gen_random_uuid()::text, @interval_id::text, 'end', now(), +SELECT gen_random_uuid()::text, @interval_id::text, 'end', clock_timestamp(), a.account_id, a.owner_user_id, @session_id::text, @runner_id::text FROM agent_accounts AS a WHERE a.account_id = @agent_account_id::text @@ -142,27 +125,16 @@ INSERT INTO compute_usage_events ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) -SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), +SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', clock_timestamp(), d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id FROM d JOIN agent_accounts AS a ON a.account_id = d.agent_account_id ON CONFLICT DO NOTHING; --- The reconnect sweep. Hub.enroll (internal/runnerhub/hub.go:905-957) clears --- every binding when a Runner (re-)enrolls: a reconnecting Runner has no live --- sessions, so a surviving binding would resolve a re-minted session id to a --- stale account. A durable table does not forget on reconnect, so the sweep must --- be explicit. --- --- :many with RETURNING, deliberately NOT :exec. enroll snapshots the bindings --- BEFORE clearing them (hub.go:912-928) because each cleared binding drives a --- presence DISCONNECTED edge (RIG-1569 T8) and each cleared session id must be --- reaped from the delivery held-deliver registry (RIG-1569 T3). A bare DELETE --- would satisfy the invariant while silently dropping both side-effects, leaving --- a long-WORKING agent stuck WORKING in the projection forever. RETURNING is --- what preserves them, so the returned rows are load-bearing, not diagnostic. --- A DELETE ... RETURNING takes no ORDER BY, so the Store method sorts the --- returned slice by session id to keep a sweep pass deterministic and diffable. +-- The reconnect sweep, run by Hub.enroll under the system role because a Runner +-- is shared across tenants. :many with RETURNING: each removed row drives a +-- presence DISCONNECTED edge, a held-deliver reap, and a tenant-scoped archive. +-- DELETE ... RETURNING takes no ORDER BY, so the Store method sorts by session id. -- name: DeleteSessionBindingsForRunner :many WITH d AS ( DELETE FROM session_bindings AS b @@ -174,14 +146,14 @@ WITH d AS ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, session_id, runner_id ) - SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', now(), + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', clock_timestamp(), d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id FROM d JOIN agent_accounts AS a ON a.account_id = d.agent_account_id ON CONFLICT DO NOTHING RETURNING 1 ) -SELECT d.session_id, d.agent_account_id FROM d; +SELECT d.tenant_id, d.session_id, d.agent_account_id FROM d; -- The one query here meant for the system role: a Runner-originated call carries -- no tenant, so the hub reads the session's tenant cross-tenant, then acts under it. diff --git a/go/internal/store/session_bindings.go b/go/internal/store/session_bindings.go index 8d60ba860..9283b59ce 100644 --- a/go/internal/store/session_bindings.go +++ b/go/internal/store/session_bindings.go @@ -27,11 +27,10 @@ import ( // PR2 adds the table and methods only; demoting the hub's in-RAM maps is PR3. -// SessionBinding is one live binding: the session, the agent account it speaks -// for, and the Runner it is attached to. Returned by -// DeleteSessionBindingsForRunner, which is the reconnect sweep — its caller -// needs every field of each binding it just removed. +// SessionBinding is one live binding, including its tenant. Returned by the +// reconnect sweep so callers can scope follow-up work to the deleted row. type SessionBinding struct { + TenantID TenantID SessionID string AccountID AccountID RunnerID string @@ -285,6 +284,7 @@ func (s *Store) DeleteSessionBindingsForRunner(ctx context.Context, runnerID str bindings := make([]SessionBinding, 0, len(rows)) for _, row := range rows { bindings = append(bindings, SessionBinding{ + TenantID: TenantID(row.TenantID), SessionID: row.SessionID, AccountID: AccountID(row.AgentAccountID), RunnerID: runnerID, diff --git a/go/internal/store/session_bindings_pgtest_test.go b/go/internal/store/session_bindings_pgtest_test.go index f313002f1..6a808614f 100644 --- a/go/internal/store/session_bindings_pgtest_test.go +++ b/go/internal/store/session_bindings_pgtest_test.go @@ -356,11 +356,12 @@ func TestDeleteSessionBindingsForRunnerReturnsEverySweptBinding(t *testing.T) { z := mustAgent(t, s, owner.ID, "agent-z") elsewhere := mustAgent(t, s, owner.ID, "agent-elsewhere") + tenant := s.EffectiveTenant(ctx) for _, bind := range []SessionBinding{ - {SessionID: "sess-z", AccountID: z.ID, RunnerID: "runner-1"}, - {SessionID: "sess-b", AccountID: b.ID, RunnerID: "runner-1"}, - {SessionID: "sess-a", AccountID: a.ID, RunnerID: "runner-1"}, - {SessionID: "sess-elsewhere", AccountID: elsewhere.ID, RunnerID: "runner-2"}, + {TenantID: tenant, SessionID: "sess-z", AccountID: z.ID, RunnerID: "runner-1"}, + {TenantID: tenant, SessionID: "sess-b", AccountID: b.ID, RunnerID: "runner-1"}, + {TenantID: tenant, SessionID: "sess-a", AccountID: a.ID, RunnerID: "runner-1"}, + {TenantID: tenant, SessionID: "sess-elsewhere", AccountID: elsewhere.ID, RunnerID: "runner-2"}, } { mustBind(t, ctx, s, bind.SessionID, bind.AccountID, bind.RunnerID) } @@ -374,9 +375,9 @@ func TestDeleteSessionBindingsForRunnerReturnsEverySweptBinding(t *testing.T) { // every swept binding, with the account each DISCONNECTED edge needs, in // sorted order (which is NOT the order they were seeded in). want := []SessionBinding{ - {SessionID: "sess-a", AccountID: a.ID, RunnerID: "runner-1"}, - {SessionID: "sess-b", AccountID: b.ID, RunnerID: "runner-1"}, - {SessionID: "sess-z", AccountID: z.ID, RunnerID: "runner-1"}, + {TenantID: s.EffectiveTenant(ctx), SessionID: "sess-a", AccountID: a.ID, RunnerID: "runner-1"}, + {TenantID: s.EffectiveTenant(ctx), SessionID: "sess-b", AccountID: b.ID, RunnerID: "runner-1"}, + {TenantID: s.EffectiveTenant(ctx), SessionID: "sess-z", AccountID: z.ID, RunnerID: "runner-1"}, } if len(swept) != len(want) { t.Fatalf("swept = %+v, want exactly the %d bindings on runner-1 (a :exec sweep would return none)", swept, len(want)) From 1a8c3a3ca8e1bfe0a883f8c5d0f090a8c4056f14 Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 05:11:38 -0400 Subject: [PATCH 03/10] fix(usage): log an estimated start for bindings an older server wrote (RIG-2872) During a rolling deploy an older server inserts bindings that have no start event. A same-session rebind, a single release, or a runner sweep now first writes an estimated start (at the binding's created_at) for such a row, so the interval stays complete. ON CONFLICT keeps any real start. Tests seed a binding without events and exercise rebind then release, rebind then runner sweep, and release alone. All three fail without this change. Refs RIG-2872 Co-authored-by: Matt Wilkinson --- .../store/compute_usage_pgtest_test.go | 56 +++++++++++++++++++ go/internal/store/db/querier.go | 5 ++ go/internal/store/db/session_bindings.sql.go | 48 +++++++++++++++- .../store/queries/session_bindings.sql | 42 +++++++++++++- go/internal/store/session_bindings.go | 6 ++ 5 files changed, 153 insertions(+), 4 deletions(-) diff --git a/go/internal/store/compute_usage_pgtest_test.go b/go/internal/store/compute_usage_pgtest_test.go index e2da565c5..6f12048bd 100644 --- a/go/internal/store/compute_usage_pgtest_test.go +++ b/go/internal/store/compute_usage_pgtest_test.go @@ -175,6 +175,62 @@ func TestComputeUsageSameSessionRunnerRebindKeepsInterval(t *testing.T) { } } +// A binding an older server wrote mid-deploy has no start event. A later rebind +// or release must still leave one complete interval, with its start estimated. +func TestComputeUsageLegacyBindingGetsEstimatedStart(t *testing.T) { + for name, release := range map[string]func(*testing.T, context.Context, *Store){ + "single release": func(t *testing.T, ctx context.Context, s *Store) { + if err := s.DeleteSessionBinding(ctx, "legacy-session"); err != nil { + t.Fatalf("DeleteSessionBinding: %v", err) + } + }, + "runner sweep": func(t *testing.T, ctx context.Context, s *Store) { + if _, err := s.DeleteSessionBindingsForRunner(ctx, "runner-new"); err != nil { + t.Fatalf("DeleteSessionBindingsForRunner: %v", err) + } + }, + } { + t.Run(name, func(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-legacy-owner") + agent := mustAgent(t, s, owner.ID, "compute-legacy-agent") + tenant := s.EffectiveTenant(ctx) + // The old INSERT omits usage_interval_id, so the column default fills it. + execAsSystem(t, s, "INSERT INTO session_bindings (tenant_id, agent_account_id, session_id, runner_id) VALUES ($1, $2, 'legacy-session', 'runner-old')", + string(tenant), string(agent.ID)) + + mustBind(t, ctx, s, "legacy-session", agent.ID, "runner-new") + release(t, ctx, s) + + events := computeEvents(t, s, tenant, agent.ID) + if len(events) != 2 || events[0].Kind != "start" || events[1].Kind != "end" || + events[0].IntervalID != events[1].IntervalID || !events[0].Estimated || events[1].Estimated { + t.Fatalf("legacy binding events = %+v, want one estimated start and one exact end", events) + } + }) + } +} + +// A legacy binding released before any new-server rebind still logs its interval. +func TestComputeUsageLegacyBindingReleaseLogsInterval(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-legacy-release-owner") + agent := mustAgent(t, s, owner.ID, "compute-legacy-release-agent") + tenant := s.EffectiveTenant(ctx) + execAsSystem(t, s, "INSERT INTO session_bindings (tenant_id, agent_account_id, session_id, runner_id) VALUES ($1, $2, 'legacy-only', 'runner-old')", + string(tenant), string(agent.ID)) + if err := s.DeleteSessionBinding(ctx, "legacy-only"); err != nil { + t.Fatalf("DeleteSessionBinding: %v", err) + } + events := computeEvents(t, s, tenant, agent.ID) + if len(events) != 2 || events[0].Kind != "start" || events[1].Kind != "end" || + events[0].IntervalID != events[1].IntervalID || !events[0].Estimated { + t.Fatalf("legacy release events = %+v, want an estimated start and its end", events) + } +} + func TestComputeUsageRunnerSweepClosesIntervalsAndReturnsBindings(t *testing.T) { ctx := context.Background() s := newTestStore(t) diff --git a/go/internal/store/db/querier.go b/go/internal/store/db/querier.go index d31c84216..9e636669e 100644 --- a/go/internal/store/db/querier.go +++ b/go/internal/store/db/querier.go @@ -123,6 +123,8 @@ type Querier interface { // identifies a row (composite PK). DeleteSecret(ctx context.Context, arg DeleteSecretParams) (int64, error) DeleteServerSecret(ctx context.Context, name string) (int64, error) + // Both deletes also write an estimated start for a binding an older server made + // without one; ON CONFLICT keeps any real start. DeleteSessionBinding(ctx context.Context, sessionID string) error // The reconnect sweep, run by Hub.enroll under the system role because a Runner // is shared across tenants. :many with RETURNING: each removed row drives a @@ -154,6 +156,9 @@ type Querier interface { // itself) makes RETURNING fire on conflict so a repeat returns the stored id. EnsureAgentForgeSubscription(ctx context.Context, arg EnsureAgentForgeSubscriptionParams) (string, error) EnsureChannelMember(ctx context.Context, arg EnsureChannelMemberParams) error + // A binding an older server wrote during a rolling deploy has no start event. + // Its created_at is the best start we hold, so the start is marked estimated. + EnsureComputeUsageIntervalStart(ctx context.Context, agentAccountID string) error EnsureForgeRepoSubscription(ctx context.Context, arg EnsureForgeRepoSubscriptionParams) error FindAskMessage(ctx context.Context, arg FindAskMessageParams) ([]FindAskMessageRow, error) // Collects the coordinate's cursor IFF no subscription for it remains (the NOT diff --git a/go/internal/store/db/session_bindings.sql.go b/go/internal/store/db/session_bindings.sql.go index 6aa17ef50..05fed17fa 100644 --- a/go/internal/store/db/session_bindings.sql.go +++ b/go/internal/store/db/session_bindings.sql.go @@ -16,7 +16,18 @@ WITH d AS ( DELETE FROM session_bindings AS b WHERE b.session_id = $1 RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, - b.session_id, b.runner_id + b.session_id, b.runner_id, b.created_at +), starts AS ( + INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id, estimated + ) + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'start', d.created_at, + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id, TRUE + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id + ON CONFLICT DO NOTHING + RETURNING 1 ) INSERT INTO compute_usage_events ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, @@ -29,6 +40,8 @@ SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'end', clock_t ON CONFLICT DO NOTHING ` +// Both deletes also write an estimated start for a binding an older server made +// without one; ON CONFLICT keeps any real start. func (q *Queries) DeleteSessionBinding(ctx context.Context, sessionID string) error { _, err := q.db.Exec(ctx, deleteSessionBinding, sessionID) return err @@ -39,7 +52,18 @@ WITH d AS ( DELETE FROM session_bindings AS b WHERE b.runner_id = $1 RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, - b.session_id, b.runner_id + b.session_id, b.runner_id, b.created_at +), starts AS ( + INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id, estimated + ) + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'start', d.created_at, + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id, TRUE + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id + ON CONFLICT DO NOTHING + RETURNING 1 ), ins AS ( INSERT INTO compute_usage_events ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, @@ -116,6 +140,26 @@ func (q *Queries) EndComputeUsageInterval(ctx context.Context, arg EndComputeUsa return err } +const ensureComputeUsageIntervalStart = `-- name: EnsureComputeUsageIntervalStart :exec +INSERT INTO compute_usage_events ( + id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, + session_id, runner_id, estimated +) +SELECT gen_random_uuid()::text, b.usage_interval_id, 'start', b.created_at, + b.agent_account_id, a.owner_user_id, b.session_id, b.runner_id, TRUE + FROM session_bindings AS b + JOIN agent_accounts AS a ON a.account_id = b.agent_account_id + WHERE b.agent_account_id = $1::text +ON CONFLICT (tenant_id, interval_id, kind) DO NOTHING +` + +// A binding an older server wrote during a rolling deploy has no start event. +// Its created_at is the best start we hold, so the start is marked estimated. +func (q *Queries) EnsureComputeUsageIntervalStart(ctx context.Context, agentAccountID string) error { + _, err := q.db.Exec(ctx, ensureComputeUsageIntervalStart, agentAccountID) + return err +} + const lockSessionBindingAccount = `-- name: LockSessionBindingAccount :exec SELECT pg_advisory_xact_lock(hashtext('binding:' || $1 || ':' || $2)) diff --git a/go/internal/store/queries/session_bindings.sql b/go/internal/store/queries/session_bindings.sql index 7c394648f..040dc8d14 100644 --- a/go/internal/store/queries/session_bindings.sql +++ b/go/internal/store/queries/session_bindings.sql @@ -98,6 +98,20 @@ SELECT gen_random_uuid()::text, @interval_id::text, 'end', clock_timestamp(), WHERE a.account_id = @agent_account_id::text ON CONFLICT (tenant_id, interval_id, kind) DO NOTHING; +-- A binding an older server wrote during a rolling deploy has no start event. +-- Its created_at is the best start we hold, so the start is marked estimated. +-- name: EnsureComputeUsageIntervalStart :exec +INSERT INTO compute_usage_events ( + id, interval_id, kind, occurred_at, agent_account_id, owner_user_id, + session_id, runner_id, estimated +) +SELECT gen_random_uuid()::text, b.usage_interval_id, 'start', b.created_at, + b.agent_account_id, a.owner_user_id, b.session_id, b.runner_id, TRUE + FROM session_bindings AS b + JOIN agent_accounts AS a ON a.account_id = b.agent_account_id + WHERE b.agent_account_id = @agent_account_id::text +ON CONFLICT (tenant_id, interval_id, kind) DO NOTHING; + -- What it DISPLACED comes from SessionBindingForUpdate above, not from a -- RETURNING here. The binding update and event writes share the Store tx. -- name: RecordSessionBinding :exec @@ -114,12 +128,25 @@ SELECT agent_account_id, runner_id FROM session_bindings WHERE session_id = $1; -- name: SessionBindingForAccount :one SELECT session_id FROM session_bindings WHERE agent_account_id = $1; +-- Both deletes also write an estimated start for a binding an older server made +-- without one; ON CONFLICT keeps any real start. -- name: DeleteSessionBinding :exec WITH d AS ( DELETE FROM session_bindings AS b WHERE b.session_id = $1 RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, - b.session_id, b.runner_id + b.session_id, b.runner_id, b.created_at +), starts AS ( + INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id, estimated + ) + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'start', d.created_at, + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id, TRUE + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id + ON CONFLICT DO NOTHING + RETURNING 1 ) INSERT INTO compute_usage_events ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, @@ -140,7 +167,18 @@ WITH d AS ( DELETE FROM session_bindings AS b WHERE b.runner_id = $1 RETURNING b.tenant_id, b.usage_interval_id, b.agent_account_id, - b.session_id, b.runner_id + b.session_id, b.runner_id, b.created_at +), starts AS ( + INSERT INTO compute_usage_events ( + tenant_id, id, interval_id, kind, occurred_at, agent_account_id, + owner_user_id, session_id, runner_id, estimated + ) + SELECT d.tenant_id, gen_random_uuid()::text, d.usage_interval_id, 'start', d.created_at, + d.agent_account_id, a.owner_user_id, d.session_id, d.runner_id, TRUE + FROM d + JOIN agent_accounts AS a ON a.account_id = d.agent_account_id + ON CONFLICT DO NOTHING + RETURNING 1 ), ins AS ( INSERT INTO compute_usage_events ( tenant_id, id, interval_id, kind, occurred_at, agent_account_id, diff --git a/go/internal/store/session_bindings.go b/go/internal/store/session_bindings.go index 9283b59ce..140041fe4 100644 --- a/go/internal/store/session_bindings.go +++ b/go/internal/store/session_bindings.go @@ -142,6 +142,12 @@ func (s *Store) RecordSessionBinding(ctx context.Context, sessionID string, acco displaced := prior.SessionID intervalID := prior.UsageIntervalID + if intervalID != "" { + // An older server may have written the prior row without a start event. + if err := qtx.EnsureComputeUsageIntervalStart(ctx, string(accountID)); err != nil { + return "", fmt.Errorf("store: ensure compute usage interval start: %w", err) + } + } if prior.SessionID != sessionID { if intervalID != "" { if err := qtx.EndComputeUsageInterval(ctx, db.EndComputeUsageIntervalParams{ From 0d78ce887f5ae731412ece7df80cca8ba5e5d254 Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 06:01:48 -0400 Subject: [PATCH 04/10] test(usage): pin lock-wait stamping and legacy re-points; document the enroll reap role (RIG-2872) - TestComputeUsageEventsStampedAfterLockWait gates a bind on the account advisory lock and asserts its events are stamped after release. It fails if the queries use now(). - TestComputeUsageLegacyBindingRepointKeepsBothIntervals re-points an event-less binding to a new session. It fails without the ensure-start call. - tenant_tx.go lists Hub.enroll's reap among system-role entrypoints; hub.go says why it is cross-tenant. - docs: inferred timestamps carry the estimated flag; drop the unqualified exact claim. Rolling-deploy coverage for old-server writes is deferred to RIG-4231. Refs RIG-2872 Co-authored-by: Matt Wilkinson --- docs/concepts/tokens-and-billing.md | 4 +- go/internal/runnerhub/hub.go | 2 + .../store/compute_usage_pgtest_test.go | 100 ++++++++++++++++++ go/internal/store/tenant_tx.go | 12 +-- 4 files changed, 111 insertions(+), 7 deletions(-) diff --git a/docs/concepts/tokens-and-billing.md b/docs/concepts/tokens-and-billing.md index 482acff20..60e7a9187 100644 --- a/docs/concepts/tokens-and-billing.md +++ b/docs/concepts/tokens-and-billing.md @@ -57,7 +57,9 @@ health/quota signal. its start to end. The server emits both events with the binding transition; an hourly sweep marks an inferred end as estimated if a binding disappears without its end event. Intervals already open when the log was introduced get - an estimated start. The event log is exact, auditable, and reconstructable. + an estimated start. Every inferred timestamp carries the `estimated` flag; + all others are recorded at the transition. The log is auditable and + reconstructable. - **LLM token usage and spend — recorded for display, not billed day-1.** Every model call's tokens-in/out and cost, captured at the gateway. It powers the in-product usage/spend charts the user sees, and it is *recorded* even though diff --git a/go/internal/runnerhub/hub.go b/go/internal/runnerhub/hub.go index 9b15e60b1..24caffda7 100644 --- a/go/internal/runnerhub/hub.go +++ b/go/internal/runnerhub/hub.go @@ -1028,6 +1028,8 @@ func (h *Hub) enroll(ctx context.Context, id string, subject store.Subject, tier var durableReaped []store.SessionBinding durableReapSucceeded := false if bindings != nil { + // Cross-tenant on purpose: a Runner serves every tenant, and id is the + // authenticated token subject, so the sweep reaches only this Runner's rows. rows, err := bindings.DeleteSessionBindingsForRunner(store.WithSystemRole(ctx), id) if err != nil { // A durable-reap fault must not wedge reconnect; fall back to the in-RAM diff --git a/go/internal/store/compute_usage_pgtest_test.go b/go/internal/store/compute_usage_pgtest_test.go index 6f12048bd..dcf4e7ddf 100644 --- a/go/internal/store/compute_usage_pgtest_test.go +++ b/go/internal/store/compute_usage_pgtest_test.go @@ -133,6 +133,80 @@ func TestComputeUsageRepointClosesPriorInterval(t *testing.T) { } } +// A bind that waits on the account lock must stamp its events after the wait, +// so the interval it opens never starts before the gate released. +func TestComputeUsageEventsStampedAfterLockWait(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-wait-owner") + agent := mustAgent(t, s, owner.ID, "compute-wait-agent") + tenant := s.EffectiveTenant(ctx) + mustBind(t, ctx, s, "wait-before", agent.ID, "runner-1") + + gate, err := s.pool.Begin(ctx) + if err != nil { + t.Fatalf("begin gate tx: %v", err) + } + defer func() { + if err := gate.Rollback(ctx); err != nil && !errors.Is(err, pgx.ErrTxClosed) { + t.Errorf("rollback gate: %v", err) + } + }() + // Character-identical to LockSessionBindingAccount. + if _, err := gate.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtext('binding:' || $1 || ':' || $2))`, + string(tenant), string(agent.ID)); err != nil { + t.Fatalf("gate advisory lock: %v", err) + } + done := make(chan error, 1) + go func() { + _, err := s.RecordSessionBinding(ctx, "wait-after", agent.ID, "runner-1") + done <- err + }() + deadline := time.After(5 * time.Second) + tick := time.NewTicker(10 * time.Millisecond) + defer tick.Stop() + for { + var waiters int + if err := gate.QueryRow(ctx, + `WITH k AS (SELECT hashtext('binding:' || $1 || ':' || $2)::bigint AS key) + SELECT count(*) FROM pg_locks, k + WHERE locktype = 'advisory' AND NOT granted + AND classid = ((k.key >> 32) & 4294967295)::oid + AND objid = (k.key & 4294967295)::oid`, + string(tenant), string(agent.ID)).Scan(&waiters); err != nil { + t.Fatalf("poll pg_locks: %v", err) + } + if waiters >= 1 { + break + } + select { + case err := <-done: + t.Fatalf("bind finished (err=%v) while the gate held the account lock", err) + case <-deadline: + t.Fatal("bind never waited on the account lock") + case <-tick.C: + } + } + var released time.Time + if err := gate.QueryRow(ctx, "SELECT clock_timestamp()").Scan(&released); err != nil { + t.Fatalf("read release time: %v", err) + } + if err := gate.Rollback(ctx); err != nil { + t.Fatalf("release gate: %v", err) + } + if err := <-done; err != nil { + t.Fatalf("RecordSessionBinding after release: %v", err) + } + + for _, event := range computeEvents(t, s, tenant, agent.ID) { + if event.SessionID == "wait-after" || event.Kind == "end" { + if event.OccurredAt.Before(released) { + t.Fatalf("%s event for %s stamped %s, before the lock release at %s", event.Kind, event.SessionID, event.OccurredAt, released) + } + } + } +} + func TestComputeUsageConflictRollsBackIntervalEvents(t *testing.T) { ctx := context.Background() s := newTestStore(t) @@ -231,6 +305,32 @@ func TestComputeUsageLegacyBindingReleaseLogsInterval(t *testing.T) { } } +// Re-pointing a legacy binding to a new session ends the legacy interval with an +// estimated start and opens an exact replacement. +func TestComputeUsageLegacyBindingRepointKeepsBothIntervals(t *testing.T) { + ctx := context.Background() + s := newTestStore(t) + owner := mustUser(t, s, "compute-legacy-repoint-owner") + agent := mustAgent(t, s, owner.ID, "compute-legacy-repoint-agent") + tenant := s.EffectiveTenant(ctx) + execAsSystem(t, s, "INSERT INTO session_bindings (tenant_id, agent_account_id, session_id, runner_id) VALUES ($1, $2, 'legacy-old', 'runner-old')", + string(tenant), string(agent.ID)) + + mustBind(t, ctx, s, "legacy-new", agent.ID, "runner-new") + + byKey := map[string]computeUsageEvent{} + for _, event := range computeEvents(t, s, tenant, agent.ID) { + byKey[event.SessionID+"/"+event.Kind] = event + } + oldStart, okOldStart := byKey["legacy-old/start"] + oldEnd, okOldEnd := byKey["legacy-old/end"] + newStart, okNewStart := byKey["legacy-new/start"] + if len(byKey) != 3 || !okOldStart || !okOldEnd || !okNewStart || + oldStart.IntervalID != oldEnd.IntervalID || !oldStart.Estimated || oldEnd.Estimated || newStart.Estimated { + t.Fatalf("legacy re-point events = %+v, want estimated legacy start, exact end, exact replacement start", byKey) + } +} + func TestComputeUsageRunnerSweepClosesIntervalsAndReturnsBindings(t *testing.T) { ctx := context.Background() s := newTestStore(t) diff --git a/go/internal/store/tenant_tx.go b/go/internal/store/tenant_tx.go index 70e2e2f44..4dc39853e 100644 --- a/go/internal/store/tenant_tx.go +++ b/go/internal/store/tenant_tx.go @@ -18,8 +18,8 @@ import ( // per-transaction compass.tenant_id GUC. // - systemRole is the narrowly-scoped BYPASSRLS role the cross-tenant // background loops (N5/OQ-4: delivery-cursor sweep, deliver-ack advance, -// reattach recovery, lag-resync, compute-usage orphan sweep) run under, and -// ONLY those. It carries no tenant GUC — it is cross-tenant by design. +// reattach recovery, lag-resync, compute-usage orphan sweep) and the Runner +// re-enroll binding reap run under, and ONLY those. It carries no tenant GUC. const ( appRole = "compass_app" systemRole = "compass_system" @@ -40,10 +40,10 @@ type systemRoleKey struct{} // WithSystemRole marks ctx as the cross-tenant background/system path: store // calls made under it run as the BYPASSRLS compass_system role and see every // tenant's rows. It is the OQ-4 (Matt-ruled option 1) exemption, applied ONLY at -// named background-loop entrypoints (the delivery consumer's Run, the hub's -// deliver-ack / forge-notification-ack arms, reattach recovery, and the -// compute-usage orphan sweep) — a request-path call NEVER sets it, so the -// request path stays tenant-scoped and fail-closed under RLS. +// named entrypoints (the delivery consumer's Run, the hub's deliver-ack / +// forge-notification-ack arms, reattach recovery, the compute-usage orphan +// sweep, and Hub.enroll's reap keyed by the authenticated Runner id). Every other +// request-path call stays tenant-scoped and fail-closed under RLS. func WithSystemRole(ctx context.Context) context.Context { return context.WithValue(ctx, systemRoleKey{}, true) } From 673db43b57bf2a172b5adfab144f0e47216fd0bd Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 06:37:57 -0400 Subject: [PATCH 05/10] test(usage): make the lock-wait test count its checks; tighten system-role and estimated wording (RIG-2872) - TestComputeUsageEventsStampedAfterLockWait fails if it finds fewer than the two events it pins. - tenant_tx.go lists the Runner session tenant lookup among system-role entrypoints. - docs: name which events are estimated, and say a reconnect reap records the reap time. Refs RIG-2872 Co-authored-by: Matt Wilkinson --- docs/concepts/tokens-and-billing.md | 7 ++++--- go/internal/store/compute_usage_pgtest_test.go | 5 +++++ go/internal/store/tenant_tx.go | 13 +++++++------ 3 files changed, 16 insertions(+), 9 deletions(-) diff --git a/docs/concepts/tokens-and-billing.md b/docs/concepts/tokens-and-billing.md index 60e7a9187..f64e09472 100644 --- a/docs/concepts/tokens-and-billing.md +++ b/docs/concepts/tokens-and-billing.md @@ -57,9 +57,10 @@ health/quota signal. its start to end. The server emits both events with the binding transition; an hourly sweep marks an inferred end as estimated if a binding disappears without its end event. Intervals already open when the log was introduced get - an estimated start. Every inferred timestamp carries the `estimated` flag; - all others are recorded at the transition. The log is auditable and - reconstructable. + an estimated start. The sweep's ends and backfilled starts carry the + `estimated` flag. Other events are stamped when the server records the + binding change; a Runner-reconnect reap records the reap time. The log is + auditable and reconstructable. - **LLM token usage and spend — recorded for display, not billed day-1.** Every model call's tokens-in/out and cost, captured at the gateway. It powers the in-product usage/spend charts the user sees, and it is *recorded* even though diff --git a/go/internal/store/compute_usage_pgtest_test.go b/go/internal/store/compute_usage_pgtest_test.go index dcf4e7ddf..f04fbc122 100644 --- a/go/internal/store/compute_usage_pgtest_test.go +++ b/go/internal/store/compute_usage_pgtest_test.go @@ -198,13 +198,18 @@ func TestComputeUsageEventsStampedAfterLockWait(t *testing.T) { t.Fatalf("RecordSessionBinding after release: %v", err) } + checked := 0 for _, event := range computeEvents(t, s, tenant, agent.ID) { if event.SessionID == "wait-after" || event.Kind == "end" { + checked++ if event.OccurredAt.Before(released) { t.Fatalf("%s event for %s stamped %s, before the lock release at %s", event.Kind, event.SessionID, event.OccurredAt, released) } } } + if checked != 2 { + t.Fatalf("checked %d events, want the displaced end and the replacement start", checked) + } } func TestComputeUsageConflictRollsBackIntervalEvents(t *testing.T) { diff --git a/go/internal/store/tenant_tx.go b/go/internal/store/tenant_tx.go index 4dc39853e..2cf9cef14 100644 --- a/go/internal/store/tenant_tx.go +++ b/go/internal/store/tenant_tx.go @@ -18,8 +18,9 @@ import ( // per-transaction compass.tenant_id GUC. // - systemRole is the narrowly-scoped BYPASSRLS role the cross-tenant // background loops (N5/OQ-4: delivery-cursor sweep, deliver-ack advance, -// reattach recovery, lag-resync, compute-usage orphan sweep) and the Runner -// re-enroll binding reap run under, and ONLY those. It carries no tenant GUC. +// reattach recovery, lag-resync, compute-usage orphan sweep), the Runner +// session tenant lookup, and the Runner re-enroll binding reap run under, +// and ONLY those. It carries no tenant GUC. const ( appRole = "compass_app" systemRole = "compass_system" @@ -33,8 +34,8 @@ const tenantGUC = "compass.tenant_id" // systemRoleKey marks a context as running the cross-tenant system path. When // present, the store arms statements with SET LOCAL ROLE compass_system -// (BYPASSRLS) and no tenant GUC, instead of the tenant-scoped app role. Only -// named background loops use it; every other path is tenant-scoped and fail-closed. +// (BYPASSRLS) and no tenant GUC, instead of the tenant-scoped app role. Only the +// named entrypoints below use it; every other path is tenant-scoped and fail-closed. type systemRoleKey struct{} // WithSystemRole marks ctx as the cross-tenant background/system path: store @@ -42,8 +43,8 @@ type systemRoleKey struct{} // tenant's rows. It is the OQ-4 (Matt-ruled option 1) exemption, applied ONLY at // named entrypoints (the delivery consumer's Run, the hub's deliver-ack / // forge-notification-ack arms, reattach recovery, the compute-usage orphan -// sweep, and Hub.enroll's reap keyed by the authenticated Runner id). Every other -// request-path call stays tenant-scoped and fail-closed under RLS. +// sweep, Hub.runnerSessionCtx's tenant lookup, and Hub.enroll's reap keyed by the +// authenticated Runner id). Every other request-path call stays tenant-scoped. func WithSystemRole(ctx context.Context) context.Context { return context.WithValue(ctx, systemRoleKey{}, true) } From dbe5ffe19d590592816d5c2eb18f070ee0844aeb Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 13:00:43 -0400 Subject: [PATCH 06/10] refactor(server): start both usage sweepers from one helper (RIG-2872) After the rebase, Serve was one line over the funlen limit. startUsageSweepers starts the retention prune and the orphan-interval close together, so Serve makes one call. Co-authored-by: Matt Wilkinson --- go/server/serve.go | 7 ++----- go/server/sinks.go | 18 +++++++----------- 2 files changed, 9 insertions(+), 16 deletions(-) diff --git a/go/server/serve.go b/go/server/serve.go index bd872a411..db11a4a74 100644 --- a/go/server/serve.go +++ b/go/server/serve.go @@ -907,11 +907,8 @@ func Serve(ctx context.Context, cfg ServeConfig) error { // projection on the comms bus, both on the serve group rooted on gctx // (cancels at shutdown; presence also ends when drainDoors closes the bus). startCommsConsumers(gctx, g, commsBus, fab, st, hub, hubLog) - // A daily prune bounds the raw usage log; the rollups keep their sums. - startUsageRetention(gctx, g, st, cfg.UsageEventRetention, hubLog) - // Close orphaned compute intervals on startup and once per hour; binding - // transitions normally close intervals in their own transaction. - startComputeUsageSweeper(gctx, g, st, hubLog) + // Background usage upkeep: the daily raw-log prune and the hourly orphan-interval close. + startUsageSweepers(gctx, g, st, cfg.UsageEventRetention, hubLog) // Drain member of the same group: wake on gctx cancellation, then hand off to // drainDoors. A drain that overruns (a handler still wedged mid-replay) // surfaces as the error rather than a false clean shutdown; a real serve diff --git a/go/server/sinks.go b/go/server/sinks.go index 3fa823d8f..3748515bb 100644 --- a/go/server/sinks.go +++ b/go/server/sinks.go @@ -184,17 +184,13 @@ func startForgeIngestLanes(gctx context.Context, g *errgroup.Group, board *board } } -// startUsageRetention starts the token-usage retention sweeper on the serve -// group. The raw log grows with activity; the rollups keep the history. -func startUsageRetention(gctx context.Context, g *errgroup.Group, st *store.Store, retention time.Duration, log *slog.Logger) { - w := usage.NewRetentionSweeper(usage.NewPostgres(st), usage.RetentionConfig{Retention: retention, Log: log}) - g.Go(func() error { return w.Run(gctx) }) -} - -// startComputeUsageSweeper closes orphaned compute intervals on the serve group. -func startComputeUsageSweeper(gctx context.Context, g *errgroup.Group, st *store.Store, log *slog.Logger) { - w := usage.NewComputeUsageSweeper(computeUsageCloser{st: st}, log) - g.Go(func() error { return w.Run(gctx) }) +// startUsageSweepers starts the usage retention prune and the orphaned compute +// interval close on the serve group; binding changes close intervals inline. +func startUsageSweepers(gctx context.Context, g *errgroup.Group, st *store.Store, retention time.Duration, log *slog.Logger) { + retain := usage.NewRetentionSweeper(usage.NewPostgres(st), usage.RetentionConfig{Retention: retention, Log: log}) + closer := usage.NewComputeUsageSweeper(computeUsageCloser{st: st}, log) + g.Go(func() error { return retain.Run(gctx) }) + g.Go(func() error { return closer.Run(gctx) }) } type computeUsageCloser struct{ st *store.Store } From 6ce1f45829a34ed8be57729fe7a43f4a47b55257 Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 18:41:20 -0400 Subject: [PATCH 07/10] fix(store): renumber the compute-usage migration to 0005 after the token-tenant migration (RIG-2872) Main added 0004_token_tenant.sql, so 0004_compute_usage.sql now collides and Open refuses the duplicate version. Co-authored-by: Matt Wilkinson --- .../{0004_compute_usage.sql => 0005_compute_usage.sql} | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) rename go/internal/store/migrations/{0004_compute_usage.sql => 0005_compute_usage.sql} (97%) diff --git a/go/internal/store/migrations/0004_compute_usage.sql b/go/internal/store/migrations/0005_compute_usage.sql similarity index 97% rename from go/internal/store/migrations/0004_compute_usage.sql rename to go/internal/store/migrations/0005_compute_usage.sql index c61af3db9..3b7aa8c76 100644 --- a/go/internal/store/migrations/0004_compute_usage.sql +++ b/go/internal/store/migrations/0005_compute_usage.sql @@ -1,4 +1,4 @@ --- 0004_compute_usage: the append-only compute interval event log. Each +-- 0005_compute_usage: the append-only compute interval event log. Each -- session binding is one billable interval; starts and ends commit with its row. CREATE TABLE compute_usage_events ( From e25fc4e919ead3c819af00c9fce0a909072d8bca Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 22:31:08 -0400 Subject: [PATCH 08/10] test(runnerhub): align tenant-B binding tests with reap on every enroll (RIG-2872) Main now reaps a Runner's bindings on every enroll, so the drop test re-records its row after enroll and the reap test needs one enroll. Co-authored-by: Matt Wilkinson --- go/internal/runnerhub/runner_tenant_pgtest_test.go | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/go/internal/runnerhub/runner_tenant_pgtest_test.go b/go/internal/runnerhub/runner_tenant_pgtest_test.go index 49dd88698..e326caa65 100644 --- a/go/internal/runnerhub/runner_tenant_pgtest_test.go +++ b/go/internal/runnerhub/runner_tenant_pgtest_test.go @@ -27,6 +27,10 @@ func TestDropLostSessionScopesToTheSessionTenant(t *testing.T) { sink := &recordingLostSink{} hub.SetSessionLostSink(sink) hub.enroll(ctx, "runner-1", runnerSubject(), compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED) + // Enroll reaps every binding on this Runner; re-record so the drop reads a live row cold. + if _, err := st.RecordSessionBinding(ctxB, "sess-b", agent.ID, "runner-1"); err != nil { + t.Fatalf("RecordSessionBinding after enroll: %v", err) + } hub.dropLostSession(ctx, "runner-1", "sess-b") @@ -38,7 +42,7 @@ func TestDropLostSessionScopesToTheSessionTenant(t *testing.T) { } } -// A re-enrolling Runner spans tenants, so its system-role sweep must reap tenant B's +// An enrolling Runner spans tenants, so its system-role sweep must reap tenant B's // binding and record its end event with tenant B's id. func TestEnrollReapsSessionBindingAcrossTenants(t *testing.T) { ctx := context.Background() @@ -46,8 +50,6 @@ func TestEnrollReapsSessionBindingAcrossTenants(t *testing.T) { hub := newHubOnly() hub.SetSessionBindingStore(st) - // The first enroll keeps durable rows; the second is the reconnect that reaps. - hub.enroll(ctx, "runner-1", runnerSubject(), compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED) hub.enroll(ctx, "runner-1", runnerSubject(), compassv1.RuntimeTier_RUNTIME_TIER_UNSPECIFIED, compassv1.EgressPosture_EGRESS_POSTURE_UNSPECIFIED) var remaining, ends int From 253273a4ce0e7cba34bcda89759ad7d2daf419ac Mon Sep 17 00:00:00 2001 From: mintaka Date: Sat, 3 Oct 2026 22:47:18 -0400 Subject: [PATCH 09/10] docs(runnerhub): say the enroll reap spans tenants (RIG-2872) Co-authored-by: Matt Wilkinson --- go/internal/runnerhub/relay_comms.go | 4 ++-- go/internal/store/tenant_tx.go | 2 +- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/go/internal/runnerhub/relay_comms.go b/go/internal/runnerhub/relay_comms.go index d28ff83dd..1803af628 100644 --- a/go/internal/runnerhub/relay_comms.go +++ b/go/internal/runnerhub/relay_comms.go @@ -251,8 +251,8 @@ func (h *Hub) unbindContainer(containerName string) { // // RIG-3108: the map is a read-through cache. A hit returns immediately. A miss // falls through to the durable binding table ONLY when a Runner is currently -// enrolled AND the ctx is request-scoped. Every enroll reaps that Runner's bootstrap- -// tenant rows, so such a session that predates the enroll stays a fail-closed miss, while a +// enrolled AND the ctx is request-scoped. Every enroll reaps that Runner's rows +// across tenants, so such a session that predates the enroll stays a fail-closed miss, while a // binding another Server instance recorded after it resolves from the row. With no // Runner enrolled the gate skips the table round-trip. store.ErrNotFound (and any // store fault) maps to ok=false, so CodeNotFound behaviour is byte-identical to today. diff --git a/go/internal/store/tenant_tx.go b/go/internal/store/tenant_tx.go index 2cf9cef14..b2b83f8e7 100644 --- a/go/internal/store/tenant_tx.go +++ b/go/internal/store/tenant_tx.go @@ -19,7 +19,7 @@ import ( // - systemRole is the narrowly-scoped BYPASSRLS role the cross-tenant // background loops (N5/OQ-4: delivery-cursor sweep, deliver-ack advance, // reattach recovery, lag-resync, compute-usage orphan sweep), the Runner -// session tenant lookup, and the Runner re-enroll binding reap run under, +// session tenant lookup, and the Runner enroll binding reap run under, // and ONLY those. It carries no tenant GUC. const ( appRole = "compass_app" From b3a2758f7e14af64e9bde9518fa7a37196a768ae Mon Sep 17 00:00:00 2001 From: mintaka Date: Sun, 4 Oct 2026 08:04:40 -0400 Subject: [PATCH 10/10] fix(store): renumber the compute-usage migration to 0006 after the issue-search migration (RIG-2872) Co-authored-by: Matt Wilkinson --- .../{0005_compute_usage.sql => 0006_compute_usage.sql} | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) rename go/internal/store/migrations/{0005_compute_usage.sql => 0006_compute_usage.sql} (97%) diff --git a/go/internal/store/migrations/0005_compute_usage.sql b/go/internal/store/migrations/0006_compute_usage.sql similarity index 97% rename from go/internal/store/migrations/0005_compute_usage.sql rename to go/internal/store/migrations/0006_compute_usage.sql index 3b7aa8c76..957fd8d60 100644 --- a/go/internal/store/migrations/0005_compute_usage.sql +++ b/go/internal/store/migrations/0006_compute_usage.sql @@ -1,4 +1,4 @@ --- 0005_compute_usage: the append-only compute interval event log. Each +-- 0006_compute_usage: the append-only compute interval event log. Each -- session binding is one billable interval; starts and ends commit with its row. CREATE TABLE compute_usage_events (