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