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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 9 additions & 5 deletions docs/concepts/tokens-and-billing.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,11 +52,15 @@ 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 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
Expand Down
6 changes: 3 additions & 3 deletions go/internal/runnerhub/binding_cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down
29 changes: 20 additions & 9 deletions go/internal/runnerhub/hub.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -1025,15 +1025,20 @@ 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)
// 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 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 {
Expand Down Expand Up @@ -1061,8 +1066,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
}
Expand Down
4 changes: 2 additions & 2 deletions go/internal/runnerhub/relay_comms.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
75 changes: 53 additions & 22 deletions go/internal/runnerhub/runner_tenant_pgtest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,58 @@ 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)
// 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")

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")
}
}

// 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()
st, agent, ctxB := openTenantBSession(t, ctx)

hub := newHubOnly()
hub.SetSessionBindingStore(st)
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 {
Expand All @@ -28,7 +80,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"})
Expand All @@ -42,27 +93,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
Expand Down
2 changes: 1 addition & 1 deletion go/internal/store/backfill_empty_tenant_pgtest_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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))
Expand Down
16 changes: 16 additions & 0 deletions go/internal/store/compute_usage.go
Original file line number Diff line number Diff line change
@@ -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
}
Loading
Loading