Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
Show all changes
36 commits
Select commit Hold shift + click to select a range
5d336df
feat: limit concurrent chat agents with pooled admission at acquisition
ibetitsmike Aug 5, 2026
edb9c68
test(coderd): cover FIFO admission, interrupt claim, and concurrent a…
ibetitsmike Aug 5, 2026
ae0122d
refactor(coderd/database): prioritize interrupts and harden candidate…
ibetitsmike Aug 5, 2026
6de8c97
chore(enterprise/coderd/x/chatd): unexport pool caps and use WaitGrou…
ibetitsmike Aug 5, 2026
dee72c0
chore: tighten comments on chatd admission branch
ibetitsmike Aug 5, 2026
0fc32c1
fix(coderd): partition acquisition by capacity pool and add dbauthz t…
ibetitsmike Aug 5, 2026
c822311
fix(coderd/x/chatd): floor acquisition batch size at two
ibetitsmike Aug 5, 2026
161c4bf
fix(coderd/x/chatd): clear persisted queue markers without an admissi…
ibetitsmike Aug 5, 2026
89603b3
docs(coderd/x/chatd): state that runtime-hours usage is not yet popul…
ibetitsmike Aug 5, 2026
1e01395
fix(site): reject stale capacity events and populate db2sdk chat fixture
ibetitsmike Aug 5, 2026
da90da4
fix: count queue entries on marker wins and require newer status events
ibetitsmike Aug 5, 2026
819439a
docs(coderd/x/chatd): correct acquisition ticker default to 1s
ibetitsmike Aug 6, 2026
64d0d69
refactor: derive chat capacity queue state instead of persisting a ma…
ibetitsmike Aug 6, 2026
e52a28d
chore: tighten comments and document capacity gauges
ibetitsmike Aug 6, 2026
b391317
test(coderd/x/chatd): deflake capacity queue event and FIFO tests
ibetitsmike Aug 6, 2026
5d3c88d
fix(coderd/x/chatd): clear queued banner when a different replica admits
ibetitsmike Aug 6, 2026
a36eddb
refactor(coderd): merge chat capacity admission and limits into one seam
ibetitsmike Aug 6, 2026
f54aed7
fix(coderd): clear capacity banner on interrupt and stale capacity ev…
ibetitsmike Aug 6, 2026
07e6306
fix(coderd/x/chatd): publish capacity clear regardless of current cap
ibetitsmike Aug 6, 2026
49beb28
fix(coderd/x/chatd): revalidate capacity queue entry from a fresh sna…
ibetitsmike Aug 6, 2026
e4d3e42
fix(coderd/x/chatd): reconcile unseen capacity queue entries on all-s…
ibetitsmike Aug 6, 2026
ab2a43e
fix(coderd): queue chats arriving behind a full-pool backlog
ibetitsmike Aug 6, 2026
d6c5e60
fix(coderd): close capacity admission bypass and stuck queued banner …
ibetitsmike Aug 6, 2026
6c99ae7
fix(site/src/api/queries): reconcile rejected capacity events via ent…
ibetitsmike Aug 6, 2026
fcefcad
chore: clean up chat capacity comments
ibetitsmike Aug 10, 2026
95806a3
fix: base chat capacity on live ownership
ibetitsmike Aug 10, 2026
e37d268
fix(enterprise/coderd/x/chatd): enforce runtime hard limit
ibetitsmike Aug 11, 2026
8676e86
refactor: simplify chat capacity admission
ibetitsmike Aug 11, 2026
fbb4d32
fix: enforce chat capacity across deployments
ibetitsmike Aug 12, 2026
8bdcf1a
chore: clean up concurrency comments
ibetitsmike Aug 12, 2026
cdadc19
fix(coderd/x/chatd): address review feedback on capacity comments and…
ibetitsmike Aug 12, 2026
84f6666
docs: simplify chatd limiter architecture
ibetitsmike Aug 12, 2026
6d1e03e
test(enterprise/coderd/x/chatd): pin capped unlock for disabled zero-…
jaaydenh Aug 13, 2026
d4e2eee
fix: align pooled chat admission with current main
ibetitsmike Aug 13, 2026
adb7efb
fix(coderd/database): renumber migration to avoid collision with main
ibetitsmike Aug 17, 2026
f56a4e2
fix(site/src/pages/AgentsPage): keep running-chat poll under the bind…
ibetitsmike Aug 17, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
fix: count queue entries on marker wins and require newer status events
  • Loading branch information
ibetitsmike committed Aug 18, 2026
commit da90da4cc55983152d1df272ab72fe1bb0c67a55
10 changes: 3 additions & 7 deletions coderd/database/querier.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 3 additions & 7 deletions coderd/database/queries.sql.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 3 additions & 7 deletions coderd/database/queries/chats.sql
Original file line number Diff line number Diff line change
Expand Up @@ -2327,13 +2327,9 @@ WHERE chat_id = @chat_id::uuid
-- Missing ownership is worker_id IS NULL. Inconsistent ownership is
-- runner_id IS NULL while worker_id is set. Stale ownership is no
-- heartbeat row for (chat_id, runner_id), or one older than
-- @stale_seconds by database time. Ordering: interrupting first (stop
-- requests skip the queue), then requires_action (bypasses capacity
-- admission), then running chats interleaved across the root/subagent
-- pools, FIFO inside each pool. Interleaving keeps each pool's oldest
-- candidate near the front so a deep backlog in one full pool cannot
-- starve the other, letting the worker end a pass once a batch makes
-- no progress. @exclude_ids advances past refusals.
-- Interrupting and requires_action bypass admission. Running chats are FIFO
-- in interleaved root/subagent pools so one full pool cannot starve the
-- other. @stale_seconds uses DB time; @exclude_ids skips refusals.
SELECT
chats_expanded.*,
chat_heartbeats.heartbeat_at AS current_heartbeat_at,
Expand Down
2 changes: 1 addition & 1 deletion coderd/x/chatd/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -959,7 +959,7 @@ Admission decisions are made per pool: root chats (`parent_chat_id IS NULL`) and

When a runner observes its chat leaving the counted statuses, it publishes a bare ownership hint so all workers immediately re-run acquisition and admit queued chats; the acquisition ticker (default 30s) is the fallback for missed nudges, entitlement changes, and heartbeat-staleness releases. Admission clears `capacity_queued_at` in the acquiring transaction and the worker publishes `capacity_change` so open chat pages clear the banner without a reload. `UpdateChatExecutionState`, `UpdateChatStatus`, archive, and auto-archive clear stray markers whenever a chat leaves the counted statuses.

The bypass rule for licensed deployments: uncapped only while `codersdk.FeatureAgentRuntimeHours` is enabled and recorded usage (`Actual`) is below the allocation (`Limit`). A nil `Actual` or `Limit` fails open, because capping paying customers on a missing usage reading is worse than bounded over-admission. Entitlements do not yet populate `Actual` for this feature (runtime-hour usage accrues externally via usage events), so today an enabled runtime-hours license is always uncapped; the exhaustion comparison activates when usage wiring lands. Queued chats admit automatically after entitlement changes via the acquisition ticker.
The bypass rule for licensed deployments: uncapped only while `codersdk.FeatureAgentRuntimeHours` is enabled and recorded usage (`Actual`) is below the allocation (`Limit`). A nil `Actual` or `Limit` fails open, because capping paying customers on a missing usage reading is worse than bounded over-admission. Runtime-hour license decoding leaves `Actual` unset, so enabled licenses take this fail-open path. Queued chats admit automatically after entitlement changes via the acquisition ticker.

## Auto-archive loop

Expand Down
4 changes: 4 additions & 0 deletions coderd/x/chatd/agentadmission.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,10 @@ type AgentAdmission interface {
// Admit reports whether the worker may acquire the chat. Refused chats
// remain unowned and are retried from the capacity queue.
Admit(ctx context.Context, store database.Store, chat database.Chat) (bool, error)
Comment thread
ibetitsmike marked this conversation as resolved.
// RecordQueued observes a chat entering the capacity queue. Workers call
// it only after their queue-marker write wins, so replicas refusing the
// same chat concurrently record one queue entry.
RecordQueued()
}

// AgentAdmissionOptions configures an AgentAdmissionFactory.
Expand Down
75 changes: 54 additions & 21 deletions coderd/x/chatd/agentadmission_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,9 @@ import (
type fakeAdmission struct {
mu sync.Mutex
refused map[uuid.UUID]bool

admitCalls int
queuedRecords int
}

func newFakeAdmission() *fakeAdmission {
Expand All @@ -39,9 +42,28 @@ func (f *fakeAdmission) allow(chatID uuid.UUID) {
func (f *fakeAdmission) Admit(_ context.Context, _ database.Store, chat database.Chat) (bool, error) {
f.mu.Lock()
defer f.mu.Unlock()
f.admitCalls++
return !f.refused[chat.ID], nil
}

func (f *fakeAdmission) RecordQueued() {
f.mu.Lock()
defer f.mu.Unlock()
f.queuedRecords++
}

func (f *fakeAdmission) admitCallCount() int {
f.mu.Lock()
defer f.mu.Unlock()
return f.admitCalls
}

func (f *fakeAdmission) queuedRecordCount() int {
f.mu.Lock()
defer f.mu.Unlock()
return f.queuedRecords
}

func capacityQueuedAt(ctx context.Context, t *testing.T, db database.Store, chatID uuid.UUID) *database.Chat {
t.Helper()
chat, err := db.GetChatByID(ctx, chatID)
Expand Down Expand Up @@ -100,6 +122,32 @@ func TestWorker_AdmissionRefusalQueuesChat(t *testing.T) {
require.Equal(t, 0, chatOwnershipMessages(t, recording, chat.ID))
}

func TestWorker_RefusalRecordsQueueEntryOnce(t *testing.T) {
t.Parallel()
f := newWorkerTestFixture(t)
starter := newRecordingTaskStarter()
admission := newFakeAdmission()
opts := testOptions(t, f, starter)
opts.AgentAdmission = admission

chat := f.createRunningChat(t)
admission.refuse(chat.ID)
worker := startWorker(t, opts)

ctx := testutil.Context(t, testutil.WaitLong)
require.Eventually(t, func() bool {
return capacityQueuedAt(ctx, t, f.db, chat.ID).CapacityQueuedAt.Valid
}, testutil.WaitLong, testutil.IntervalFast)
require.Equal(t, 1, admission.queuedRecordCount())

admitCallsAfterMark := admission.admitCallCount()
worker.Wake()
require.Eventually(t, func() bool {
return admission.admitCallCount() > admitCallsAfterMark
}, testutil.WaitLong, testutil.IntervalFast)
require.Equal(t, 1, admission.queuedRecordCount())
}

func TestWorker_AdmissionAdmitClearsQueueMark(t *testing.T) {
t.Parallel()
f := newWorkerTestFixture(t)
Expand Down Expand Up @@ -159,8 +207,6 @@ func TestWorker_InterruptingSortsBeforeCapacityQueue(t *testing.T) {
require.Equal(t, queued.ID, rows[2].ID)
}

// rootRefusingAdmission simulates a full root pool with free subagent
// capacity: running root chats refuse, everything else admits.
type rootRefusingAdmission struct {
mu sync.Mutex
calls int
Expand All @@ -176,6 +222,8 @@ func (a *rootRefusingAdmission) Admit(_ context.Context, _ database.Store, chat
return chat.ParentChatID.Valid, nil
}

func (*rootRefusingAdmission) RecordQueued() {}

func (a *rootRefusingAdmission) callCount() int {
a.mu.Lock()
defer a.mu.Unlock()
Expand All @@ -190,9 +238,8 @@ func TestWorker_FullPoolDoesNotStarveOtherPool(t *testing.T) {
opts.AgentAdmission = &rootRefusingAdmission{}
ctx := testutil.Context(t, testutil.WaitLong)

// A pre-marked root backlog deeper than two batches; the subagent is
// created last, so an ordering that ignores pools buries it behind a
// batch of pure re-skips and the pass ends before reaching it.
// The root backlog is deep enough to hide the later subagent without
// pool-interleaved ordering.
roots := make([]database.Chat, 0, 2*int(opts.AcquisitionBatchSize)+5)
for range cap(roots) {
chat := f.createRunningChat(t)
Expand All @@ -212,11 +259,6 @@ func TestWorker_FullPoolDoesNotStarveOtherPool(t *testing.T) {
require.Equal(t, sub.ID, call.input.ChatID)
}

// A configured batch size of 1 would only ever surface the
// tie-break-favored root pool: with two already-marked queued roots the
// pass refuses one, re-skips the other without progress, and ends
// before examining the subagent pool. The floor of 2 keeps both pool
// heads in every batch.
func TestWorker_BatchSizeOneCannotHideAPool(t *testing.T) {
t.Parallel()
f := newWorkerTestFixture(t)
Expand All @@ -243,10 +285,6 @@ func TestWorker_BatchSizeOneCannotHideAPool(t *testing.T) {
require.Equal(t, sub.ID, call.input.ChatID)
}

// After the first refusal proves a pool full, the rest of the pass must
// queue-mark that pool's chats without opening refusal transactions.
// All-marked is the pass-completion signal, so the count is stable when
// read.
func TestWorker_FullPoolSkipsRefusalsAfterFirst(t *testing.T) {
t.Parallel()
f := newWorkerTestFixture(t)
Expand Down Expand Up @@ -274,10 +312,6 @@ func TestWorker_FullPoolSkipsRefusalsAfterFirst(t *testing.T) {
"a full pool must be skipped after one refusal, not re-refused per chat")
}

// A persisted queue marker (for example left by an enterprise
// deployment) must clear on acquisition even when no admission hook is
// configured, or the chat page keeps reporting a generating agent as
// queued.
func TestWorker_AcquisitionWithoutAdmissionClearsQueueMark(t *testing.T) {
t.Parallel()
f := newWorkerTestFixture(t)
Expand Down Expand Up @@ -336,6 +370,8 @@ func (runningRefusingAdmission) Admit(_ context.Context, _ database.Store, chat
return chat.Status != database.ChatStatusRunning, nil
}

func (runningRefusingAdmission) RecordQueued() {}

func TestWorker_InterruptClaimsCapacityQueuedChat(t *testing.T) {
t.Parallel()
f := newWorkerTestFixture(t)
Expand All @@ -358,9 +394,6 @@ func TestWorker_InterruptClaimsCapacityQueuedChat(t *testing.T) {
require.Equal(t, chat.ID, call.input.ChatID)
}

// One acquisition pass must page past a full pool's backlog (via
// @exclude_ids) and queue-mark every skipped chat, not just the first
// batch.
func TestWorker_AdmissionPassReachesChatsBeyondRefusedBatch(t *testing.T) {
t.Parallel()
f := newWorkerTestFixture(t)
Expand Down
7 changes: 2 additions & 5 deletions coderd/x/chatd/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -231,11 +231,8 @@ func (o chatWorkerOptions) withDefaults() (chatWorkerOptions, error) {
if o.AcquisitionBatchSize <= 0 {
o.AcquisitionBatchSize = defaultAcquisitionBatchSize
}
// The acquisition query interleaves the two capacity pools, so any
// batch of two or more contains each pool's oldest candidate and an
// all-skipped batch proves nothing admittable remains. A batch of
// one sees only the tie-break-favored pool and would end the pass
// with the other pool unexamined.
// A batch size of one can end the pass on the tie-break-favored pool
// before the other pool's head is examined.
if o.AcquisitionBatchSize < 2 {
o.AcquisitionBatchSize = 2
}
Expand Down
5 changes: 2 additions & 3 deletions coderd/x/chatd/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -178,9 +178,8 @@ func (r *runner) processState(state runnerStateUpdate) {
r.acceptState(state)
}

// occupiesCapacitySlot reports whether a chat in this state counts
// against its concurrency pool. The zero runnerStateUpdate does not, so
// callers need no hasAcceptedState guard.
// The zero value is not capacity-counted, so callers need no
// hasAcceptedState guard.
func occupiesCapacitySlot(state runnerStateUpdate) bool {
return !state.Archived &&
(state.Status == database.ChatStatusRunning || state.Status == database.ChatStatusInterrupting)
Expand Down
19 changes: 8 additions & 11 deletions coderd/x/chatd/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -191,10 +191,8 @@ func (w *chatWorker) acquireOnce(ctx context.Context, workerID uuid.UUID, manage
// Capacity-refused chats remain candidates, so each batch must exclude
// prior attempts or a full pool would hide all later candidates.
excludeIDs := []uuid.UUID{}
// One refusal proves a pool full for the rest of the pass: later
// running chats in that pool are queue-marked without an acquisition
// attempt. Keyed by parent_chat_id presence (false = root pool,
// true = subagent pool).
// One refusal marks the pool full for the rest of the pass. The key is
// parent_chat_id presence: false for roots, true for subagents.
refusedPools := map[bool]bool{}
for {
rows, err := w.opts.Store.GetChatWorkerAcquisitionCandidates(ctx, database.GetChatWorkerAcquisitionCandidatesParams{
Expand Down Expand Up @@ -234,9 +232,7 @@ func (w *chatWorker) acquireOnce(ctx context.Context, workerID uuid.UUID, manage
w.opts.Logger.Warn(ctx, "chatworker acquisition candidate failed", slogError(err))
}
}
// A batch of nothing but re-skipped queued chats means every
// admittable candidate was already tried: the query sorts
// bypassing statuses and each pool's oldest chats first.
// An all-skipped batch means the available pool heads were tried.
if !progressed || len(rows) < int(w.opts.AcquisitionBatchSize) {
return
}
Expand Down Expand Up @@ -310,10 +306,8 @@ func (w *chatWorker) acquireCandidate(
return errCapacityRefused
}
}
// Clear independently of the admission hook: a persisted marker
// must not survive acquisition on a deployment that later runs
// without admission (for example an AGPL build over an
// enterprise database).
// Clear markers left by deployments with admission enabled, even when
// this worker has no admission hook.
if chat.CapacityQueuedAt.Valid {
if _, err := store.ClearChatCapacityQueued(ctx, chat.ID); err != nil {
return xerrors.Errorf("clear capacity queue marker: %w", err)
Expand Down Expand Up @@ -361,6 +355,9 @@ func (w *chatWorker) markCapacityQueued(ctx context.Context, chatID uuid.UUID) {
if marked == 0 {
return
}
if w.opts.AgentAdmission != nil {
w.opts.AgentAdmission.RecordQueued()
Comment thread
ibetitsmike marked this conversation as resolved.
Outdated
}
w.publishCapacityChange(ctx, chatID)
}

Expand Down
14 changes: 6 additions & 8 deletions enterprise/coderd/x/chatd/agentadmission.go
Original file line number Diff line number Diff line change
Expand Up @@ -160,21 +160,15 @@ func (a *admission) Admit(ctx context.Context, store database.Store, chat databa
used, capacity = counts.SubagentCount, a.subagentCapacity
}
if used >= capacity {
if !chat.CapacityQueuedAt.Valid {
a.queueTotal.Inc()
}
return false, nil
}
a.observeAdmission(chat)
return true, nil
}

// uncapped returns true for an enabled entitlement with remaining hours.
// Missing limit or usage data also fails open to avoid capping on incomplete
// entitlement data. Entitlements do not yet populate Actual for this feature
// (usage accrues externally via usage events), so today every enabled
// runtime-hours license is uncapped; the exhaustion branch activates when
// usage wiring lands.
// Missing limit or usage data fails open to avoid capping on incomplete
// entitlement data; runtime-hour license decoding leaves Actual unset.
func (a *admission) uncapped() bool {
f, ok := a.entitlements.Feature(codersdk.FeatureAgentRuntimeHours)
if !ok || !f.Enabled {
Expand All @@ -186,6 +180,10 @@ func (a *admission) uncapped() bool {
return *f.Actual < *f.Limit
Comment thread
ibetitsmike marked this conversation as resolved.
Outdated
}

func (a *admission) RecordQueued() {
a.queueTotal.Inc()
}

func (a *admission) observeAdmission(chat database.Chat) {
if !chat.CapacityQueuedAt.Valid {
return
Expand Down
23 changes: 23 additions & 0 deletions enterprise/coderd/x/chatd/agentadmission_internal_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"testing"

"github.com/google/uuid"
promtestutil "github.com/prometheus/client_golang/prometheus/testutil"
"github.com/stretchr/testify/require"
"golang.org/x/xerrors"

Expand Down Expand Up @@ -118,6 +119,28 @@ func TestAdmission_RootPoolCap(t *testing.T) {
require.True(t, admitted)
}

func TestAdmission_QueueCounterCountsMarkerWinsOnly(t *testing.T) {
t.Parallel()
f := newAdmissionFixture(t)
ctx := testutil.Context(t, testutil.WaitLong)
a := testAdmission(f, nil)

f.occupiedRoot(t)
f.occupiedRoot(t)

chat := f.chat(t, database.Chat{})
for range 2 {
admitted, err := a.Admit(ctx, f.db, chat)
require.NoError(t, err)
require.False(t, admitted)
}
require.Zero(t, promtestutil.ToFloat64(a.queueTotal),
"refusals must not count queue entries: concurrent replicas can refuse the same chat before one marker write wins")

a.RecordQueued()
require.Equal(t, float64(1), promtestutil.ToFloat64(a.queueTotal))
}

func TestAdmission_SubagentPoolCap(t *testing.T) {
t.Parallel()
f := newAdmissionFixture(t)
Expand Down
Loading