mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: add hourly hb_agent_runtime_v1 usage events for Coder Agent runtime (#27312)
closes CODAGT-839 closes CODAGT-843 closes CODAGT-773 ## Summary Adds a new heartbeat usage event type, `hb_agent_runtime_v1`, measuring the total agent-loop runtime of Coder Agents (chats) per UTC hour, plus a reconciler that generates one event per hour with self-healing backfill over a trailing 7-day window. Events flow to Tallyman through the existing publisher unchanged. This measures the new Coder Agents (the `chats` tables), not the deprecated Tasks counted by `dc_managed_agents_v1`. Independent of #27508, which fixes the dead ai-seats cron registration. Both PRs carry the identical `usage_event` create permission hunk for the usage-publisher subject (this feature's generator and the ai-seats cron each need it for heartbeat inserts), so they can land in either order and the overlap merges cleanly. > [!WARNING] > **Do not include this in a release until Tallyman accepts `hb_agent_runtime_v1`.** The publisher marks permanently rejected events as done-forever, and the generator then sees those buckets as complete locally, so their usage would be silently and permanently lost. ## Details Each event's payload is `{"runtime_ms": N}`: the sum of `chat_messages.runtime_ms` for messages created in the hour bucket `[H, H+1)`, across all chats (sub-agents, API-created, archived, and soft-deleted messages included). Events use deterministic IDs (`hb_agent_runtime_v1:<bucket start>`) with `created_at` set to the bucket start, so concurrent replicas race safely via `ON CONFLICT (id) DO NOTHING` without locking, and daily rollups attribute backfilled hours to the correct day. Idle hours produce zero-valued events. A bucket becomes eligible 5 minutes after it closes; hours missing for longer than the 7-day window are forfeited, which can only undercount. Note that this makes `usage_events.created_at` explicitly the *event occurrence time* rather than the row insertion time; the two only diverge for backfilled events. It already behaved as the occurrence timestamp (it drives the daily rollup day and is shipped to Tallyman/Metronome as the event timestamp), and the migration now documents this with a `COMMENT ON COLUMN`, which also surfaces as a Go doc comment on `UsageEvent.CreatedAt`. The new `usage.Generator` runs unconditionally in enterprise builds; the `publish_usage_data` license flag continues to gate egress only, so air-gapped deployments still fill their local ledger. The `aggregate_usage_event()` trigger sums `runtime_ms` per day into `usage_events_daily` (unlike `hb_ai_seats_v1`, which takes the daily max). `InsertHeartbeatUsageEvent` now takes an explicit `createdAt` so generators can backfill historical buckets; the cron passes `clock.Now()` to preserve its existing behavior. ## Tallyman follow-up <details> <summary>Prompt for the Tallyman-repo agent</summary> > **Task**: Add support for the new Coder usage event type `hb_agent_runtime_v1` so Tallyman accepts, validates, and forwards it to Metronome. > > **Background**: coder/coder PR (this PR) adds hourly heartbeat events measuring Coder Agent runtime. Events arrive via the existing `/api/v1/events/ingest` endpoint with: `event_type: "hb_agent_runtime_v1"`, `event_data: {"runtime_ms": <int64 >= 0>}`, deterministic `id` of the form `hb_agent_runtime_v1:2026-07-15_14:00:00` (UTC hour bucket start), and `created_at` set to the bucket start (may be up to ~8 days in the past due to backfill; within Metronome's 34-day dedup window). Zero-value events are normal (idle hours). > > **Work**: > 1. Update Tallyman's vendored/imported `coderd/usage/usagetypes` (or equivalent) to the coder/coder commit that adds `UsageEventTypeHBAgentRuntimeV1` and `HBAgentRuntime`. > 2. Ensure ingestion validation accepts the type (`Valid()` switches) and rejects negative `runtime_ms`. > 3. Ensure Metronome forwarding maps the event with transaction ID derived from the event `id` as for existing types, passing `runtime_ms` through as the property for a SUM-aggregated billable metric ("Coder Agent Hours" = `SUM(runtime_ms) / 3,600,000`). > 4. Do NOT permanently reject unknown-but-well-formed future `hb_*` types if avoidable; at minimum confirm current behavior for unknown types (temporary vs permanent rejection) and report it. > 5. Tests: ingest accept/validate, dedup by ID, Metronome payload mapping. > > **Constraint**: this must be deployed to tallyman-prod **before** any coder/coder release containing the event generator; coderd treats permanent rejections as terminal per event. </details>
This commit is contained in:
@@ -23,7 +23,10 @@ import (
|
||||
var epoch = time.Date(2023, 1, 1, 0, 0, 0, 0, time.UTC)
|
||||
|
||||
const (
|
||||
cronDateFormat = "2006-01-02_15:04:05"
|
||||
// usageEventIDTimeFormat is the timestamp layout used in every
|
||||
// deterministic usage event ID, both the cron's boundary IDs and the
|
||||
// generator's bucket IDs.
|
||||
usageEventIDTimeFormat = "2006-01-02_15:04:05"
|
||||
)
|
||||
|
||||
// HeartbeatFunc generates a heartbeat event and its stable ID.
|
||||
@@ -145,7 +148,7 @@ func (c *Cron) run(ctx context.Context, job CronJob) {
|
||||
// Use the boundary (not wall-clock "now") for the stable ID
|
||||
// so all replicas targeting the same boundary produce the
|
||||
// same key.
|
||||
stableID := string(job.EventType) + ":" + boundary.UTC().Format(cronDateFormat)
|
||||
stableID := string(job.EventType) + ":" + boundary.UTC().Format(usageEventIDTimeFormat)
|
||||
|
||||
// Skip if this bucket was already recorded — avoids running
|
||||
// the potentially expensive heartbeat function for a
|
||||
@@ -184,7 +187,7 @@ func (c *Cron) run(ctx context.Context, job CronJob) {
|
||||
continue
|
||||
}
|
||||
|
||||
if err := c.ins.InsertHeartbeatUsageEvent(ctx, c.db, stableID, event); err != nil {
|
||||
if err := c.ins.InsertHeartbeatUsageEvent(ctx, c.db, stableID, c.clock.Now(), event); err != nil {
|
||||
c.log.Warn(ctx, "cron heartbeat insert failed",
|
||||
slog.F("job", job.Name),
|
||||
slog.Error(err),
|
||||
|
||||
@@ -0,0 +1,242 @@
|
||||
package usage
|
||||
|
||||
import (
|
||||
"context"
|
||||
"math/rand"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"cdr.dev/slog/v3"
|
||||
"github.com/coder/coder/v2/coderd/database"
|
||||
"github.com/coder/coder/v2/coderd/database/dbauthz"
|
||||
"github.com/coder/coder/v2/coderd/pproflabel"
|
||||
agplusage "github.com/coder/coder/v2/coderd/usage"
|
||||
"github.com/coder/coder/v2/coderd/usage/usagetypes"
|
||||
"github.com/coder/quartz"
|
||||
)
|
||||
|
||||
const (
|
||||
// AgentRuntimeInterval is the bucket size of hb_agent_runtime_v1 events.
|
||||
AgentRuntimeInterval = time.Hour
|
||||
// AgentRuntimeWindow is the trailing window scanned for missing buckets.
|
||||
// Buckets still missing beyond this window (e.g. because the deployment
|
||||
// was down for longer) are forfeited, which can only ever undercount
|
||||
// usage.
|
||||
AgentRuntimeWindow = 7 * 24 * time.Hour
|
||||
// AgentRuntimeEligibilityLag is how long after a bucket closes before it
|
||||
// becomes eligible for generation, giving replicas time to commit
|
||||
// in-flight chat messages with timestamps inside the bucket. This
|
||||
// assumes chat message inserts commit within the lag of their
|
||||
// statement-time created_at; a message committing later than the lag
|
||||
// after its bucket closes lands in an already-sealed bucket and its
|
||||
// runtime is dropped (undercount-only).
|
||||
AgentRuntimeEligibilityLag = 5 * time.Minute
|
||||
// agentRuntimeJitter staggers replicas after each hour boundary so one
|
||||
// is likely to complete the work before others attempt it.
|
||||
agentRuntimeJitter = 4 * time.Minute
|
||||
// agentRuntimeStartupDelay is the floor on the first pass after start,
|
||||
// giving the deployment time to finish booting before the generator
|
||||
// competes for database work. Jitter is added on top of it.
|
||||
agentRuntimeStartupDelay = time.Minute
|
||||
// generatorTimerName tags the quartz timer so tests can trap it.
|
||||
generatorTimerName = "agent-runtime-generator"
|
||||
)
|
||||
|
||||
// Generator reconciles hb_agent_runtime_v1 heartbeat usage events. Unlike
|
||||
// Cron jobs, which sample live state when they fire, the Generator derives
|
||||
// events from data already persisted in the database, so it can
|
||||
// deterministically backfill hours missed while the deployment was down,
|
||||
// zero-filling idle hours. Deterministic event IDs plus the database's
|
||||
// ON CONFLICT (id) DO NOTHING make concurrent replicas safe without locking.
|
||||
//
|
||||
// Events are generated unconditionally in enterprise builds; the
|
||||
// publish_usage_data license flag only gates publishing to Tallyman.
|
||||
type Generator struct {
|
||||
clock quartz.Clock
|
||||
log slog.Logger
|
||||
db database.Store
|
||||
ins agplusage.Inserter
|
||||
|
||||
cancel context.CancelFunc
|
||||
wg sync.WaitGroup
|
||||
startOnce sync.Once
|
||||
}
|
||||
|
||||
// NewGenerator creates an unstarted Generator.
|
||||
func NewGenerator(clock quartz.Clock, log slog.Logger, db database.Store, ins agplusage.Inserter) *Generator {
|
||||
return &Generator{
|
||||
clock: clock,
|
||||
log: log,
|
||||
db: db,
|
||||
ins: ins,
|
||||
}
|
||||
}
|
||||
|
||||
// Start launches the reconciliation goroutine. Subsequent calls are no-ops;
|
||||
// a closed Generator cannot be restarted.
|
||||
func (g *Generator) Start(ctx context.Context) {
|
||||
g.startOnce.Do(func() {
|
||||
ctx, g.cancel = context.WithCancel(ctx)
|
||||
g.wg.Add(1)
|
||||
pproflabel.Go(ctx, pproflabel.Service(pproflabel.ServiceUsageEventGenerator), func(ctx context.Context) {
|
||||
g.run(ctx)
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
// Close stops the Generator and waits for its goroutine to exit.
|
||||
// It always returns nil; the error return exists to satisfy io.Closer, as the
|
||||
// Generator is registered with the server's closer list.
|
||||
func (g *Generator) Close() error {
|
||||
if g.cancel != nil {
|
||||
g.cancel()
|
||||
}
|
||||
g.wg.Wait()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (g *Generator) run(ctx context.Context) {
|
||||
//nolint:gocritic // We are a publisher in this function.
|
||||
ctx = dbauthz.AsUsagePublisher(ctx)
|
||||
defer g.wg.Done()
|
||||
|
||||
// The random initial delay staggers replicas that start simultaneously.
|
||||
//nolint:gosec // Jitter does not need cryptographic randomness.
|
||||
delay := agentRuntimeStartupDelay + time.Duration(rand.Int63n(int64(agentRuntimeJitter)))
|
||||
for {
|
||||
timer := g.clock.NewTimer(delay, generatorTimerName)
|
||||
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
if !timer.Stop() {
|
||||
<-timer.C
|
||||
}
|
||||
return
|
||||
case <-timer.C:
|
||||
}
|
||||
|
||||
err := g.generateAgentRuntimeEvents(ctx)
|
||||
if ctx.Err() != nil {
|
||||
return
|
||||
}
|
||||
if err != nil {
|
||||
g.log.Warn(ctx, "generate agent runtime usage events", slog.Error(err))
|
||||
}
|
||||
|
||||
// Wake at the next eligibility instant (hour boundary + lag), not
|
||||
// the next hour boundary. Computing the tick against the
|
||||
// lag-shifted clock keeps a bucket whose eligibility is still
|
||||
// pending in this hour (e.g. a pass that ran just before HH:05)
|
||||
// from waiting a whole extra hour.
|
||||
_, delay = nextTick(g.clock.Now().Add(-AgentRuntimeEligibilityLag), AgentRuntimeInterval, agentRuntimeJitter)
|
||||
}
|
||||
}
|
||||
|
||||
// generateAgentRuntimeEvents inserts one hb_agent_runtime_v1 event per
|
||||
// missing hourly bucket in the trailing window. Per-bucket errors skip only
|
||||
// that bucket; the next tick rescans the whole window, so transient failures
|
||||
// self-heal.
|
||||
func (g *Generator) generateAgentRuntimeEvents(ctx context.Context) error {
|
||||
now := g.clock.Now().UTC()
|
||||
// Bucket [H, H+1) becomes eligible at H + interval + lag.
|
||||
latestEligible := now.Add(-AgentRuntimeInterval - AgentRuntimeEligibilityLag).Truncate(AgentRuntimeInterval)
|
||||
earliest := now.Truncate(AgentRuntimeInterval).Add(-AgentRuntimeWindow)
|
||||
if latestEligible.Before(earliest) {
|
||||
return nil
|
||||
}
|
||||
|
||||
existingTimes, err := g.db.ListUsageEventCreatedAtsByTypeSince(ctx, database.ListUsageEventCreatedAtsByTypeSinceParams{
|
||||
EventType: string(usagetypes.UsageEventTypeHBAgentRuntimeV1),
|
||||
Since: earliest,
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("list existing agent runtime events: %w", err)
|
||||
}
|
||||
// A row marks its bucket complete regardless of publish outcome, so a
|
||||
// bucket whose event Tallyman permanently rejected is never
|
||||
// regenerated (re-inserting under the deterministic ID is a no-op via
|
||||
// ON CONFLICT (id) DO NOTHING).
|
||||
//
|
||||
// The runtime is not lost locally: the row still holds it, and the
|
||||
// event can be re-queued for publishing with
|
||||
//
|
||||
// UPDATE usage_events
|
||||
// SET published_at = NULL, publish_started_at = NULL, failure_message = NULL
|
||||
// WHERE id = 'hb_agent_runtime_v1:<bucket start, e.g. 2026-07-15_14:00:00>';
|
||||
//
|
||||
// That re-arm only has an effect while the bucket is inside the
|
||||
// publisher's 30-day cutoff: SelectUsageEventsForPublishing also
|
||||
// filters created_at > now - INTERVAL '30 days', and created_at is the
|
||||
// bucket start, so past that the UPDATE reports success but the row is
|
||||
// never picked up again. The release gate (Tallyman must accept this
|
||||
// event type before coderd ships it) is what keeps permanent
|
||||
// rejections exceptional.
|
||||
existing := make(map[time.Time]struct{}, len(existingTimes))
|
||||
for _, ts := range existingTimes {
|
||||
// created_at is always the exact bucket start for this event type;
|
||||
// truncation just normalizes timezone and precision.
|
||||
existing[ts.UTC().Truncate(AgentRuntimeInterval)] = struct{}{}
|
||||
}
|
||||
|
||||
var filled, failed int
|
||||
for bucket := earliest; !bucket.After(latestEligible); bucket = bucket.Add(AgentRuntimeInterval) {
|
||||
if _, ok := existing[bucket]; ok {
|
||||
continue
|
||||
}
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
|
||||
err := g.generateBucket(ctx, bucket)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
// A cancel landing mid-query surfaces as a bucket error.
|
||||
// Report the cancellation instead of logging shutdown as
|
||||
// a bucket failure.
|
||||
return ctx.Err()
|
||||
}
|
||||
// Skip only the failed bucket so one bad bucket (e.g. an
|
||||
// invalid runtime sum) cannot stall every later bucket until
|
||||
// it ages out of the window.
|
||||
g.log.Warn(ctx, "generate agent runtime usage event for bucket",
|
||||
slog.F("bucket", bucket),
|
||||
slog.Error(err),
|
||||
)
|
||||
failed++
|
||||
continue
|
||||
}
|
||||
filled++
|
||||
}
|
||||
if filled > 0 || failed > 0 {
|
||||
g.log.Info(ctx, "generated agent runtime usage events",
|
||||
slog.F("buckets_filled", filled),
|
||||
slog.F("buckets_failed", failed),
|
||||
slog.F("window_start", earliest),
|
||||
slog.F("latest_eligible", latestEligible),
|
||||
)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// generateBucket computes and inserts the event for a single hourly bucket.
|
||||
func (g *Generator) generateBucket(ctx context.Context, bucket time.Time) error {
|
||||
runtimeMs, err := g.db.GetTotalChatMessageRuntimeMsInRange(ctx, database.GetTotalChatMessageRuntimeMsInRangeParams{
|
||||
StartTime: bucket,
|
||||
EndTime: bucket.Add(AgentRuntimeInterval),
|
||||
})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("sum chat message runtime: %w", err)
|
||||
}
|
||||
|
||||
// The deterministic ID makes concurrent inserts of the same bucket
|
||||
// idempotent, and created_at is the bucket start (not the insertion
|
||||
// time) so daily rollups attribute backfilled hours to the correct day.
|
||||
stableID := string(usagetypes.UsageEventTypeHBAgentRuntimeV1) + ":" + bucket.Format(usageEventIDTimeFormat)
|
||||
err = g.ins.InsertHeartbeatUsageEvent(ctx, g.db, stableID, bucket, usagetypes.HBAgentRuntime{RuntimeMs: runtimeMs})
|
||||
if err != nil {
|
||||
return xerrors.Errorf("insert usage event: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,481 @@
|
||||
package usage_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"cdr.dev/slog/v3"
|
||||
"cdr.dev/slog/v3/sloggers/slogtest"
|
||||
"github.com/coder/coder/v2/coderd/coderdtest"
|
||||
"github.com/coder/coder/v2/coderd/database"
|
||||
"github.com/coder/coder/v2/coderd/database/dbauthz"
|
||||
"github.com/coder/coder/v2/coderd/database/dbgen"
|
||||
"github.com/coder/coder/v2/coderd/database/dbtestutil"
|
||||
"github.com/coder/coder/v2/coderd/rbac"
|
||||
"github.com/coder/coder/v2/coderd/usage/usagetypes"
|
||||
"github.com/coder/coder/v2/enterprise/coderd/usage"
|
||||
"github.com/coder/coder/v2/testutil"
|
||||
"github.com/coder/quartz"
|
||||
)
|
||||
|
||||
// generatorTimerName must match the tag the Generator passes to
|
||||
// clock.NewTimer so tests can trap its timers.
|
||||
const generatorTimerName = "agent-runtime-generator"
|
||||
|
||||
// warnSink counts log entries at Warn or above. The generator downgrades
|
||||
// per-bucket failures to Warn logs (which slogtest tolerates), so tests that
|
||||
// must prove the error paths stayed quiet assert on this counter instead.
|
||||
type warnSink struct{ count atomic.Int64 }
|
||||
|
||||
func (s *warnSink) LogEntry(_ context.Context, e slog.SinkEntry) {
|
||||
if e.Level >= slog.LevelWarn {
|
||||
s.count.Add(1)
|
||||
}
|
||||
}
|
||||
|
||||
func (*warnSink) Sync() {}
|
||||
|
||||
// generatorHarness runs the generator against a dbauthz-wrapped store so the
|
||||
// tests also verify that the usage publisher subject holds the permissions
|
||||
// the generator's queries require.
|
||||
type generatorHarness struct {
|
||||
db database.Store
|
||||
authzDB database.Store
|
||||
rawDB *sql.DB
|
||||
|
||||
user database.User
|
||||
modelConfig database.ChatModelConfig
|
||||
chat database.Chat
|
||||
chat2 database.Chat
|
||||
}
|
||||
|
||||
func newGeneratorHarness(t *testing.T) *generatorHarness {
|
||||
t.Helper()
|
||||
db, _, rawDB := dbtestutil.NewDBWithSQLDB(t)
|
||||
log := slogtest.Make(t, nil)
|
||||
authzDB := dbauthz.New(db, rbac.NewStrictAuthorizer(prometheus.NewRegistry()), log, coderdtest.AccessControlStorePointer())
|
||||
|
||||
user := dbgen.User(t, db, database.User{})
|
||||
org := dbgen.Organization(t, db, database.Organization{})
|
||||
_ = dbgen.OrganizationMember(t, db, database.OrganizationMember{UserID: user.ID, OrganizationID: org.ID})
|
||||
_ = dbgen.ChatProvider(t, db, database.ChatProvider{
|
||||
Provider: "openai",
|
||||
DisplayName: "OpenAI",
|
||||
})
|
||||
mc := dbgen.ChatModelConfig(t, db, database.ChatModelConfig{
|
||||
Model: "test-model",
|
||||
ContextLimit: 8192,
|
||||
})
|
||||
chat := dbgen.Chat(t, db, database.Chat{
|
||||
OrganizationID: org.ID,
|
||||
OwnerID: user.ID,
|
||||
LastModelConfigID: mc.ID,
|
||||
})
|
||||
chat2 := dbgen.Chat(t, db, database.Chat{
|
||||
OrganizationID: org.ID,
|
||||
OwnerID: user.ID,
|
||||
LastModelConfigID: mc.ID,
|
||||
})
|
||||
return &generatorHarness{
|
||||
db: db,
|
||||
authzDB: authzDB,
|
||||
rawDB: rawDB,
|
||||
user: user,
|
||||
modelConfig: mc,
|
||||
chat: chat,
|
||||
chat2: chat2,
|
||||
}
|
||||
}
|
||||
|
||||
func (h *generatorHarness) insertRuntimeMessage(ctx context.Context, t *testing.T, chatID uuid.UUID, runtimeMs int64, createdAt time.Time, deleted bool) {
|
||||
t.Helper()
|
||||
msg := dbgen.ChatMessage(t, h.db, database.ChatMessage{
|
||||
ChatID: chatID,
|
||||
CreatedBy: uuid.NullUUID{UUID: h.user.ID, Valid: true},
|
||||
ModelConfigID: uuid.NullUUID{UUID: h.modelConfig.ID, Valid: true},
|
||||
Role: database.ChatMessageRoleAssistant,
|
||||
RuntimeMs: sql.NullInt64{Int64: runtimeMs, Valid: true},
|
||||
})
|
||||
_, err := h.rawDB.ExecContext(ctx, "UPDATE chat_messages SET created_at = $1, deleted = $2 WHERE id = $3", createdAt, deleted, msg.ID)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
// fetchRuntimeEvents fails the test if any bucket has more than one event.
|
||||
func (h *generatorHarness) fetchRuntimeEvents(ctx context.Context, t *testing.T) (map[time.Time]int64, map[time.Time]string) {
|
||||
t.Helper()
|
||||
rows, err := h.rawDB.QueryContext(ctx, `
|
||||
SELECT id, (event_data->>'runtime_ms')::bigint, created_at
|
||||
FROM usage_events
|
||||
WHERE event_type = 'hb_agent_runtime_v1'
|
||||
`)
|
||||
require.NoError(t, err)
|
||||
defer rows.Close()
|
||||
|
||||
runtimes := make(map[time.Time]int64)
|
||||
ids := make(map[time.Time]string)
|
||||
for rows.Next() {
|
||||
var (
|
||||
id string
|
||||
runtimeMs int64
|
||||
createdAt time.Time
|
||||
)
|
||||
require.NoError(t, rows.Scan(&id, &runtimeMs, &createdAt))
|
||||
bucket := createdAt.UTC()
|
||||
_, ok := runtimes[bucket]
|
||||
require.False(t, ok, "duplicate event for bucket %s", bucket)
|
||||
runtimes[bucket] = runtimeMs
|
||||
ids[bucket] = id
|
||||
}
|
||||
require.NoError(t, rows.Err())
|
||||
return runtimes, ids
|
||||
}
|
||||
|
||||
func expectedBuckets(first, last time.Time, overrides map[time.Time]int64) map[time.Time]int64 {
|
||||
expected := make(map[time.Time]int64)
|
||||
for bucket := first; !bucket.After(last); bucket = bucket.Add(time.Hour) {
|
||||
expected[bucket] = 0
|
||||
}
|
||||
for bucket, runtimeMs := range overrides {
|
||||
expected[bucket] = runtimeMs
|
||||
}
|
||||
return expected
|
||||
}
|
||||
|
||||
func TestGenerator(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// startTime is exactly on an hour boundary, so the first tick (which
|
||||
// fires 1-5 minutes later) always lands before the just-closed bucket
|
||||
// [13:00, 14:00) becomes eligible at 14:05.
|
||||
startTime := time.Date(2025, 3, 10, 14, 0, 0, 0, time.UTC)
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
log := slogtest.Make(t, nil)
|
||||
h := newGeneratorHarness(t)
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(startTime)
|
||||
|
||||
var (
|
||||
bucketA = time.Date(2025, 3, 10, 10, 0, 0, 0, time.UTC)
|
||||
bucketB = time.Date(2025, 3, 10, 11, 0, 0, 0, time.UTC)
|
||||
// The most recent closed bucket; not eligible at the first tick.
|
||||
bucketC = time.Date(2025, 3, 10, 13, 0, 0, 0, time.UTC)
|
||||
// Window bounds at the first tick.
|
||||
windowFirst = startTime.Add(-usage.AgentRuntimeWindow) // 2025-03-03 14:00
|
||||
windowLast = startTime.Add(-2 * time.Hour) // 2025-03-10 12:00
|
||||
)
|
||||
|
||||
// Bucket A: a message on the bucket start boundary, a message from a
|
||||
// second chat, and a soft-deleted message. All must be counted.
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 1000, bucketA, false)
|
||||
h.insertRuntimeMessage(ctx, t, h.chat2.ID, 2000, bucketA.Add(15*time.Minute), false)
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 4000, bucketA.Add(30*time.Minute), true)
|
||||
// Bucket B: a message exactly on the A/B boundary belongs to B.
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 8000, bucketB, false)
|
||||
// Older than the window: must never be generated.
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 16000, windowFirst.Add(-30*time.Minute), false)
|
||||
// Bucket C: only becomes eligible at 14:05, after the first tick.
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 32000, bucketC.Add(30*time.Minute), false)
|
||||
|
||||
trap := clock.Trap().NewTimer(generatorTimerName)
|
||||
defer trap.Close()
|
||||
|
||||
gen := usage.NewGenerator(clock, log, h.authzDB, usage.NewDBInserter())
|
||||
gen.Start(ctx)
|
||||
defer gen.Close()
|
||||
|
||||
call := trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
require.GreaterOrEqual(t, call.Duration, time.Minute)
|
||||
require.Less(t, call.Duration, 5*time.Minute)
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
|
||||
// The generator creates the next timer only after the tick completes,
|
||||
// so trapping it synchronizes with the end of the pass.
|
||||
call = trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
|
||||
// The first pass fills every bucket in [windowFirst, windowLast]: bucket
|
||||
// C is not yet eligible, the pre-window message is excluded, and idle
|
||||
// hours are zero-filled.
|
||||
runtimes, ids := h.fetchRuntimeEvents(ctx, t)
|
||||
require.Equal(t, expectedBuckets(windowFirst, windowLast, map[time.Time]int64{
|
||||
bucketA: 7000,
|
||||
bucketB: 8000,
|
||||
}), runtimes)
|
||||
require.Equal(t, "hb_agent_runtime_v1:2025-03-10_10:00:00", ids[bucketA])
|
||||
|
||||
// The next tick fires at bucket C's eligibility instant (14:05 plus
|
||||
// jitter), in the same hour the first pass ran, rather than waiting for
|
||||
// the next hour boundary.
|
||||
fireTime := clock.Now().Add(call.Duration)
|
||||
require.Equal(t, startTime, fireTime.Truncate(usage.AgentRuntimeInterval))
|
||||
require.GreaterOrEqual(t, fireTime.Sub(fireTime.Truncate(usage.AgentRuntimeInterval)), usage.AgentRuntimeEligibilityLag)
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
call = trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
|
||||
// The second tick fills only the newly-eligible bucket C (13:00). The
|
||||
// window's start advanced by an hour, but the bucket at the old
|
||||
// windowFirst is kept (rows are never deleted).
|
||||
runtimes, _ = h.fetchRuntimeEvents(ctx, t)
|
||||
require.Equal(t, expectedBuckets(windowFirst, windowLast.Add(time.Hour), map[time.Time]int64{
|
||||
bucketA: 7000,
|
||||
bucketB: 8000,
|
||||
bucketC: 32000,
|
||||
}), runtimes)
|
||||
|
||||
// The third tick fires after the next hour boundary (15:05 plus jitter)
|
||||
// and fills the idle 14:00 bucket.
|
||||
fireTime = clock.Now().Add(call.Duration)
|
||||
require.Equal(t, startTime.Add(time.Hour), fireTime.Truncate(usage.AgentRuntimeInterval))
|
||||
require.GreaterOrEqual(t, fireTime.Sub(fireTime.Truncate(usage.AgentRuntimeInterval)), usage.AgentRuntimeEligibilityLag)
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
call = trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
|
||||
runtimes, _ = h.fetchRuntimeEvents(ctx, t)
|
||||
require.Equal(t, expectedBuckets(windowFirst, windowLast.Add(2*time.Hour), map[time.Time]int64{
|
||||
bucketA: 7000,
|
||||
bucketB: 8000,
|
||||
bucketC: 32000,
|
||||
}), runtimes)
|
||||
|
||||
// A separate generator started later (e.g. another replica restarting)
|
||||
// finds nothing to do: all buckets in its window already exist.
|
||||
gen2 := usage.NewGenerator(clock, log, h.authzDB, usage.NewDBInserter())
|
||||
gen2.Start(ctx)
|
||||
defer gen2.Close()
|
||||
call = trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
call = trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
|
||||
runtimes2, _ := h.fetchRuntimeEvents(ctx, t)
|
||||
require.Equal(t, runtimes, runtimes2)
|
||||
}
|
||||
|
||||
// TestGeneratorBackfillAfterDowntime simulates a deployment that was down
|
||||
// for several hours: a fresh generator's first pass backfills exactly the
|
||||
// missing buckets.
|
||||
func TestGeneratorBackfillAfterDowntime(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
startTime := time.Date(2025, 3, 10, 14, 0, 0, 0, time.UTC)
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
log := slogtest.Make(t, nil)
|
||||
h := newGeneratorHarness(t)
|
||||
|
||||
windowFirst := startTime.Add(-usage.AgentRuntimeWindow)
|
||||
gapBucket := time.Date(2025, 3, 10, 11, 0, 0, 0, time.UTC)
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 5000, gapBucket.Add(45*time.Minute), false)
|
||||
|
||||
// Simulate events generated before the deployment went down at 09:00:
|
||||
// every bucket in [windowFirst, 08:00] exists with runtime 1.
|
||||
inserter := usage.NewDBInserter()
|
||||
for bucket := windowFirst; !bucket.After(time.Date(2025, 3, 10, 8, 0, 0, 0, time.UTC)); bucket = bucket.Add(time.Hour) {
|
||||
err := inserter.InsertHeartbeatUsageEvent(ctx, h.db, "hb_agent_runtime_v1:"+bucket.Format("2006-01-02_15:04:05"), bucket, usagetypes.HBAgentRuntime{RuntimeMs: 1})
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(startTime)
|
||||
trap := clock.Trap().NewTimer(generatorTimerName)
|
||||
defer trap.Close()
|
||||
|
||||
gen := usage.NewGenerator(clock, log, h.authzDB, usage.NewDBInserter())
|
||||
gen.Start(ctx)
|
||||
defer gen.Close()
|
||||
|
||||
call := trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
call = trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
|
||||
// The pass must fill exactly the gap [09:00, 12:00] and leave the
|
||||
// pre-existing rows untouched.
|
||||
runtimes, _ := h.fetchRuntimeEvents(ctx, t)
|
||||
expected := expectedBuckets(windowFirst, startTime.Add(-2*time.Hour), map[time.Time]int64{
|
||||
gapBucket: 5000,
|
||||
})
|
||||
for bucket := windowFirst; !bucket.After(time.Date(2025, 3, 10, 8, 0, 0, 0, time.UTC)); bucket = bucket.Add(time.Hour) {
|
||||
expected[bucket] = 1
|
||||
}
|
||||
require.Equal(t, expected, runtimes)
|
||||
}
|
||||
|
||||
func TestGeneratorInserterArguments(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
startTime := time.Date(2025, 3, 10, 14, 0, 0, 0, time.UTC)
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
log := slogtest.Make(t, nil)
|
||||
h := newGeneratorHarness(t)
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(startTime)
|
||||
|
||||
var (
|
||||
bucketA = time.Date(2025, 3, 10, 10, 0, 0, 0, time.UTC)
|
||||
windowFirst = startTime.Add(-usage.AgentRuntimeWindow)
|
||||
// The first pass stops at 12:00; the second pass fires after 14:05
|
||||
// and adds the newly eligible 13:00 bucket.
|
||||
windowLast = startTime.Add(-time.Hour)
|
||||
)
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 1000, bucketA.Add(10*time.Minute), false)
|
||||
|
||||
ins := coderdtest.NewUsageInserter()
|
||||
trap := clock.Trap().NewTimer(generatorTimerName)
|
||||
defer trap.Close()
|
||||
|
||||
gen := usage.NewGenerator(clock, log, h.authzDB, ins)
|
||||
gen.Start(ctx)
|
||||
defer gen.Close()
|
||||
|
||||
// Each pass requests its next timer only after finishing, so trapping
|
||||
// that request synchronizes with the end of the pass.
|
||||
for range 2 {
|
||||
call := trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
}
|
||||
call := trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
|
||||
var expected []coderdtest.HeartbeatEvent
|
||||
for bucket := windowFirst; !bucket.After(windowLast); bucket = bucket.Add(usage.AgentRuntimeInterval) {
|
||||
var runtimeMs int64
|
||||
if bucket.Equal(bucketA) {
|
||||
runtimeMs = 1000
|
||||
}
|
||||
expected = append(expected, coderdtest.HeartbeatEvent{
|
||||
ID: "hb_agent_runtime_v1:" + bucket.Format("2006-01-02_15:04:05"),
|
||||
CreatedAt: bucket,
|
||||
Event: usagetypes.HBAgentRuntime{RuntimeMs: runtimeMs},
|
||||
})
|
||||
}
|
||||
require.Equal(t, expected, ins.GetHeartbeatEvents())
|
||||
}
|
||||
|
||||
// TestGeneratorConcurrentReplicas runs two generators against the same
|
||||
// database concurrently and verifies exactly one event is produced per
|
||||
// bucket. Each replica gets its own mock clock (as real replicas have their
|
||||
// own wall clocks) so their first passes can be fired independently and run
|
||||
// at the same time.
|
||||
func TestGeneratorConcurrentReplicas(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
startTime := time.Date(2025, 3, 10, 14, 0, 0, 0, time.UTC)
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
// If insert idempotency regressed (e.g. ON CONFLICT DO NOTHING was
|
||||
// dropped), the losing replica's inserts would error and surface as
|
||||
// Warn logs, which slogtest alone would not catch.
|
||||
sink := &warnSink{}
|
||||
log := slogtest.Make(t, nil).AppendSinks(sink)
|
||||
h := newGeneratorHarness(t)
|
||||
|
||||
bucketA := time.Date(2025, 3, 10, 10, 0, 0, 0, time.UTC)
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 1000, bucketA.Add(10*time.Minute), false)
|
||||
|
||||
// Both replicas' first passes fire between 14:01 and 14:05, so they
|
||||
// compute identical windows regardless of their random startup jitter.
|
||||
var traps []*quartz.Trap
|
||||
for range 2 {
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(startTime)
|
||||
trap := clock.Trap().NewTimer(generatorTimerName)
|
||||
t.Cleanup(trap.Close)
|
||||
traps = append(traps, trap)
|
||||
|
||||
gen := usage.NewGenerator(clock, log, h.authzDB, usage.NewDBInserter())
|
||||
gen.Start(ctx)
|
||||
// Cleanups run LIFO, so each generator is closed before its trap.
|
||||
t.Cleanup(func() { _ = gen.Close() })
|
||||
|
||||
call := trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
// Fire the initial timer without waiting for the pass to complete so
|
||||
// both replicas' passes overlap.
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
}
|
||||
|
||||
// Wait for both passes to complete (each requests its next timer only
|
||||
// after the pass finishes).
|
||||
for _, trap := range traps {
|
||||
call := trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
}
|
||||
|
||||
// fetchRuntimeEvents fails on duplicate buckets; the expected map proves
|
||||
// both replicas raced without double-inserting, and the warn counter
|
||||
// proves the duplicate inserts were deduplicated rather than rejected
|
||||
// with errors.
|
||||
runtimes, _ := h.fetchRuntimeEvents(ctx, t)
|
||||
require.Equal(t, expectedBuckets(
|
||||
startTime.Add(-usage.AgentRuntimeWindow),
|
||||
startTime.Add(-2*time.Hour),
|
||||
map[time.Time]int64{bucketA: 1000},
|
||||
), runtimes)
|
||||
require.Zero(t, sink.count.Load(), "no replica may log Warn or above during the race")
|
||||
}
|
||||
|
||||
// TestGeneratorPoisonBucket verifies that a bucket which fails
|
||||
// deterministically on every tick does not stall generation of later
|
||||
// buckets.
|
||||
func TestGeneratorPoisonBucket(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
startTime := time.Date(2025, 3, 10, 14, 0, 0, 0, time.UTC)
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
sink := &warnSink{}
|
||||
log := slogtest.Make(t, nil).AppendSinks(sink)
|
||||
h := newGeneratorHarness(t)
|
||||
clock := quartz.NewMock(t)
|
||||
clock.Set(startTime)
|
||||
|
||||
var (
|
||||
poisonBucket = time.Date(2025, 3, 10, 10, 0, 0, 0, time.UTC)
|
||||
laterBucket = time.Date(2025, 3, 10, 11, 0, 0, 0, time.UTC)
|
||||
)
|
||||
// Nothing prevents a negative runtime_ms row in the database, and a
|
||||
// negative sum fails event validation on every tick.
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, -5, poisonBucket.Add(10*time.Minute), false)
|
||||
h.insertRuntimeMessage(ctx, t, h.chat.ID, 1000, laterBucket.Add(10*time.Minute), false)
|
||||
|
||||
trap := clock.Trap().NewTimer(generatorTimerName)
|
||||
defer trap.Close()
|
||||
|
||||
gen := usage.NewGenerator(clock, log, h.authzDB, usage.NewDBInserter())
|
||||
gen.Start(ctx)
|
||||
defer gen.Close()
|
||||
|
||||
call := trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
clock.Advance(call.Duration).MustWait(ctx)
|
||||
call = trap.MustWait(ctx)
|
||||
call.MustRelease(ctx)
|
||||
|
||||
// Every bucket except the poison bucket is generated, including buckets
|
||||
// after it, and the failure is logged.
|
||||
runtimes, _ := h.fetchRuntimeEvents(ctx, t)
|
||||
expected := expectedBuckets(
|
||||
startTime.Add(-usage.AgentRuntimeWindow),
|
||||
startTime.Add(-2*time.Hour),
|
||||
map[time.Time]int64{laterBucket: 1000},
|
||||
)
|
||||
delete(expected, poisonBucket)
|
||||
require.Equal(t, expected, runtimes)
|
||||
require.NotZero(t, sink.count.Load(), "the poison bucket failure must be logged")
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package usage
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"golang.org/x/xerrors"
|
||||
@@ -68,10 +69,16 @@ func (i *dbInserter) InsertDiscreteUsageEvent(ctx context.Context, tx database.S
|
||||
}
|
||||
|
||||
// InsertHeartbeatUsageEvent implements agplusage.Inserter.
|
||||
func (i *dbInserter) InsertHeartbeatUsageEvent(ctx context.Context, tx database.Store, id string, event usagetypes.HeartbeatEvent) error {
|
||||
func (*dbInserter) InsertHeartbeatUsageEvent(ctx context.Context, tx database.Store, id string, createdAt time.Time, event usagetypes.HeartbeatEvent) error {
|
||||
if !event.EventType().IsHeartbeat() {
|
||||
return xerrors.Errorf("event type %q is not a heartbeat event", event.EventType())
|
||||
}
|
||||
// A zero createdAt stores the row at year 1, where bucket reconciliation
|
||||
// can never match it again while the deterministic id turns every retry
|
||||
// into a no-op, silently forfeiting the bucket's usage.
|
||||
if createdAt.IsZero() {
|
||||
return xerrors.Errorf("createdAt must be set for %q event", event.EventType())
|
||||
}
|
||||
if err := event.Valid(); err != nil {
|
||||
return xerrors.Errorf("invalid %q event: %w", event.EventType(), err)
|
||||
}
|
||||
@@ -87,6 +94,6 @@ func (i *dbInserter) InsertHeartbeatUsageEvent(ctx context.Context, tx database.
|
||||
ID: id,
|
||||
EventType: string(event.EventType()),
|
||||
EventData: jsonData,
|
||||
CreatedAt: dbtime.Time(i.clock.Now()),
|
||||
CreatedAt: dbtime.Time(createdAt),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -68,6 +68,35 @@ func TestInserter(t *testing.T) {
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Heartbeat", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
ctrl := gomock.NewController(t)
|
||||
db := dbmock.NewMockStore(ctrl)
|
||||
inserter := usage.NewDBInserter()
|
||||
|
||||
// Heartbeat inserts must store the provided id and createdAt
|
||||
// verbatim.
|
||||
event := usagetypes.HBAgentRuntime{RuntimeMs: 1234}
|
||||
eventJSON := jsoninate(t, event)
|
||||
id := "hb_agent_runtime_v1:2025-01-02_03:00:00"
|
||||
createdAt := time.Date(2025, 1, 2, 3, 0, 0, 0, time.UTC)
|
||||
|
||||
db.EXPECT().InsertUsageEvent(gomock.Any(), gomock.Any()).DoAndReturn(
|
||||
func(ctx interface{}, params database.InsertUsageEventParams) error {
|
||||
assert.Equal(t, id, params.ID)
|
||||
assert.Equal(t, event.EventType(), usagetypes.UsageEventType(params.EventType))
|
||||
assert.JSONEq(t, eventJSON, string(params.EventData))
|
||||
assert.Equal(t, dbtime.Time(createdAt), params.CreatedAt)
|
||||
return nil
|
||||
},
|
||||
).Times(1)
|
||||
|
||||
err := inserter.InsertHeartbeatUsageEvent(ctx, db, id, createdAt, event)
|
||||
require.NoError(t, err)
|
||||
})
|
||||
|
||||
t.Run("InvalidEvent", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
@@ -81,5 +110,26 @@ func TestInserter(t *testing.T) {
|
||||
Count: 0, // invalid
|
||||
})
|
||||
assert.ErrorContains(t, err, `invalid "dc_managed_agents_v1" event: count must be greater than 0`)
|
||||
|
||||
err = inserter.InsertHeartbeatUsageEvent(ctx, db, "some-id", time.Now(), usagetypes.HBAgentRuntime{
|
||||
RuntimeMs: -1, // invalid
|
||||
})
|
||||
assert.ErrorContains(t, err, `invalid "hb_agent_runtime_v1" event: runtime_ms cannot be negative`)
|
||||
})
|
||||
|
||||
t.Run("ZeroCreatedAt", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := testutil.Context(t, testutil.WaitLong)
|
||||
ctrl := gomock.NewController(t)
|
||||
// The mock store fails the test on any unexpected call, so no insert
|
||||
// may reach the database.
|
||||
db := dbmock.NewMockStore(ctrl)
|
||||
|
||||
inserter := usage.NewDBInserter()
|
||||
err := inserter.InsertHeartbeatUsageEvent(ctx, db, "some-id", time.Time{}, usagetypes.HBAgentRuntime{
|
||||
RuntimeMs: 1,
|
||||
})
|
||||
assert.ErrorContains(t, err, `createdAt must be set for "hb_agent_runtime_v1" event`)
|
||||
})
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user