PR #3465
Sections
Review

feat(queue): move automation toggles to session scope

main ← feature/move-queue-toggles-p-dxs 197 files +4821 −1184 PR #3465 ↗

Queue Auto-merge moves from a global-only switch to a per-session override that inherits the global value until the user changes it in the queue panel, with incarnation-bound persistence and ordering.

Why this change

Auto-merge was a single install-wide setting. Every session shared one value, so a user could not keep one session merged and another separate without changing the global default for all sessions.

What it does

Architecture, end to end

Global settings publish one immutable snapshot. Each session resolves effective policy from override or global, then admits and reports status under task-then-session locks with incarnation validation.

flowchart LR
  Admin[Settings > Task Behavior] --> GlobalStore[(queuesettings Store)]
  GlobalStore --> Snapshot[messagequeue.Service atomic snapshot]
  Snapshot --> Resolve{Resolve effective policy}
  SessionState[(queue_session_state)] --> Resolve
  Resolve --> Admission[Admission with claim and fold]
  Admission --> Queue[(queue rows)]
  Queue --> Status[QueueStatus snapshot]
  Status --> WS[WS status_changed]
  Status --> Panel[QueuePanelHeader pills]
  Browser[Browser useQueue] --> Admission
  Browser --> Status

Key code changes

Drag to pan. Use the + and − buttons to zoom. Click a node to open the full code. The arrows show how the parts interact.

drag to pan · +/− to zoom · click a node for details

func provideOrchestrator(...) (*orchestrator.Service, error)
Click for details →

Hydrates the queue service with the persisted global value and its monotonic revision before admissions start.

Before and after
queueResolution := resolveQueueSettingsWithStore(settingsStore, pool, log, queueConfiguration(cfg))
queueSettings := queueResolution.Effective
maxPerSession := queueSettings.MaxPerSession
mergeEnabled := queueSettings.MergeEnabled
autoMergeEnabled := queueSettings.AutoMergeEnabled
msgQueue := messagequeue.NewService(queueRepo, maxPerSession, log)
msgQueue.SetMergeEnabled(mergeEnabled)
msgQueue.SetAutoMergePolicy(autoMergeEnabled, queueResolution.Settings.AutoMergeRevision)
func (s *Service) ResolveAutoMergePolicy(ctx context.Context, identity QueueSessionIdentity) (AutoMergePolicy, error)
Click for details →

Resolves session override or global snapshot, then admits with policy-aware insert and full-queue candidate fold.

Policy types
type QueueSessionIdentity struct {
  TaskID               string `json:"task_id"`
  SessionID            string `json:"session_id"`
  SessionIncarnationID string `json:"session_incarnation_id"`
}
type AutoMergePolicy struct {
  Enabled  bool
  Source   AutoMergeSource
  Revision int64
}
Resolve
func (s *Service) ResolveAutoMergePolicy(ctx context.Context, identity QueueSessionIdentity) (AutoMergePolicy, error) {
  override, err := s.repo.GetAutoMergeOverride(ctx, identity)
  if err != nil {
    return AutoMergePolicy{}, err
  }
  if override == nil {
    return s.loadGlobalAutoMergePolicy(ctx)
  }
  return AutoMergePolicy{Enabled: override.Enabled, Source: AutoMergeSourceSession, Revision: override.Revision}, nil
}
Admission loop
policy := s.resolveAdmissionAutoMergePolicy(admittedCtx, identity, sessionID)
source, insertErr := s.insertQueueMessageWithMetadata(admittedCtx, identity, sessionID, taskID, content, model, userID, planMode, attachments, metadata, claim, s.MaxPerSession(), policy)
if errors.Is(insertErr, ErrAutoMergePolicyChanged) {
  continue
}
if errors.Is(insertErr, ErrQueueFull) && policy != nil && policy.Enabled && claim == nil {
  _, merged, _ := s.admitQueueFullMessage(admittedCtx, identity, insertErr, sessionID, taskID, content, model, userID, planMode, attachments, metadata, claim, s.MaxPerSession(), policy)
  if merged != nil { queued = merged; return nil }
}
queued = s.finalizeAutoMerge(admittedCtx, identity, source, policy)
Session override and incarnation persistenceapps/backend/internal/orchestrator/messagequeue/types.go ↗
type QueueStatus struct
Click for details →

Carries effective policy, source, revision, and incarnation so clients can order and validate status.

QueueStatus
type QueueStatus struct {
  Entries              []QueuedMessage `json:"entries"`
  Count                int             `json:"count"`
  Max                  int             `json:"max"`
  TaskID               string          `json:"task_id,omitempty"`
  SessionID            string          `json:"session_id,omitempty"`
  SessionIncarnationID string          `json:"session_incarnation_id,omitempty"`
  StatusEpoch          string          `json:"status_epoch,omitempty"`
  StatusGeneration     int64           `json:"status_generation,omitempty"`
  AutoRun              bool            `json:"auto_run"`
  MergeEnabled         bool            `json:"merge_enabled"`
  AutoMergeAvailable   bool            `json:"auto_merge_available"`
  AutoMergeEnabled     *bool           `json:"auto_merge_enabled,omitempty"`
  AutoMergeSource      AutoMergeSource `json:"auto_merge_source,omitempty"`
  AutoMergeRevision    *int64          `json:"auto_merge_revision,omitempty"`
}
Repository contract
GetAutoMergeOverride(ctx context.Context, identity QueueSessionIdentity) (*AutoMergeOverride, error)
SetAutoMergeOverride(ctx context.Context, identity QueueSessionIdentity, enabled bool) (AutoMergeOverride, error)
Snapshot(ctx context.Context, identity QueueSessionIdentity) (RepositorySnapshot, error)
func (h *QueueHandlers) admitQueuedMessage(ctx context.Context, req *wsQueueMessageRequest, queuedBy string, metadata map[string]interface{}) (*messagequeue.QueuedMessage, error)
Click for details →

Requires the task/session/incarnation triplet, authorizes it, and passes the unchanged identity to persistence.

Admission
func (h *QueueHandlers) admitQueuedMessage(ctx context.Context, req *wsQueueMessageRequest, queuedBy string, metadata map[string]interface{}) (*messagequeue.QueuedMessage, error) {
  if !h.requiresQueueIdentity() {
    return h.queueService.QueueMessageWithMetadata(ctx, req.SessionID, req.TaskID, req.Content, req.Model, queuedBy, req.PlanMode, req.Attachments, metadata)
  }
  identity := messagequeue.QueueSessionIdentity{TaskID: req.TaskID, SessionID: req.SessionID, SessionIncarnationID: req.SessionIncarnationID}
  if h.attachmentClaimer == nil || len(req.Attachments) == 0 {
    return admissions.QueueMessageWithMetadataForSession(ctx, identity, req.Content, req.Model, queuedBy, req.PlanMode, req.Attachments, metadata)
  }
  claim, _ := preparer.PrepareQueueAttachmentClaim(ctx, req.TaskID, queueAttachmentsToV1(req.Attachments))
  return atomicAdmissions.QueueMessageWithMetadataForSessionWithClaim(ctx, identity, req.Content, req.Model, queuedBy, req.PlanMode, req.Attachments, metadata, claim)
}
Validation
if h.requiresQueueIdentity() && req.SessionIncarnationID == "" {
  return ws.NewError(msg.ID, msg.Action, ws.ErrorCodeValidation, "task_id, session_id, and session_incarnation_id are required", nil)
}
if denied := h.authorizeQueueIdentity(ctx, msg, req.TaskID, req.SessionID, req.SessionIncarnationID); denied != nil {
  return denied, nil
}
function QueueAutomationPills(props: QueueAutomationPillsProps)
Click for details →

Renders Auto-run and Auto-merge as compact pills with shared disabled state and availability gating.

Pills
function QueueAutomationPills({ autoRun, autoMerge, autoMergeAvailable, isLoading, cancellationPending, onAutoRunChange, onAutoMergeChange }: QueueAutomationPillsProps) {
  const controlsDisabled = isLoading || cancellationPending
  return (
    <div className="flex min-w-0 flex-wrap items-center gap-1.5">
      <label className="flex min-h-7 items-center gap-1.5 rounded-full border px-2 text-xs">
        <span>{t("chat:queueAutoRun")}</span>
        <Switch checked={autoRun} disabled={controlsDisabled} onCheckedChange={onAutoRunChange} />
      </label>
      <label className={autoMergeAvailable ? "cursor-pointer" : "opacity-50 pointer-events-none"}>
        <span>{t("chat:queueAutoMerge")}</span>
        <Switch checked={autoMerge} disabled={controlsDisabled || !autoMergeAvailable} onCheckedChange={onAutoMergeChange} />
      </label>
    </div>
  )
}
Frontend incarnation and orderingapps/web/hooks/domains/session/use-queue.ts ↗
export function useQueue(sessionId: string | null)
Click for details →

Binds every queue mutation to the current incarnation and orders status by epoch, generation, source, and revision.

Identity
function useCurrentQueueIdentity(sessionId: string | null, taskSession: TaskSession | undefined) {
  return useMemo(() => {
    if (!sessionId || !taskSession?.task_id || !taskSession.queue_incarnation_id) return null
    return { task_id: taskSession.task_id, session_id: sessionId, session_incarnation_id: taskSession.queue_incarnation_id }
  }, [sessionId, taskSession?.task_id, taskSession?.queue_incarnation_id])
}
Mutation token
async function runQueueMutation(identity, beginQueueOperation, finishQueueOperation, refetch, mutate) {
  if (!identity) return
  const token = beginQueueOperation(sessionId, incarnationId)
  if (!token) return
  try {
    await mutate(identity)
    await refetch(sessionId, token)
  } finally {
    finishQueueOperation(sessionId, token)
  }
}
Read the changes as a list

Startup hydration with revision

apps/backend/internal/backendapp/orchestrator.go

Hydrates the queue service with the persisted global value and its monotonic revision before admissions start.

Before and after
queueResolution := resolveQueueSettingsWithStore(settingsStore, pool, log, queueConfiguration(cfg))
queueSettings := queueResolution.Effective
maxPerSession := queueSettings.MaxPerSession
mergeEnabled := queueSettings.MergeEnabled
autoMergeEnabled := queueSettings.AutoMergeEnabled
msgQueue := messagequeue.NewService(queueRepo, maxPerSession, log)
msgQueue.SetMergeEnabled(mergeEnabled)
msgQueue.SetAutoMergePolicy(autoMergeEnabled, queueResolution.Settings.AutoMergeRevision)

Effective policy resolution and admission

apps/backend/internal/orchestrator/messagequeue/service.go

Resolves session override or global snapshot, then admits with policy-aware insert and full-queue candidate fold.

Policy types
type QueueSessionIdentity struct {
  TaskID               string `json:"task_id"`
  SessionID            string `json:"session_id"`
  SessionIncarnationID string `json:"session_incarnation_id"`
}
type AutoMergePolicy struct {
  Enabled  bool
  Source   AutoMergeSource
  Revision int64
}
Resolve
func (s *Service) ResolveAutoMergePolicy(ctx context.Context, identity QueueSessionIdentity) (AutoMergePolicy, error) {
  override, err := s.repo.GetAutoMergeOverride(ctx, identity)
  if err != nil {
    return AutoMergePolicy{}, err
  }
  if override == nil {
    return s.loadGlobalAutoMergePolicy(ctx)
  }
  return AutoMergePolicy{Enabled: override.Enabled, Source: AutoMergeSourceSession, Revision: override.Revision}, nil
}
Admission loop
policy := s.resolveAdmissionAutoMergePolicy(admittedCtx, identity, sessionID)
source, insertErr := s.insertQueueMessageWithMetadata(admittedCtx, identity, sessionID, taskID, content, model, userID, planMode, attachments, metadata, claim, s.MaxPerSession(), policy)
if errors.Is(insertErr, ErrAutoMergePolicyChanged) {
  continue
}
if errors.Is(insertErr, ErrQueueFull) && policy != nil && policy.Enabled && claim == nil {
  _, merged, _ := s.admitQueueFullMessage(admittedCtx, identity, insertErr, sessionID, taskID, content, model, userID, planMode, attachments, metadata, claim, s.MaxPerSession(), policy)
  if merged != nil { queued = merged; return nil }
}
queued = s.finalizeAutoMerge(admittedCtx, identity, source, policy)

Session override and incarnation persistence

apps/backend/internal/orchestrator/messagequeue/types.go

Carries effective policy, source, revision, and incarnation so clients can order and validate status.

QueueStatus
type QueueStatus struct {
  Entries              []QueuedMessage `json:"entries"`
  Count                int             `json:"count"`
  Max                  int             `json:"max"`
  TaskID               string          `json:"task_id,omitempty"`
  SessionID            string          `json:"session_id,omitempty"`
  SessionIncarnationID string          `json:"session_incarnation_id,omitempty"`
  StatusEpoch          string          `json:"status_epoch,omitempty"`
  StatusGeneration     int64           `json:"status_generation,omitempty"`
  AutoRun              bool            `json:"auto_run"`
  MergeEnabled         bool            `json:"merge_enabled"`
  AutoMergeAvailable   bool            `json:"auto_merge_available"`
  AutoMergeEnabled     *bool           `json:"auto_merge_enabled,omitempty"`
  AutoMergeSource      AutoMergeSource `json:"auto_merge_source,omitempty"`
  AutoMergeRevision    *int64          `json:"auto_merge_revision,omitempty"`
}
Repository contract
GetAutoMergeOverride(ctx context.Context, identity QueueSessionIdentity) (*AutoMergeOverride, error)
SetAutoMergeOverride(ctx context.Context, identity QueueSessionIdentity, enabled bool) (AutoMergeOverride, error)
Snapshot(ctx context.Context, identity QueueSessionIdentity) (RepositorySnapshot, error)

Identity-bound queue handlers

apps/backend/internal/orchestrator/handlers/queue_handlers.go

Requires the task/session/incarnation triplet, authorizes it, and passes the unchanged identity to persistence.

Admission
func (h *QueueHandlers) admitQueuedMessage(ctx context.Context, req *wsQueueMessageRequest, queuedBy string, metadata map[string]interface{}) (*messagequeue.QueuedMessage, error) {
  if !h.requiresQueueIdentity() {
    return h.queueService.QueueMessageWithMetadata(ctx, req.SessionID, req.TaskID, req.Content, req.Model, queuedBy, req.PlanMode, req.Attachments, metadata)
  }
  identity := messagequeue.QueueSessionIdentity{TaskID: req.TaskID, SessionID: req.SessionID, SessionIncarnationID: req.SessionIncarnationID}
  if h.attachmentClaimer == nil || len(req.Attachments) == 0 {
    return admissions.QueueMessageWithMetadataForSession(ctx, identity, req.Content, req.Model, queuedBy, req.PlanMode, req.Attachments, metadata)
  }
  claim, _ := preparer.PrepareQueueAttachmentClaim(ctx, req.TaskID, queueAttachmentsToV1(req.Attachments))
  return atomicAdmissions.QueueMessageWithMetadataForSessionWithClaim(ctx, identity, req.Content, req.Model, queuedBy, req.PlanMode, req.Attachments, metadata, claim)
}
Validation
if h.requiresQueueIdentity() && req.SessionIncarnationID == "" {
  return ws.NewError(msg.ID, msg.Action, ws.ErrorCodeValidation, "task_id, session_id, and session_incarnation_id are required", nil)
}
if denied := h.authorizeQueueIdentity(ctx, msg, req.TaskID, req.SessionID, req.SessionIncarnationID); denied != nil {
  return denied, nil
}

Queue automation pills

apps/web/components/task/chat/queued-ghost-panel-header.tsx

Renders Auto-run and Auto-merge as compact pills with shared disabled state and availability gating.

Pills
function QueueAutomationPills({ autoRun, autoMerge, autoMergeAvailable, isLoading, cancellationPending, onAutoRunChange, onAutoMergeChange }: QueueAutomationPillsProps) {
  const controlsDisabled = isLoading || cancellationPending
  return (
    <div className="flex min-w-0 flex-wrap items-center gap-1.5">
      <label className="flex min-h-7 items-center gap-1.5 rounded-full border px-2 text-xs">
        <span>{t("chat:queueAutoRun")}</span>
        <Switch checked={autoRun} disabled={controlsDisabled} onCheckedChange={onAutoRunChange} />
      </label>
      <label className={autoMergeAvailable ? "cursor-pointer" : "opacity-50 pointer-events-none"}>
        <span>{t("chat:queueAutoMerge")}</span>
        <Switch checked={autoMerge} disabled={controlsDisabled || !autoMergeAvailable} onCheckedChange={onAutoMergeChange} />
      </label>
    </div>
  )
}

Frontend incarnation and ordering

apps/web/hooks/domains/session/use-queue.ts

Binds every queue mutation to the current incarnation and orders status by epoch, generation, source, and revision.

Identity
function useCurrentQueueIdentity(sessionId: string | null, taskSession: TaskSession | undefined) {
  return useMemo(() => {
    if (!sessionId || !taskSession?.task_id || !taskSession.queue_incarnation_id) return null
    return { task_id: taskSession.task_id, session_id: sessionId, session_incarnation_id: taskSession.queue_incarnation_id }
  }, [sessionId, taskSession?.task_id, taskSession?.queue_incarnation_id])
}
Mutation token
async function runQueueMutation(identity, beginQueueOperation, finishQueueOperation, refetch, mutate) {
  if (!identity) return
  const token = beginQueueOperation(sessionId, incarnationId)
  if (!token) return
  try {
    await mutate(identity)
    await refetch(sessionId, token)
  } finally {
    finishQueueOperation(sessionId, token)
  }
}

Data and storage

Per-session state lives beside the queue. Global revision and session revision order effective values without scanning sessions.

FieldTypeNotes
queue_session_state.session_idTEXT PKone row per session
queue_session_state.auto_merge_overrideINTEGER NULLNULL means inherit global, 0 is OFF, 1 is ON
queue_session_state.auto_merge_revisionINTEGERmonotonic per session, increments on each override write
task_sessions.queue_incarnation_idTEXTopaque UUID, stable for row lifetime, changes on reuse
system_settings.auto_merge_revisionINTEGERglobal revision, increments only when global value changes
QueueStatus.auto_merge_availableboolfalse when override read fails, client preserves last tuple
QueueStatus.auto_merge_sourceenum global|sessiontransport ordering only, not shown in UI
QueueStatus.status_epoch/generationstring/int64process epoch plus atomic generation orders delayed status

Risk

7 / 10 High
1 low5 medium10 high

Why this score

  • 197 files changed across backend, frontend, and docs with new persistence and WS contract.
  • Two new migrations must backfill incarnation UUIDs and override columns without breaking replay.
  • Task-then-session lock order must hold for admission, override, transfer, and deletion to avoid deadlock.
  • Stale incarnation must fail closed on every queue path, including attachment claims and deferred moves.
  • Delayed status with epoch and generation must not regress newer authoritative policy.

Trade-offs and review notes

Where to look first

  1. Verify ResolveAutoMergePolicy and admission retry on ErrAutoMergePolicyChanged in service.go.
  2. Check SQLite and memory repository lock order and incarnation validation in repository_sqlite.go and repository_memory.go.
  3. Confirm queue_handlers.go requires and authorizes the triplet on every browser queue action.
  4. Review task_sessions incarnation migration and queue_session_state override columns in base_migrations.go and base_schema.go.
  5. Check use-queue.ts and session-slice.ts incarnation and generation ordering and token CAS.
  6. Confirm task_notifications.go drops task-scoped recounts without browser broadcast.