mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: remove unused chat statuses pending, paused, and completed (#27064)
The chatd state machine only recognizes `waiting`, `running`, `error`, `requires_action`, and `interrupting`. Remove the unused `pending`, `paused`, and `completed` values from the database enum, backend, SDK, frontend, generated queries, and API docs. Migration `000543_chat_status_remove_unused` remaps existing `pending` rows to `running`, remaps `paused` and `completed` rows to `waiting`, drops the obsolete `idx_chats_pending` index, and recreates `chats_expanded` around the enum swap. It also removes the dead `AcquireChats` query and all remaining query literals for the deleted statuses. **NOTE**: The enum swap can break chat queries from older replicas during a mixed-version rollout because they still reference `'pending'::chat_status`. Chats are experimental, so this PR accepts that limited rollout window instead of adding a two-release expand and contract sequence. > This PR was authored by Mux (AI agent) on Mike's behalf.
This commit is contained in:
@@ -1669,15 +1669,6 @@ func scopedOrgRoleIdentifiers(names []string, orgID uuid.UUID) []rbac.RoleIdenti
|
||||
return out
|
||||
}
|
||||
|
||||
func (q *querier) AcquireChats(ctx context.Context, arg database.AcquireChatsParams) ([]database.Chat, error) {
|
||||
// AcquireChats is a system-level operation used by the chat processor.
|
||||
// Authorization is done at the system level, not per-user.
|
||||
if err := q.authorizeContext(ctx, policy.ActionUpdate, rbac.ResourceChat); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return q.db.AcquireChats(ctx, arg)
|
||||
}
|
||||
|
||||
func (q *querier) AcquireLock(ctx context.Context, id int64) error {
|
||||
return q.db.AcquireLock(ctx, id)
|
||||
}
|
||||
|
||||
@@ -546,16 +546,6 @@ func (s *MethodTestSuite) TestConnectionLogs() {
|
||||
}
|
||||
|
||||
func (s *MethodTestSuite) TestChats() {
|
||||
s.Run("AcquireChats", s.Mocked(func(dbm *dbmock.MockStore, faker *gofakeit.Faker, check *expects) {
|
||||
arg := database.AcquireChatsParams{
|
||||
StartedAt: dbtime.Now(),
|
||||
WorkerID: uuid.New(),
|
||||
NumChats: 1,
|
||||
}
|
||||
chat := testutil.Fake(s.T(), faker, database.Chat{})
|
||||
dbm.EXPECT().AcquireChats(gomock.Any(), arg).Return([]database.Chat{chat}, nil).AnyTimes()
|
||||
check.Args(arg).Asserts(rbac.ResourceChat, policy.ActionUpdate).Returns([]database.Chat{chat})
|
||||
}))
|
||||
s.Run("HydrateAgentChatsContext", s.Mocked(func(dbm *dbmock.MockStore, faker *gofakeit.Faker, check *expects) {
|
||||
arg := database.HydrateAgentChatsContextParams{AgentID: uuid.New()}
|
||||
dbm.EXPECT().HydrateAgentChatsContext(gomock.Any(), arg).Return(nil).AnyTimes()
|
||||
|
||||
-8
@@ -105,14 +105,6 @@ func (m queryMetricsStore) DeleteOrganization(ctx context.Context, id uuid.UUID)
|
||||
return r0
|
||||
}
|
||||
|
||||
func (m queryMetricsStore) AcquireChats(ctx context.Context, arg database.AcquireChatsParams) ([]database.Chat, error) {
|
||||
start := time.Now()
|
||||
r0, r1 := m.s.AcquireChats(ctx, arg)
|
||||
m.queryLatencies.WithLabelValues("AcquireChats").Observe(time.Since(start).Seconds())
|
||||
m.queryCounts.WithLabelValues(httpmw.ExtractHTTPRoute(ctx), httpmw.ExtractHTTPMethod(ctx), "AcquireChats").Inc()
|
||||
return r0, r1
|
||||
}
|
||||
|
||||
func (m queryMetricsStore) AcquireLock(ctx context.Context, pgAdvisoryXactLock int64) error {
|
||||
start := time.Now()
|
||||
r0 := m.s.AcquireLock(ctx, pgAdvisoryXactLock)
|
||||
|
||||
Generated
-15
@@ -45,21 +45,6 @@ func (m *MockStore) EXPECT() *MockStoreMockRecorder {
|
||||
return m.recorder
|
||||
}
|
||||
|
||||
// AcquireChats mocks base method.
|
||||
func (m *MockStore) AcquireChats(ctx context.Context, arg database.AcquireChatsParams) ([]database.Chat, error) {
|
||||
m.ctrl.T.Helper()
|
||||
ret := m.ctrl.Call(m, "AcquireChats", ctx, arg)
|
||||
ret0, _ := ret[0].([]database.Chat)
|
||||
ret1, _ := ret[1].(error)
|
||||
return ret0, ret1
|
||||
}
|
||||
|
||||
// AcquireChats indicates an expected call of AcquireChats.
|
||||
func (mr *MockStoreMockRecorder) AcquireChats(ctx, arg any) *gomock.Call {
|
||||
mr.mock.ctrl.T.Helper()
|
||||
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "AcquireChats", reflect.TypeOf((*MockStore)(nil).AcquireChats), ctx, arg)
|
||||
}
|
||||
|
||||
// AcquireLock mocks base method.
|
||||
func (m *MockStore) AcquireLock(ctx context.Context, pgAdvisoryXactLock int64) error {
|
||||
m.ctrl.T.Helper()
|
||||
|
||||
Generated
-5
@@ -362,10 +362,7 @@ CREATE TYPE chat_reasoning_effort AS ENUM (
|
||||
|
||||
CREATE TYPE chat_status AS ENUM (
|
||||
'waiting',
|
||||
'pending',
|
||||
'running',
|
||||
'paused',
|
||||
'completed',
|
||||
'error',
|
||||
'requires_action',
|
||||
'interrupting'
|
||||
@@ -4763,8 +4760,6 @@ CREATE INDEX idx_chats_owner ON chats USING btree (owner_id);
|
||||
|
||||
CREATE INDEX idx_chats_parent_chat_id ON chats USING btree (parent_chat_id);
|
||||
|
||||
CREATE INDEX idx_chats_pending ON chats USING btree (status) WHERE (status = 'pending'::chat_status);
|
||||
|
||||
CREATE INDEX idx_chats_root_chat_id ON chats USING btree (root_chat_id);
|
||||
|
||||
CREATE INDEX idx_chats_worker_acquisition_candidates ON chats USING btree (status, updated_at, id) WHERE (archived = false);
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
-- No-op: the removed enum values are not restored, matching prior art such
|
||||
-- as 000377 and 000384. Restoring them would require another
|
||||
-- rename-create-cast-drop cycle, and the data cannot be restored anyway:
|
||||
-- rows remapped to 'running' or 'waiting' by the up migration keep their
|
||||
-- new status.
|
||||
@@ -0,0 +1,91 @@
|
||||
-- Remove legacy chat statuses that the chatd state machine treats as
|
||||
-- invalid. 'pending', 'paused', and 'completed' are never written by the
|
||||
-- backend anymore; the valid set is exactly what the state machine
|
||||
-- recognizes: waiting, running, error, requires_action, interrupting.
|
||||
|
||||
-- Remap any historical rows to the closest valid status. The column type
|
||||
-- is still the original chat_status here.
|
||||
--
|
||||
-- 'pending' meant queued work that no runner had picked up yet, so remap
|
||||
-- it to 'running': the worker acquisition query picks up 'running' chats
|
||||
-- without a worker and services them.
|
||||
UPDATE chats SET status = 'running'
|
||||
WHERE status = 'pending';
|
||||
|
||||
-- 'paused' and 'completed' were settled states; 'waiting' is the idle
|
||||
-- resting state and the column default.
|
||||
UPDATE chats SET status = 'waiting'
|
||||
WHERE status IN ('paused', 'completed');
|
||||
|
||||
-- The partial index's WHERE clause references 'pending', which is being
|
||||
-- removed. The index is obsolete now that the legacy AcquireChats query
|
||||
-- is gone.
|
||||
DROP INDEX idx_chats_pending;
|
||||
|
||||
-- The view selects c.status, so it must be dropped before the column's
|
||||
-- type can be altered. It is recreated verbatim below.
|
||||
DROP VIEW chats_expanded;
|
||||
|
||||
-- Recreate the enum without the removed values using the
|
||||
-- rename-create-cast-drop pattern.
|
||||
ALTER TYPE chat_status RENAME TO chat_status_old;
|
||||
CREATE TYPE chat_status AS ENUM (
|
||||
'waiting',
|
||||
'running',
|
||||
'error',
|
||||
'requires_action',
|
||||
'interrupting'
|
||||
);
|
||||
ALTER TABLE chats ALTER COLUMN status DROP DEFAULT;
|
||||
ALTER TABLE chats ALTER COLUMN status TYPE chat_status USING status::text::chat_status;
|
||||
ALTER TABLE chats ALTER COLUMN status SET DEFAULT 'waiting';
|
||||
DROP TYPE chat_status_old;
|
||||
|
||||
CREATE VIEW chats_expanded AS
|
||||
SELECT c.id,
|
||||
c.owner_id,
|
||||
c.workspace_id,
|
||||
c.title,
|
||||
c.status,
|
||||
c.worker_id,
|
||||
c.started_at,
|
||||
c.heartbeat_at,
|
||||
c.created_at,
|
||||
c.updated_at,
|
||||
c.parent_chat_id,
|
||||
c.root_chat_id,
|
||||
c.last_model_config_id,
|
||||
c.last_reasoning_effort,
|
||||
c.archived,
|
||||
c.last_error,
|
||||
c.mode,
|
||||
c.mcp_server_ids,
|
||||
c.labels,
|
||||
c.build_id,
|
||||
c.agent_id,
|
||||
c.pin_order,
|
||||
c.last_read_message_id,
|
||||
c.dynamic_tools,
|
||||
c.organization_id,
|
||||
c.plan_mode,
|
||||
c.client_type,
|
||||
c.last_turn_summary,
|
||||
c.snapshot_version,
|
||||
c.history_version,
|
||||
c.queue_version,
|
||||
c.generation_attempt,
|
||||
c.retry_state,
|
||||
c.retry_state_version,
|
||||
c.runner_id,
|
||||
c.requires_action_deadline_at,
|
||||
COALESCE(root.user_acl, c.user_acl) AS user_acl,
|
||||
COALESCE(root.group_acl, c.group_acl) AS group_acl,
|
||||
owner.username AS owner_username,
|
||||
owner.name AS owner_name,
|
||||
c.context_aggregate_hash,
|
||||
c.context_dirty_since,
|
||||
c.context_dirty_resources,
|
||||
c.context_error
|
||||
FROM ((chats c
|
||||
LEFT JOIN chats root ON ((root.id = COALESCE(c.root_chat_id, c.parent_chat_id))))
|
||||
JOIN visible_users owner ON ((owner.id = c.owner_id)));
|
||||
Generated
-9
@@ -1725,10 +1725,7 @@ type ChatStatus string
|
||||
|
||||
const (
|
||||
ChatStatusWaiting ChatStatus = "waiting"
|
||||
ChatStatusPending ChatStatus = "pending"
|
||||
ChatStatusRunning ChatStatus = "running"
|
||||
ChatStatusPaused ChatStatus = "paused"
|
||||
ChatStatusCompleted ChatStatus = "completed"
|
||||
ChatStatusError ChatStatus = "error"
|
||||
ChatStatusRequiresAction ChatStatus = "requires_action"
|
||||
ChatStatusInterrupting ChatStatus = "interrupting"
|
||||
@@ -1772,10 +1769,7 @@ func (ns NullChatStatus) Value() (driver.Value, error) {
|
||||
func (e ChatStatus) Valid() bool {
|
||||
switch e {
|
||||
case ChatStatusWaiting,
|
||||
ChatStatusPending,
|
||||
ChatStatusRunning,
|
||||
ChatStatusPaused,
|
||||
ChatStatusCompleted,
|
||||
ChatStatusError,
|
||||
ChatStatusRequiresAction,
|
||||
ChatStatusInterrupting:
|
||||
@@ -1787,10 +1781,7 @@ func (e ChatStatus) Valid() bool {
|
||||
func AllChatStatusValues() []ChatStatus {
|
||||
return []ChatStatus{
|
||||
ChatStatusWaiting,
|
||||
ChatStatusPending,
|
||||
ChatStatusRunning,
|
||||
ChatStatusPaused,
|
||||
ChatStatusCompleted,
|
||||
ChatStatusError,
|
||||
ChatStatusRequiresAction,
|
||||
ChatStatusInterrupting,
|
||||
|
||||
Generated
-3
@@ -13,9 +13,6 @@ import (
|
||||
)
|
||||
|
||||
type sqlcQuerier interface {
|
||||
// Acquires up to @num_chats pending chats for processing. Uses SKIP LOCKED
|
||||
// to prevent multiple replicas from acquiring the same chat.
|
||||
AcquireChats(ctx context.Context, arg AcquireChatsParams) ([]Chat, error)
|
||||
// Blocks until the lock is acquired.
|
||||
//
|
||||
// This must be called from within a transaction. The lock will be automatically
|
||||
|
||||
@@ -1280,11 +1280,11 @@ func TestChatContextHydration(t *testing.T) {
|
||||
hashH := []byte{0x01, 0x02, 0x03}
|
||||
hashOther := []byte{0xff, 0xee}
|
||||
|
||||
chatNull := newChat(database.ChatStatusWaiting, agent.ID) // never hydrated
|
||||
chatMatch := newChat(database.ChatStatusRunning, agent.ID) // already at hashH
|
||||
chatDrift := newChat(database.ChatStatusRunning, agent.ID) // drifted, active
|
||||
chatTerminal := newChat(database.ChatStatusCompleted, agent.ID) // drifted, terminal
|
||||
chatArchived := newChat(database.ChatStatusRunning, agent.ID) // drifted, archived
|
||||
chatNull := newChat(database.ChatStatusWaiting, agent.ID) // never hydrated
|
||||
chatMatch := newChat(database.ChatStatusRunning, agent.ID) // already at hashH
|
||||
chatDrift := newChat(database.ChatStatusRunning, agent.ID) // drifted, active
|
||||
chatTerminal := newChat(database.ChatStatusError, agent.ID) // drifted, terminal
|
||||
chatArchived := newChat(database.ChatStatusRunning, agent.ID) // drifted, archived
|
||||
chatOtherAgent := newChat(database.ChatStatusRunning, otherAgent.ID)
|
||||
|
||||
// Pin starting hashes; chatNull is intentionally left NULL.
|
||||
@@ -13019,7 +13019,7 @@ func TestChatPinOrderConstraints(t *testing.T) {
|
||||
|
||||
parent, err := db.InsertChat(ctx, database.InsertChatParams{
|
||||
OrganizationID: org.ID,
|
||||
Status: database.ChatStatusCompleted,
|
||||
Status: database.ChatStatusWaiting,
|
||||
ClientType: database.ChatClientTypeUi,
|
||||
OwnerID: owner.ID,
|
||||
LastModelConfigID: modelCfg.ID,
|
||||
@@ -13029,7 +13029,7 @@ func TestChatPinOrderConstraints(t *testing.T) {
|
||||
|
||||
child, err := db.InsertChat(ctx, database.InsertChatParams{
|
||||
OrganizationID: org.ID,
|
||||
Status: database.ChatStatusCompleted,
|
||||
Status: database.ChatStatusWaiting,
|
||||
ClientType: database.ChatClientTypeUi,
|
||||
OwnerID: owner.ID,
|
||||
LastModelConfigID: modelCfg.ID,
|
||||
@@ -13050,7 +13050,7 @@ func TestChatPinOrderConstraints(t *testing.T) {
|
||||
|
||||
chat, err := db.InsertChat(ctx, database.InsertChatParams{
|
||||
OrganizationID: org.ID,
|
||||
Status: database.ChatStatusCompleted,
|
||||
Status: database.ChatStatusWaiting,
|
||||
ClientType: database.ChatClientTypeUi,
|
||||
OwnerID: owner.ID,
|
||||
LastModelConfigID: modelCfg.ID,
|
||||
|
||||
Generated
+5
-167
@@ -5606,165 +5606,6 @@ func (q *sqlQuerier) UpdateChatModelConfig(ctx context.Context, arg UpdateChatMo
|
||||
return i, err
|
||||
}
|
||||
|
||||
const acquireChats = `-- name: AcquireChats :many
|
||||
WITH acquired_chats AS (
|
||||
UPDATE
|
||||
chats
|
||||
SET
|
||||
status = 'running'::chat_status,
|
||||
started_at = $1::timestamptz,
|
||||
heartbeat_at = $1::timestamptz,
|
||||
updated_at = $1::timestamptz,
|
||||
worker_id = $2::uuid
|
||||
WHERE
|
||||
id = ANY(
|
||||
SELECT
|
||||
id
|
||||
FROM
|
||||
chats
|
||||
WHERE
|
||||
status = 'pending'::chat_status
|
||||
AND archived = false
|
||||
ORDER BY
|
||||
updated_at ASC
|
||||
FOR UPDATE
|
||||
SKIP LOCKED
|
||||
LIMIT
|
||||
$3::int
|
||||
)
|
||||
RETURNING id, owner_id, workspace_id, title, status, worker_id, started_at, heartbeat_at, created_at, updated_at, parent_chat_id, root_chat_id, last_model_config_id, archived, last_error, mode, mcp_server_ids, labels, build_id, agent_id, pin_order, last_read_message_id, dynamic_tools, organization_id, plan_mode, client_type, last_turn_summary, user_acl, group_acl, snapshot_version, history_version, queue_version, generation_attempt, retry_state, retry_state_version, runner_id, requires_action_deadline_at, context_aggregate_hash, context_dirty_since, context_dirty_resources, context_error, last_reasoning_effort
|
||||
),
|
||||
chats_expanded AS (
|
||||
SELECT
|
||||
acquired_chats.id,
|
||||
acquired_chats.owner_id,
|
||||
acquired_chats.workspace_id,
|
||||
acquired_chats.title,
|
||||
acquired_chats.status,
|
||||
acquired_chats.worker_id,
|
||||
acquired_chats.started_at,
|
||||
acquired_chats.heartbeat_at,
|
||||
acquired_chats.created_at,
|
||||
acquired_chats.updated_at,
|
||||
acquired_chats.parent_chat_id,
|
||||
acquired_chats.root_chat_id,
|
||||
acquired_chats.last_model_config_id,
|
||||
acquired_chats.last_reasoning_effort,
|
||||
acquired_chats.archived,
|
||||
acquired_chats.last_error,
|
||||
acquired_chats.mode,
|
||||
acquired_chats.mcp_server_ids,
|
||||
acquired_chats.labels,
|
||||
acquired_chats.build_id,
|
||||
acquired_chats.agent_id,
|
||||
acquired_chats.pin_order,
|
||||
acquired_chats.last_read_message_id,
|
||||
acquired_chats.dynamic_tools,
|
||||
acquired_chats.organization_id,
|
||||
acquired_chats.plan_mode,
|
||||
acquired_chats.client_type,
|
||||
acquired_chats.last_turn_summary,
|
||||
acquired_chats.snapshot_version,
|
||||
acquired_chats.history_version,
|
||||
acquired_chats.queue_version,
|
||||
acquired_chats.generation_attempt,
|
||||
acquired_chats.retry_state,
|
||||
acquired_chats.retry_state_version,
|
||||
acquired_chats.runner_id,
|
||||
acquired_chats.requires_action_deadline_at,
|
||||
COALESCE(root.user_acl, acquired_chats.user_acl) AS user_acl,
|
||||
COALESCE(root.group_acl, acquired_chats.group_acl) AS group_acl,
|
||||
owner.username AS owner_username,
|
||||
owner.name AS owner_name,
|
||||
acquired_chats.context_aggregate_hash,
|
||||
acquired_chats.context_dirty_since,
|
||||
acquired_chats.context_dirty_resources,
|
||||
acquired_chats.context_error
|
||||
FROM
|
||||
acquired_chats
|
||||
LEFT JOIN chats root ON root.id = COALESCE(acquired_chats.root_chat_id, acquired_chats.parent_chat_id)
|
||||
JOIN visible_users owner ON owner.id = acquired_chats.owner_id
|
||||
)
|
||||
SELECT id, owner_id, workspace_id, title, status, worker_id, started_at, heartbeat_at, created_at, updated_at, parent_chat_id, root_chat_id, last_model_config_id, last_reasoning_effort, archived, last_error, mode, mcp_server_ids, labels, build_id, agent_id, pin_order, last_read_message_id, dynamic_tools, organization_id, plan_mode, client_type, last_turn_summary, snapshot_version, history_version, queue_version, generation_attempt, retry_state, retry_state_version, runner_id, requires_action_deadline_at, user_acl, group_acl, owner_username, owner_name, context_aggregate_hash, context_dirty_since, context_dirty_resources, context_error
|
||||
FROM chats_expanded
|
||||
`
|
||||
|
||||
type AcquireChatsParams struct {
|
||||
StartedAt time.Time `db:"started_at" json:"started_at"`
|
||||
WorkerID uuid.UUID `db:"worker_id" json:"worker_id"`
|
||||
NumChats int32 `db:"num_chats" json:"num_chats"`
|
||||
}
|
||||
|
||||
// Acquires up to @num_chats pending chats for processing. Uses SKIP LOCKED
|
||||
// to prevent multiple replicas from acquiring the same chat.
|
||||
func (q *sqlQuerier) AcquireChats(ctx context.Context, arg AcquireChatsParams) ([]Chat, error) {
|
||||
rows, err := q.db.QueryContext(ctx, acquireChats, arg.StartedAt, arg.WorkerID, arg.NumChats)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
var items []Chat
|
||||
for rows.Next() {
|
||||
var i Chat
|
||||
if err := rows.Scan(
|
||||
&i.ID,
|
||||
&i.OwnerID,
|
||||
&i.WorkspaceID,
|
||||
&i.Title,
|
||||
&i.Status,
|
||||
&i.WorkerID,
|
||||
&i.StartedAt,
|
||||
&i.HeartbeatAt,
|
||||
&i.CreatedAt,
|
||||
&i.UpdatedAt,
|
||||
&i.ParentChatID,
|
||||
&i.RootChatID,
|
||||
&i.LastModelConfigID,
|
||||
&i.LastReasoningEffort,
|
||||
&i.Archived,
|
||||
&i.LastError,
|
||||
&i.Mode,
|
||||
pq.Array(&i.MCPServerIDs),
|
||||
&i.Labels,
|
||||
&i.BuildID,
|
||||
&i.AgentID,
|
||||
&i.PinOrder,
|
||||
&i.LastReadMessageID,
|
||||
&i.DynamicTools,
|
||||
&i.OrganizationID,
|
||||
&i.PlanMode,
|
||||
&i.ClientType,
|
||||
&i.LastTurnSummary,
|
||||
&i.SnapshotVersion,
|
||||
&i.HistoryVersion,
|
||||
&i.QueueVersion,
|
||||
&i.GenerationAttempt,
|
||||
&i.RetryState,
|
||||
&i.RetryStateVersion,
|
||||
&i.RunnerID,
|
||||
&i.RequiresActionDeadlineAt,
|
||||
&i.UserACL,
|
||||
&i.GroupACL,
|
||||
&i.OwnerUsername,
|
||||
&i.OwnerName,
|
||||
&i.ContextAggregateHash,
|
||||
&i.ContextDirtySince,
|
||||
&i.ContextDirtyResources,
|
||||
&i.ContextError,
|
||||
); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
items = append(items, i)
|
||||
}
|
||||
if err := rows.Close(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return items, nil
|
||||
}
|
||||
|
||||
const acquireStaleChatDiffStatuses = `-- name: AcquireStaleChatDiffStatuses :many
|
||||
WITH acquired AS (
|
||||
UPDATE
|
||||
@@ -6036,7 +5877,7 @@ WITH to_archive AS (
|
||||
-- Redundant filter helps the planner use the partial index on created_at.
|
||||
AND c.created_at < $1::timestamptz
|
||||
-- New active statuses must be added here to prevent archiving.
|
||||
AND c.status NOT IN ('running', 'pending', 'paused', 'requires_action')
|
||||
AND c.status NOT IN ('running', 'requires_action')
|
||||
AND COALESCE(activity.last_activity_at, c.created_at) < $1::timestamptz
|
||||
-- Sorting by created_at lets Postgres drive the scan from the
|
||||
-- partial index instead of evaluating every LATERAL subquery
|
||||
@@ -6443,10 +6284,9 @@ SELECT id, owner_id, workspace_id, title, status, worker_id, started_at, heartbe
|
||||
FROM chats_expanded
|
||||
WHERE agent_id = $1::uuid
|
||||
AND archived = false
|
||||
-- Active statuses only: waiting, pending, running, paused,
|
||||
-- requires_action.
|
||||
-- Excludes completed and error (terminal states).
|
||||
AND status IN ('waiting', 'running', 'paused', 'pending', 'requires_action')
|
||||
-- Active statuses only: waiting, running, requires_action.
|
||||
-- Excludes error (terminal state) and interrupting.
|
||||
AND status IN ('waiting', 'running', 'requires_action')
|
||||
ORDER BY updated_at DESC
|
||||
`
|
||||
|
||||
@@ -6538,8 +6378,6 @@ WHERE
|
||||
AND chats_expanded.status NOT IN (
|
||||
'running'::chat_status,
|
||||
'interrupting'::chat_status,
|
||||
'pending'::chat_status,
|
||||
'paused'::chat_status,
|
||||
'requires_action'::chat_status
|
||||
)
|
||||
AND COALESCE(activity.last_activity_at, chats_expanded.created_at) < $1::timestamptz
|
||||
@@ -10400,7 +10238,7 @@ UPDATE chats
|
||||
SET context_dirty_since = $1
|
||||
WHERE agent_id = $2::uuid
|
||||
AND archived = false
|
||||
AND status IN ('waiting', 'running', 'paused', 'pending', 'requires_action')
|
||||
AND status IN ('waiting', 'running', 'requires_action')
|
||||
AND context_aggregate_hash IS NOT NULL
|
||||
AND context_aggregate_hash IS DISTINCT FROM $3
|
||||
AND context_dirty_since IS NULL
|
||||
|
||||
@@ -1469,7 +1469,7 @@ UPDATE chats
|
||||
SET context_dirty_since = @dirty_since
|
||||
WHERE agent_id = @agent_id::uuid
|
||||
AND archived = false
|
||||
AND status IN ('waiting', 'running', 'paused', 'pending', 'requires_action')
|
||||
AND status IN ('waiting', 'running', 'requires_action')
|
||||
AND context_aggregate_hash IS NOT NULL
|
||||
AND context_aggregate_hash IS DISTINCT FROM @aggregate_hash
|
||||
AND context_dirty_since IS NULL
|
||||
@@ -1540,90 +1540,6 @@ SELECT
|
||||
(SELECT COUNT(*)::int FROM genuinely_new) -
|
||||
(SELECT COUNT(*)::int FROM inserted) AS rejected_new_files;
|
||||
|
||||
-- name: AcquireChats :many
|
||||
-- Acquires up to @num_chats pending chats for processing. Uses SKIP LOCKED
|
||||
-- to prevent multiple replicas from acquiring the same chat.
|
||||
WITH acquired_chats AS (
|
||||
UPDATE
|
||||
chats
|
||||
SET
|
||||
status = 'running'::chat_status,
|
||||
started_at = @started_at::timestamptz,
|
||||
heartbeat_at = @started_at::timestamptz,
|
||||
updated_at = @started_at::timestamptz,
|
||||
worker_id = @worker_id::uuid
|
||||
WHERE
|
||||
id = ANY(
|
||||
SELECT
|
||||
id
|
||||
FROM
|
||||
chats
|
||||
WHERE
|
||||
status = 'pending'::chat_status
|
||||
AND archived = false
|
||||
ORDER BY
|
||||
updated_at ASC
|
||||
FOR UPDATE
|
||||
SKIP LOCKED
|
||||
LIMIT
|
||||
@num_chats::int
|
||||
)
|
||||
RETURNING *
|
||||
),
|
||||
chats_expanded AS (
|
||||
SELECT
|
||||
acquired_chats.id,
|
||||
acquired_chats.owner_id,
|
||||
acquired_chats.workspace_id,
|
||||
acquired_chats.title,
|
||||
acquired_chats.status,
|
||||
acquired_chats.worker_id,
|
||||
acquired_chats.started_at,
|
||||
acquired_chats.heartbeat_at,
|
||||
acquired_chats.created_at,
|
||||
acquired_chats.updated_at,
|
||||
acquired_chats.parent_chat_id,
|
||||
acquired_chats.root_chat_id,
|
||||
acquired_chats.last_model_config_id,
|
||||
acquired_chats.last_reasoning_effort,
|
||||
acquired_chats.archived,
|
||||
acquired_chats.last_error,
|
||||
acquired_chats.mode,
|
||||
acquired_chats.mcp_server_ids,
|
||||
acquired_chats.labels,
|
||||
acquired_chats.build_id,
|
||||
acquired_chats.agent_id,
|
||||
acquired_chats.pin_order,
|
||||
acquired_chats.last_read_message_id,
|
||||
acquired_chats.dynamic_tools,
|
||||
acquired_chats.organization_id,
|
||||
acquired_chats.plan_mode,
|
||||
acquired_chats.client_type,
|
||||
acquired_chats.last_turn_summary,
|
||||
acquired_chats.snapshot_version,
|
||||
acquired_chats.history_version,
|
||||
acquired_chats.queue_version,
|
||||
acquired_chats.generation_attempt,
|
||||
acquired_chats.retry_state,
|
||||
acquired_chats.retry_state_version,
|
||||
acquired_chats.runner_id,
|
||||
acquired_chats.requires_action_deadline_at,
|
||||
COALESCE(root.user_acl, acquired_chats.user_acl) AS user_acl,
|
||||
COALESCE(root.group_acl, acquired_chats.group_acl) AS group_acl,
|
||||
owner.username AS owner_username,
|
||||
owner.name AS owner_name,
|
||||
acquired_chats.context_aggregate_hash,
|
||||
acquired_chats.context_dirty_since,
|
||||
acquired_chats.context_dirty_resources,
|
||||
acquired_chats.context_error
|
||||
FROM
|
||||
acquired_chats
|
||||
LEFT JOIN chats root ON root.id = COALESCE(acquired_chats.root_chat_id, acquired_chats.parent_chat_id)
|
||||
JOIN visible_users owner ON owner.id = acquired_chats.owner_id
|
||||
)
|
||||
SELECT *
|
||||
FROM chats_expanded;
|
||||
|
||||
-- name: UpdateChatStatus :one
|
||||
WITH updated_chat AS (
|
||||
UPDATE
|
||||
@@ -2534,10 +2450,9 @@ SELECT *
|
||||
FROM chats_expanded
|
||||
WHERE agent_id = @agent_id::uuid
|
||||
AND archived = false
|
||||
-- Active statuses only: waiting, pending, running, paused,
|
||||
-- requires_action.
|
||||
-- Excludes completed and error (terminal states).
|
||||
AND status IN ('waiting', 'running', 'paused', 'pending', 'requires_action')
|
||||
-- Active statuses only: waiting, running, requires_action.
|
||||
-- Excludes error (terminal state) and interrupting.
|
||||
AND status IN ('waiting', 'running', 'requires_action')
|
||||
ORDER BY updated_at DESC;
|
||||
|
||||
-- name: SoftDeleteContextFileMessages :exec
|
||||
@@ -2630,8 +2545,6 @@ WHERE
|
||||
AND chats_expanded.status NOT IN (
|
||||
'running'::chat_status,
|
||||
'interrupting'::chat_status,
|
||||
'pending'::chat_status,
|
||||
'paused'::chat_status,
|
||||
'requires_action'::chat_status
|
||||
)
|
||||
AND COALESCE(activity.last_activity_at, chats_expanded.created_at) < @archive_cutoff::timestamptz
|
||||
@@ -2999,7 +2912,7 @@ WITH to_archive AS (
|
||||
-- Redundant filter helps the planner use the partial index on created_at.
|
||||
AND c.created_at < @archive_cutoff::timestamptz
|
||||
-- New active statuses must be added here to prevent archiving.
|
||||
AND c.status NOT IN ('running', 'pending', 'paused', 'requires_action')
|
||||
AND c.status NOT IN ('running', 'requires_action')
|
||||
AND COALESCE(activity.last_activity_at, c.created_at) < @archive_cutoff::timestamptz
|
||||
-- Sorting by created_at lets Postgres drive the scan from the
|
||||
-- partial index instead of evaluating every LATERAL subquery
|
||||
|
||||
Reference in New Issue
Block a user