mirror of
https://github.com/multica-ai/multica.git
synced 2026-08-08 11:53:43 +02:00
Layer 1 (#6383) raised the shared scanner cap to 32 MiB, which moved the cliff without removing it — a codex rollout is append-only, so a thread that outgrows any fixed cap still fails its resume forever. This adds the recovery path. Layer 2: an oversized thread/resume response is reported as Result.ResumeRejected rather than a crash, which is the positive evidence shouldRetryWithFreshSession needs and reconnects the recovery #5715 had unintentionally cut off. Reaching it required one real fix — the reader wrapped the scanner error with %v, so bufio.ErrTooLong never survived to a caller. Layer 3: a new codex_resume_oversized reason marks the session resume-unsafe, and both resume lookups block by TIME rather than by matching the failed row. That shape is required, not stylistic: the failure happens before the turn starts, so the row lands with session_id NULL and is dropped by latest_per_session before any error-text filter runs. The block expires once a thread terminates after the overflow, so an issue recovers instead of starting cold forever. Also splits the MUL-4424 continuity notice by what each surface can still read — issue comments and Slack channel history can be re-read, web chat and Feishu cannot — so only the last group tells the user a loss happened. The backend no longer holds any of that wording; it receives ExecOptions.ResumeContinuityNotice from the caller and stays silent when the prompt already carries it, which makes the duplicate injection on the retry path structurally impossible. Known gap, tracked not claimed: for chat, the claim handler reads chat_session.session_id before the fallback query, so a daemon predating this PR leaves that pointer naming the oversized thread. Clearing it needs the attempted session recorded at claim time, which is a schema change.
2718 lines
101 KiB
Go
2718 lines
101 KiB
Go
// Code generated by sqlc. DO NOT EDIT.
|
|
// versions:
|
|
// sqlc v1.31.1
|
|
// source: chat.sql
|
|
|
|
package db
|
|
|
|
import (
|
|
"context"
|
|
|
|
"github.com/jackc/pgx/v5/pgtype"
|
|
)
|
|
|
|
const advanceCancelledChatSessionPointer = `-- name: AdvanceCancelledChatSessionPointer :exec
|
|
UPDATE chat_session cs
|
|
SET session_id = t.session_id,
|
|
runtime_id = t.runtime_id,
|
|
work_dir = COALESCE(t.work_dir, cs.work_dir),
|
|
updated_at = now()
|
|
FROM agent_task_queue t
|
|
WHERE t.id = $1
|
|
AND t.chat_session_id = cs.id
|
|
AND t.status = 'cancelled'
|
|
AND t.session_id IS NOT NULL
|
|
AND t.runtime_id IS NOT NULL
|
|
AND NOT EXISTS (
|
|
SELECT 1 FROM agent_task_queue newer
|
|
WHERE newer.chat_session_id = t.chat_session_id
|
|
AND newer.id <> t.id
|
|
AND newer.session_id IS NOT NULL
|
|
AND newer.created_at > t.created_at
|
|
)
|
|
`
|
|
|
|
// Moves a chat's resume pointer onto the session a CANCELLED task recorded
|
|
// (GH #6340).
|
|
//
|
|
// Cancellation is the one terminal state that never reports back: the daemon
|
|
// discards its result and only sends a cancel-ack, so neither CompleteTask nor
|
|
// FailTask — the only other writers of chat_session.session_id — ever runs. The
|
|
// claim handler reads this pointer BEFORE falling back to
|
|
// GetLastChatTaskSession, so on a chat that already has history a pointer left
|
|
// on the previous turn shadows the cancelled turn's session no matter what the
|
|
// fallback would have found.
|
|
//
|
|
// Two callers, one statement, because both are races the other cannot cover:
|
|
// the cancel path runs it inside the status-flip transaction (so no follow-up
|
|
// can observe `cancelled` while the pointer still names the older session), and
|
|
// the pin path runs it after a mid-flight pin lands on an already-cancelled row
|
|
// (Codex waits for its rollout, so the pin routinely arrives after the cancel —
|
|
// at which point the cancel path saw no session to publish).
|
|
//
|
|
// Everything it decides on is read from the task row inside the statement, so
|
|
// neither caller can act on a stale in-memory copy. The NOT EXISTS guard is what
|
|
// makes the late pin safe: a NEWER task on this chat that already recorded a
|
|
// session owns the pointer, and a straggler must not drag the conversation
|
|
// backwards onto the turn the user interrupted.
|
|
func (q *Queries) AdvanceCancelledChatSessionPointer(ctx context.Context, taskID pgtype.UUID) error {
|
|
_, err := q.db.Exec(ctx, advanceCancelledChatSessionPointer, taskID)
|
|
return err
|
|
}
|
|
|
|
const chatSessionHasUserMessage = `-- name: ChatSessionHasUserMessage :one
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM chat_message
|
|
WHERE chat_session_id = $1 AND role = 'user'
|
|
) AS has_user_message
|
|
`
|
|
|
|
// Reports whether a session has any human (role='user') message yet. Used to
|
|
// scope the is_agent_intro self-introduction prompt to the very first,
|
|
// server-driven turn: an intro session starts with zero user messages, so the
|
|
// opening run gets the "introduce yourself" prompt. Once the creator replies,
|
|
// later turns in the same session must fall back to the normal reply prompt
|
|
// instead of repeating the introduction every turn (MUL-4259).
|
|
func (q *Queries) ChatSessionHasUserMessage(ctx context.Context, chatSessionID pgtype.UUID) (bool, error) {
|
|
row := q.db.QueryRow(ctx, chatSessionHasUserMessage, chatSessionID)
|
|
var has_user_message bool
|
|
err := row.Scan(&has_user_message)
|
|
return has_user_message, err
|
|
}
|
|
|
|
const clearChatMessageChannelMediaPending = `-- name: ClearChatMessageChannelMediaPending :exec
|
|
UPDATE chat_message
|
|
SET channel_media_pending_until = NULL
|
|
WHERE id = $1 AND chat_session_id = $2
|
|
`
|
|
|
|
type ClearChatMessageChannelMediaPendingParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
}
|
|
|
|
func (q *Queries) ClearChatMessageChannelMediaPending(ctx context.Context, arg ClearChatMessageChannelMediaPendingParams) error {
|
|
_, err := q.db.Exec(ctx, clearChatMessageChannelMediaPending, arg.ID, arg.ChatSessionID)
|
|
return err
|
|
}
|
|
|
|
const clearChatSessionProjectByProject = `-- name: ClearChatSessionProjectByProject :exec
|
|
UPDATE chat_session
|
|
SET project_id = NULL
|
|
WHERE project_id = $1 AND workspace_id = $2
|
|
`
|
|
|
|
type ClearChatSessionProjectByProjectParams struct {
|
|
ProjectID pgtype.UUID `json:"project_id"`
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
}
|
|
|
|
// Project references are intentionally soft (no database FK). Keep chat
|
|
// history while removing the context selection when a project is deleted.
|
|
// Do not touch updated_at: context cleanup is not chat activity.
|
|
func (q *Queries) ClearChatSessionProjectByProject(ctx context.Context, arg ClearChatSessionProjectByProjectParams) error {
|
|
_, err := q.db.Exec(ctx, clearChatSessionProjectByProject, arg.ProjectID, arg.WorkspaceID)
|
|
return err
|
|
}
|
|
|
|
const clearChatSessionSessionIfMatches = `-- name: ClearChatSessionSessionIfMatches :exec
|
|
UPDATE chat_session
|
|
SET session_id = NULL,
|
|
runtime_id = NULL,
|
|
updated_at = now()
|
|
WHERE id = $1
|
|
AND session_id = $2
|
|
AND runtime_id = $3
|
|
`
|
|
|
|
type ClearChatSessionSessionIfMatchesParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
SessionID pgtype.Text `json:"session_id"`
|
|
RuntimeID pgtype.UUID `json:"runtime_id"`
|
|
}
|
|
|
|
// Drops the chat session's resume pointer, but only while it still points at
|
|
// the exact session the caller proved unresumable.
|
|
//
|
|
// The claim handler reads chat_session.session_id FIRST and only falls back to
|
|
// GetLastChatTaskSession when it is empty, so a poisoned pointer here bypasses
|
|
// every filter that query applies. Declining to OVERWRITE the pointer on a
|
|
// resume-unsafe failure — which is all the fail path used to do — leaves the
|
|
// dead session in place and the next turn resumes it (GH #6066).
|
|
//
|
|
// The session_id + runtime_id predicate is what makes this safe to run in the
|
|
// fail transaction: a concurrent turn that has already written a NEW pointer
|
|
// does not match, so its healthy session survives instead of being cleared by
|
|
// a slower sibling's failure. work_dir is deliberately left alone — the
|
|
// directory is still reusable, only the conversation is not.
|
|
func (q *Queries) ClearChatSessionSessionIfMatches(ctx context.Context, arg ClearChatSessionSessionIfMatchesParams) error {
|
|
_, err := q.db.Exec(ctx, clearChatSessionSessionIfMatches, arg.ID, arg.SessionID, arg.RuntimeID)
|
|
return err
|
|
}
|
|
|
|
const createChatDraftRestore = `-- name: CreateChatDraftRestore :one
|
|
INSERT INTO chat_draft_restore (id, chat_session_id, task_id, content, attachment_ids)
|
|
VALUES ($1, $2, $3, $4, $5)
|
|
RETURNING id, chat_session_id, task_id, content, attachment_ids, created_at
|
|
`
|
|
|
|
type CreateChatDraftRestoreParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
Content string `json:"content"`
|
|
AttachmentIds []pgtype.UUID `json:"attachment_ids"`
|
|
}
|
|
|
|
// Persists the deferred-cancellation draft restore (#5219) in the same tx
|
|
// that deletes the triggering user message: the chat:cancel_finalized
|
|
// broadcast is best-effort, so an offline client recovers the draft from
|
|
// this row instead. id is the deleted message's id.
|
|
func (q *Queries) CreateChatDraftRestore(ctx context.Context, arg CreateChatDraftRestoreParams) (ChatDraftRestore, error) {
|
|
row := q.db.QueryRow(ctx, createChatDraftRestore,
|
|
arg.ID,
|
|
arg.ChatSessionID,
|
|
arg.TaskID,
|
|
arg.Content,
|
|
arg.AttachmentIds,
|
|
)
|
|
var i ChatDraftRestore
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.TaskID,
|
|
&i.Content,
|
|
&i.AttachmentIds,
|
|
&i.CreatedAt,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const createChatMessage = `-- name: CreateChatMessage :one
|
|
INSERT INTO chat_message (
|
|
chat_session_id, role, content, task_id, failure_reason, elapsed_ms,
|
|
message_kind, quick_actions, channel_media_pending_until, channel_ingested
|
|
)
|
|
VALUES (
|
|
$1, $2, $3, $4, $5, $6,
|
|
COALESCE($7::text, 'message'),
|
|
COALESCE($8::jsonb, '[]'::jsonb),
|
|
-- The media deadline is DB-clock time: every consumer compares it against
|
|
-- SQL now() (GetChannelMediaPendingUntil, the deferred promote, the
|
|
-- trailing-message guard), so the writer must use the same clock. The
|
|
-- caller passes a relative budget in seconds; an application-clock
|
|
-- timestamp here would let a skewed app node shrink or stretch the
|
|
-- fallback window.
|
|
CASE WHEN $9::float8 IS NULL THEN NULL
|
|
ELSE now() + make_interval(secs => $9::float8) END,
|
|
COALESCE($10::boolean, FALSE)
|
|
)
|
|
RETURNING id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions
|
|
`
|
|
|
|
type CreateChatMessageParams struct {
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
Role string `json:"role"`
|
|
Content string `json:"content"`
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
FailureReason pgtype.Text `json:"failure_reason"`
|
|
ElapsedMs pgtype.Int8 `json:"elapsed_ms"`
|
|
MessageKind pgtype.Text `json:"message_kind"`
|
|
QuickActions []byte `json:"quick_actions"`
|
|
ChannelMediaPendingSecs pgtype.Float8 `json:"channel_media_pending_secs"`
|
|
ChannelIngested pgtype.Bool `json:"channel_ingested"`
|
|
}
|
|
|
|
// message_kind and quick_actions default via COALESCE so every existing caller
|
|
// (which omits it) keeps writing ordinary messages; the empty-reply path passes
|
|
// 'no_response' to mark a visible turn with no text output (MUL-4351).
|
|
func (q *Queries) CreateChatMessage(ctx context.Context, arg CreateChatMessageParams) (ChatMessage, error) {
|
|
row := q.db.QueryRow(ctx, createChatMessage,
|
|
arg.ChatSessionID,
|
|
arg.Role,
|
|
arg.Content,
|
|
arg.TaskID,
|
|
arg.FailureReason,
|
|
arg.ElapsedMs,
|
|
arg.MessageKind,
|
|
arg.QuickActions,
|
|
arg.ChannelMediaPendingSecs,
|
|
arg.ChannelIngested,
|
|
)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const createChatSession = `-- name: CreateChatSession :one
|
|
INSERT INTO chat_session (workspace_id, agent_id, creator_id, title, runtime_id, is_agent_intro, project_id)
|
|
VALUES ($1, $2, $3, $4, (SELECT runtime_id FROM agent WHERE id = $2), $5, $6)
|
|
RETURNING id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id
|
|
`
|
|
|
|
type CreateChatSessionParams struct {
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
AgentID pgtype.UUID `json:"agent_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
Title string `json:"title"`
|
|
IsAgentIntro bool `json:"is_agent_intro"`
|
|
ProjectID pgtype.UUID `json:"project_id"`
|
|
}
|
|
|
|
func (q *Queries) CreateChatSession(ctx context.Context, arg CreateChatSessionParams) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, createChatSession,
|
|
arg.WorkspaceID,
|
|
arg.AgentID,
|
|
arg.CreatorID,
|
|
arg.Title,
|
|
arg.IsAgentIntro,
|
|
arg.ProjectID,
|
|
)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const createChatTask = `-- name: CreateChatTask :one
|
|
INSERT INTO agent_task_queue (
|
|
agent_id, runtime_id, issue_id, status, priority, chat_session_id,
|
|
initiator_user_id, originator_user_id, accountable_user_id, force_fresh_session, runtime_mcp_overlay,
|
|
runtime_connected_apps, originator_source, trigger_evidence_kind, trigger_evidence_ref_id,
|
|
fire_at
|
|
)
|
|
VALUES (
|
|
$1, $2, NULL,
|
|
CASE WHEN $6::timestamptz IS NULL THEN 'queued' ELSE 'deferred' END,
|
|
$3, $4, $5,
|
|
$7,
|
|
$8,
|
|
COALESCE($9::boolean, FALSE),
|
|
$10,
|
|
$11,
|
|
$12,
|
|
$13,
|
|
$14,
|
|
$6::timestamptz
|
|
)
|
|
RETURNING id, agent_id, issue_id, status, priority, dispatched_at, started_at, completed_at, result, error, created_at, context, runtime_id, session_id, work_dir, trigger_comment_id, chat_session_id, autopilot_run_id, attempt, max_attempts, parent_task_id, failure_reason, trigger_summary, force_fresh_session, is_leader_task, wait_reason, initiator_user_id, handoff_note, prepare_lease_expires_at, squad_id, runtime_mcp_overlay, escalation_for_task_id, fire_at, originator_user_id, runtime_connected_apps, coalesced_comment_ids, delivered_comment_ids, chat_input_task_id, chat_finalize_deferred_at, originator_source, delegated_from_task_id, retry_of_task_id, rerun_of_task_id, rule_version_id, trigger_evidence_kind, trigger_evidence_ref_id, accountable_user_id, session_rollout_missing, retired_session_id, quick_actions_disabled, regenerate_quick_actions_for
|
|
`
|
|
|
|
type CreateChatTaskParams struct {
|
|
AgentID pgtype.UUID `json:"agent_id"`
|
|
RuntimeID pgtype.UUID `json:"runtime_id"`
|
|
Priority int32 `json:"priority"`
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
InitiatorUserID pgtype.UUID `json:"initiator_user_id"`
|
|
FireAt pgtype.Timestamptz `json:"fire_at"`
|
|
OriginatorUserID pgtype.UUID `json:"originator_user_id"`
|
|
AccountableUserID pgtype.UUID `json:"accountable_user_id"`
|
|
ForceFreshSession pgtype.Bool `json:"force_fresh_session"`
|
|
RuntimeMcpOverlay []byte `json:"runtime_mcp_overlay"`
|
|
RuntimeConnectedApps []byte `json:"runtime_connected_apps"`
|
|
OriginatorSource pgtype.Text `json:"originator_source"`
|
|
TriggerEvidenceKind pgtype.Text `json:"trigger_evidence_kind"`
|
|
TriggerEvidenceRefID pgtype.UUID `json:"trigger_evidence_ref_id"`
|
|
}
|
|
|
|
// The chat sender (initiator) is a direct_human originator and accountable;
|
|
// attribution provenance is stamped so this path is not a NULL-source enqueue
|
|
// bypass (MUL-4302 §2).
|
|
func (q *Queries) CreateChatTask(ctx context.Context, arg CreateChatTaskParams) (AgentTaskQueue, error) {
|
|
row := q.db.QueryRow(ctx, createChatTask,
|
|
arg.AgentID,
|
|
arg.RuntimeID,
|
|
arg.Priority,
|
|
arg.ChatSessionID,
|
|
arg.InitiatorUserID,
|
|
arg.FireAt,
|
|
arg.OriginatorUserID,
|
|
arg.AccountableUserID,
|
|
arg.ForceFreshSession,
|
|
arg.RuntimeMcpOverlay,
|
|
arg.RuntimeConnectedApps,
|
|
arg.OriginatorSource,
|
|
arg.TriggerEvidenceKind,
|
|
arg.TriggerEvidenceRefID,
|
|
)
|
|
var i AgentTaskQueue
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.AgentID,
|
|
&i.IssueID,
|
|
&i.Status,
|
|
&i.Priority,
|
|
&i.DispatchedAt,
|
|
&i.StartedAt,
|
|
&i.CompletedAt,
|
|
&i.Result,
|
|
&i.Error,
|
|
&i.CreatedAt,
|
|
&i.Context,
|
|
&i.RuntimeID,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.TriggerCommentID,
|
|
&i.ChatSessionID,
|
|
&i.AutopilotRunID,
|
|
&i.Attempt,
|
|
&i.MaxAttempts,
|
|
&i.ParentTaskID,
|
|
&i.FailureReason,
|
|
&i.TriggerSummary,
|
|
&i.ForceFreshSession,
|
|
&i.IsLeaderTask,
|
|
&i.WaitReason,
|
|
&i.InitiatorUserID,
|
|
&i.HandoffNote,
|
|
&i.PrepareLeaseExpiresAt,
|
|
&i.SquadID,
|
|
&i.RuntimeMcpOverlay,
|
|
&i.EscalationForTaskID,
|
|
&i.FireAt,
|
|
&i.OriginatorUserID,
|
|
&i.RuntimeConnectedApps,
|
|
&i.CoalescedCommentIds,
|
|
&i.DeliveredCommentIds,
|
|
&i.ChatInputTaskID,
|
|
&i.ChatFinalizeDeferredAt,
|
|
&i.OriginatorSource,
|
|
&i.DelegatedFromTaskID,
|
|
&i.RetryOfTaskID,
|
|
&i.RerunOfTaskID,
|
|
&i.RuleVersionID,
|
|
&i.TriggerEvidenceKind,
|
|
&i.TriggerEvidenceRefID,
|
|
&i.AccountableUserID,
|
|
&i.SessionRolloutMissing,
|
|
&i.RetiredSessionID,
|
|
&i.QuickActionsDisabled,
|
|
&i.RegenerateQuickActionsFor,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const deferChatTaskForSealedPendingMedia = `-- name: DeferChatTaskForSealedPendingMedia :one
|
|
UPDATE agent_task_queue AS task
|
|
SET status = 'deferred', fire_at = pending.max_until
|
|
FROM (
|
|
SELECT max(message.channel_media_pending_until) AS max_until
|
|
FROM chat_message AS message
|
|
WHERE message.task_id = $1
|
|
AND message.role = 'user'
|
|
AND message.channel_media_pending_until > now()
|
|
) AS pending
|
|
WHERE task.id = $1
|
|
AND pending.max_until IS NOT NULL
|
|
AND (task.fire_at IS NULL OR task.fire_at < pending.max_until)
|
|
RETURNING task.id, task.agent_id, task.issue_id, task.status, task.priority, task.dispatched_at, task.started_at, task.completed_at, task.result, task.error, task.created_at, task.context, task.runtime_id, task.session_id, task.work_dir, task.trigger_comment_id, task.chat_session_id, task.autopilot_run_id, task.attempt, task.max_attempts, task.parent_task_id, task.failure_reason, task.trigger_summary, task.force_fresh_session, task.is_leader_task, task.wait_reason, task.initiator_user_id, task.handoff_note, task.prepare_lease_expires_at, task.squad_id, task.runtime_mcp_overlay, task.escalation_for_task_id, task.fire_at, task.originator_user_id, task.runtime_connected_apps, task.coalesced_comment_ids, task.delivered_comment_ids, task.chat_input_task_id, task.chat_finalize_deferred_at, task.originator_source, task.delegated_from_task_id, task.retry_of_task_id, task.rerun_of_task_id, task.rule_version_id, task.trigger_evidence_kind, task.trigger_evidence_ref_id, task.accountable_user_id, task.session_rollout_missing, task.retired_session_id, task.quick_actions_disabled, task.regenerate_quick_actions_for
|
|
`
|
|
|
|
// Closes the enqueue-vs-append race: under READ COMMITTED a media message can
|
|
// commit between GetChannelMediaPendingUntil and the batch seal above, landing
|
|
// an unexpired media marker inside a task the deadline read decided was
|
|
// 'queued'. Re-derive the deferral from the sealed batch itself, in the same
|
|
// transaction, so a task is never claimable while its own input still has an
|
|
// unexpired marker. No row (ErrNoRows) means no correction was needed.
|
|
func (q *Queries) DeferChatTaskForSealedPendingMedia(ctx context.Context, taskID pgtype.UUID) (AgentTaskQueue, error) {
|
|
row := q.db.QueryRow(ctx, deferChatTaskForSealedPendingMedia, taskID)
|
|
var i AgentTaskQueue
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.AgentID,
|
|
&i.IssueID,
|
|
&i.Status,
|
|
&i.Priority,
|
|
&i.DispatchedAt,
|
|
&i.StartedAt,
|
|
&i.CompletedAt,
|
|
&i.Result,
|
|
&i.Error,
|
|
&i.CreatedAt,
|
|
&i.Context,
|
|
&i.RuntimeID,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.TriggerCommentID,
|
|
&i.ChatSessionID,
|
|
&i.AutopilotRunID,
|
|
&i.Attempt,
|
|
&i.MaxAttempts,
|
|
&i.ParentTaskID,
|
|
&i.FailureReason,
|
|
&i.TriggerSummary,
|
|
&i.ForceFreshSession,
|
|
&i.IsLeaderTask,
|
|
&i.WaitReason,
|
|
&i.InitiatorUserID,
|
|
&i.HandoffNote,
|
|
&i.PrepareLeaseExpiresAt,
|
|
&i.SquadID,
|
|
&i.RuntimeMcpOverlay,
|
|
&i.EscalationForTaskID,
|
|
&i.FireAt,
|
|
&i.OriginatorUserID,
|
|
&i.RuntimeConnectedApps,
|
|
&i.CoalescedCommentIds,
|
|
&i.DeliveredCommentIds,
|
|
&i.ChatInputTaskID,
|
|
&i.ChatFinalizeDeferredAt,
|
|
&i.OriginatorSource,
|
|
&i.DelegatedFromTaskID,
|
|
&i.RetryOfTaskID,
|
|
&i.RerunOfTaskID,
|
|
&i.RuleVersionID,
|
|
&i.TriggerEvidenceKind,
|
|
&i.TriggerEvidenceRefID,
|
|
&i.AccountableUserID,
|
|
&i.SessionRolloutMissing,
|
|
&i.RetiredSessionID,
|
|
&i.QuickActionsDisabled,
|
|
&i.RegenerateQuickActionsFor,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const deleteChatDraftRestore = `-- name: DeleteChatDraftRestore :execrows
|
|
DELETE FROM chat_draft_restore
|
|
WHERE id = $1 AND chat_session_id = $2
|
|
`
|
|
|
|
type DeleteChatDraftRestoreParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
}
|
|
|
|
// Idempotent consume: deleting an already-consumed restore matches no row.
|
|
func (q *Queries) DeleteChatDraftRestore(ctx context.Context, arg DeleteChatDraftRestoreParams) (int64, error) {
|
|
result, err := q.db.Exec(ctx, deleteChatDraftRestore, arg.ID, arg.ChatSessionID)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return result.RowsAffected(), nil
|
|
}
|
|
|
|
const deleteChatDraftRestoresBySession = `-- name: DeleteChatDraftRestoresBySession :exec
|
|
DELETE FROM chat_draft_restore
|
|
WHERE chat_session_id = $1
|
|
`
|
|
|
|
// chat_draft_restore carries no chat_session FK (MUL-3515), so DeleteChatSession
|
|
// prunes its pending restores in the same tx that deletes the session.
|
|
func (q *Queries) DeleteChatDraftRestoresBySession(ctx context.Context, chatSessionID pgtype.UUID) error {
|
|
_, err := q.db.Exec(ctx, deleteChatDraftRestoresBySession, chatSessionID)
|
|
return err
|
|
}
|
|
|
|
const deleteChatDraftRestoresBySystemRuntimeAgents = `-- name: DeleteChatDraftRestoresBySystemRuntimeAgents :exec
|
|
DELETE FROM chat_draft_restore
|
|
WHERE chat_session_id IN (
|
|
SELECT cs.id FROM chat_session cs
|
|
JOIN agent a ON a.id = cs.agent_id
|
|
WHERE a.runtime_id = $1 AND a.kind = 'system'
|
|
)
|
|
`
|
|
|
|
// chat_session cascades from agent, so hard-deleting a runtime's system agents
|
|
// silently drops their sessions — and, without an FK, would strand the pending
|
|
// restores (which still hold the user's prompt text) forever. Prune them in the
|
|
// same tx, BEFORE the agent rows go: the join below needs them. Mirrors
|
|
// DeleteChannelInstallationsBySystemRuntimeAgents.
|
|
//
|
|
// Only system agents are hard-deleted on runtime teardown since MUL-5559; user
|
|
// agents (archived or not) are unbound and keep their sessions and restores.
|
|
func (q *Queries) DeleteChatDraftRestoresBySystemRuntimeAgents(ctx context.Context, runtimeID pgtype.UUID) error {
|
|
_, err := q.db.Exec(ctx, deleteChatDraftRestoresBySystemRuntimeAgents, runtimeID)
|
|
return err
|
|
}
|
|
|
|
const deleteChatSession = `-- name: DeleteChatSession :exec
|
|
DELETE FROM chat_session WHERE id = $1 AND workspace_id = $2
|
|
`
|
|
|
|
type DeleteChatSessionParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
}
|
|
|
|
// Hard delete. chat_message rows cascade via FK ON DELETE CASCADE; the
|
|
// chat_session_id on agent_task_queue is set NULL by FK so completed/failed
|
|
// task history survives the session being removed. Callers MUST run inside
|
|
// the same transaction that holds LockChatSessionForDelete and that has
|
|
// already cancelled any in-flight tasks (see CancelAgentTasksByChatSession)
|
|
// so the daemon does not keep running work whose result has nowhere to
|
|
// land. workspace_id in the WHERE clause is a SQL-layer tenant guard; see
|
|
// DeleteIssue.
|
|
func (q *Queries) DeleteChatSession(ctx context.Context, arg DeleteChatSessionParams) error {
|
|
_, err := q.db.Exec(ctx, deleteChatSession, arg.ID, arg.WorkspaceID)
|
|
return err
|
|
}
|
|
|
|
const deleteUserChatMessageByTask = `-- name: DeleteUserChatMessageByTask :one
|
|
DELETE FROM chat_message
|
|
WHERE task_id = $1 AND role = 'user'
|
|
RETURNING id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions
|
|
`
|
|
|
|
func (q *Queries) DeleteUserChatMessageByTask(ctx context.Context, taskID pgtype.UUID) (ChatMessage, error) {
|
|
row := q.db.QueryRow(ctx, deleteUserChatMessageByTask, taskID)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getChannelMediaPendingUntil = `-- name: GetChannelMediaPendingUntil :one
|
|
SELECT channel_media_pending_until
|
|
FROM chat_message
|
|
WHERE chat_session_id = $1
|
|
AND role = 'user'
|
|
AND message_kind != 'channel_command'
|
|
AND channel_media_pending_until > now()
|
|
ORDER BY channel_media_pending_until DESC
|
|
LIMIT 1
|
|
`
|
|
|
|
// The latest unexpired media deadline gates a channel task. Using a durable
|
|
// task fire_at means a process restart still produces the placeholder fallback.
|
|
// Only a turn that can join a task's input batch may gate one: channel_command
|
|
// turns are excluded from the seal below, so an unrelated later message would
|
|
// otherwise wait out a command's media (or its whole fallback budget, since a
|
|
// command whose create failed never runs the binder that clears the marker).
|
|
func (q *Queries) GetChannelMediaPendingUntil(ctx context.Context, chatSessionID pgtype.UUID) (pgtype.Timestamptz, error) {
|
|
row := q.db.QueryRow(ctx, getChannelMediaPendingUntil, chatSessionID)
|
|
var channel_media_pending_until pgtype.Timestamptz
|
|
err := row.Scan(&channel_media_pending_until)
|
|
return channel_media_pending_until, err
|
|
}
|
|
|
|
const getChatMessage = `-- name: GetChatMessage :one
|
|
SELECT id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions FROM chat_message
|
|
WHERE id = $1
|
|
`
|
|
|
|
func (q *Queries) GetChatMessage(ctx context.Context, id pgtype.UUID) (ChatMessage, error) {
|
|
row := q.db.QueryRow(ctx, getChatMessage, id)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getChatMessageByTaskAssistant = `-- name: GetChatMessageByTaskAssistant :one
|
|
SELECT id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions FROM chat_message
|
|
WHERE task_id = $1 AND role = 'assistant'
|
|
ORDER BY created_at DESC
|
|
LIMIT 1
|
|
`
|
|
|
|
// The completed turn's assistant outcome row, for the quick-actions
|
|
// supplement path (daemon suggestion pass finishing after chat:done).
|
|
func (q *Queries) GetChatMessageByTaskAssistant(ctx context.Context, taskID pgtype.UUID) (ChatMessage, error) {
|
|
row := q.db.QueryRow(ctx, getChatMessageByTaskAssistant, taskID)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getChatSession = `-- name: GetChatSession :one
|
|
SELECT id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id FROM chat_session
|
|
WHERE id = $1
|
|
`
|
|
|
|
func (q *Queries) GetChatSession(ctx context.Context, id pgtype.UUID) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, getChatSession, id)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getChatSessionInWorkspace = `-- name: GetChatSessionInWorkspace :one
|
|
SELECT id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id FROM chat_session
|
|
WHERE id = $1 AND workspace_id = $2
|
|
`
|
|
|
|
type GetChatSessionInWorkspaceParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
}
|
|
|
|
func (q *Queries) GetChatSessionInWorkspace(ctx context.Context, arg GetChatSessionInWorkspaceParams) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, getChatSessionInWorkspace, arg.ID, arg.WorkspaceID)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getLastChatTaskSession = `-- name: GetLastChatTaskSession :one
|
|
WITH retired_sessions AS (
|
|
SELECT DISTINCT r.retired_session_id AS session_id
|
|
FROM agent_task_queue r
|
|
WHERE r.chat_session_id = $1
|
|
AND r.retired_session_id IS NOT NULL
|
|
), resume_overflow_at AS (
|
|
-- completed_at alone, where the issue-side twin coalesces four columns:
|
|
-- this query already selects and orders by bare completed_at throughout,
|
|
-- so the cutoff has to be measured on the same clock as the values it is
|
|
-- compared against. Change both halves together if that ever moves.
|
|
SELECT MAX(t.completed_at) AS at
|
|
FROM agent_task_queue t
|
|
WHERE t.chat_session_id = $1
|
|
AND t.status = 'failed'
|
|
AND (
|
|
COALESCE(t.failure_reason, '') = 'codex_resume_oversized'
|
|
OR (COALESCE(t.error, '') ILIKE '%thread/resume failed%' AND COALESCE(t.error, '') ILIKE '%token too long%')
|
|
)
|
|
), latest_per_session AS (
|
|
SELECT DISTINCT ON (t.session_id)
|
|
t.session_id, t.work_dir, t.runtime_id, t.status, t.failure_reason, t.error, t.completed_at
|
|
FROM agent_task_queue t
|
|
WHERE t.chat_session_id = $1
|
|
AND t.session_id IS NOT NULL
|
|
AND t.status IN ('completed', 'failed', 'cancelled')
|
|
ORDER BY t.session_id, t.completed_at DESC
|
|
)
|
|
SELECT session_id, work_dir, runtime_id FROM latest_per_session
|
|
WHERE session_id NOT IN (SELECT session_id FROM retired_sessions)
|
|
AND (
|
|
status IN ('completed', 'cancelled')
|
|
OR (
|
|
status = 'failed'
|
|
AND COALESCE(failure_reason, '') NOT IN ('iteration_limit', 'agent_fallback_message', 'api_invalid_request', 'codex_semantic_inactivity', 'agent_error.context_overflow', 'codex_resume_oversized')
|
|
AND NOT (COALESCE(error, '') ILIKE '%400%' AND COALESCE(error, '') ILIKE '%invalid_request_error%')
|
|
AND NOT (COALESCE(error, '') ~* 'must not be empty|must be non-?empty|must have non-?empty|non-?empty content|cannot be empty|should not be empty'
|
|
AND COALESCE(error, '') ~* 'role[^a-z0-9]{0,2}assistant|assistant message|message at position|messages\.[0-9]|messages\[[0-9]')
|
|
)
|
|
)
|
|
-- MUL-5722, mirroring GetLastTaskSession: an overflowed resume records no
|
|
-- session, so exclude by time instead of by matching the failed row. Note
|
|
-- this only guards the FALLBACK — the claim handler reads
|
|
-- chat_session.session_id first, so a pointer still naming the oversized
|
|
-- thread has to be cleared at fail time (see FailTask) to be covered.
|
|
AND (
|
|
(SELECT at FROM resume_overflow_at) IS NULL
|
|
OR completed_at > (SELECT at FROM resume_overflow_at)
|
|
)
|
|
ORDER BY completed_at DESC
|
|
LIMIT 1
|
|
`
|
|
|
|
type GetLastChatTaskSessionRow struct {
|
|
SessionID pgtype.Text `json:"session_id"`
|
|
WorkDir pgtype.Text `json:"work_dir"`
|
|
RuntimeID pgtype.UUID `json:"runtime_id"`
|
|
}
|
|
|
|
// Returns the most recent task in this chat session that managed to record a
|
|
// session_id. Includes completed, failed AND cancelled tasks: each of them may
|
|
// have established a real agent session, and we'd rather resume there than
|
|
// start over and lose conversation memory. Used as a fallback when
|
|
// chat_session.session_id is NULL. Resume-unsafe failures are excluded because
|
|
// replaying those sessions deterministically reproduces the same terminal
|
|
// state. Keep this list in sync with resumeUnsafeFailureReason and
|
|
// GetLastTaskSession.
|
|
//
|
|
// 'cancelled' is resumable and its absence was GH #6340: the user stops a turn
|
|
// the agent had already started answering, and the next message starts from
|
|
// nothing. A cancelled row only carries a session_id because the daemon pinned
|
|
// one mid-flight (UpdateAgentTaskSession), which means the provider really did
|
|
// emit that session — either the resume loaded or it opened a fresh one. The
|
|
// user interrupted it; the provider did not reject it, so it is no more
|
|
// suspect than a completed one. Cancellation records no failure_reason/error,
|
|
// so the poison filters below have nothing to match and cancelled rows pass
|
|
// them the way completed rows do. The remaining risk — a transcript killed
|
|
// mid-tool-call that the provider later refuses — is caught downstream by
|
|
// taskfailure.UnresumableHistory and retires the session on the next turn.
|
|
//
|
|
// The regex pair mirrors GetLastTaskSession's provider-agnostic guard for an
|
|
// empty message baked into the conversation history: both must match, and
|
|
// both track emptyContentRe / historyMessageLocatorRe in
|
|
// pkg/taskfailure/resume.go (GH #6066).
|
|
//
|
|
// Selection is per-session, not per-row, and retired sessions are excluded —
|
|
// both mirroring GetLastTaskSession, which this query had drifted away from.
|
|
// A plain row-level filter reopens the poisoning wormhole GH #5975 closed on
|
|
// the issue side: it drops the newest poisoned row for a session and then
|
|
// happily falls back to an OLDER completed row carrying the same dead
|
|
// session_id. Judging each session by its LATEST terminal state means a newer
|
|
// poisoned row invalidates the whole session, while a genuinely different
|
|
// healthy session stays eligible.
|
|
func (q *Queries) GetLastChatTaskSession(ctx context.Context, chatSessionID pgtype.UUID) (GetLastChatTaskSessionRow, error) {
|
|
row := q.db.QueryRow(ctx, getLastChatTaskSession, chatSessionID)
|
|
var i GetLastChatTaskSessionRow
|
|
err := row.Scan(&i.SessionID, &i.WorkDir, &i.RuntimeID)
|
|
return i, err
|
|
}
|
|
|
|
const getLatestAssistantChatMessageForSession = `-- name: GetLatestAssistantChatMessageForSession :one
|
|
SELECT id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions FROM chat_message
|
|
WHERE chat_session_id = $1 AND role = 'assistant' AND task_id IS NOT NULL
|
|
ORDER BY created_at DESC
|
|
LIMIT 1
|
|
`
|
|
|
|
// The session's most recent assistant turn, used as the regeneration target
|
|
// when the user clicks "refresh" on the quick-actions row (MUL-5149). Only rows
|
|
// with a task_id qualify — the daemon suggest supplement keys off task_id and
|
|
// a resume needs a real completed turn to resume from.
|
|
func (q *Queries) GetLatestAssistantChatMessageForSession(ctx context.Context, chatSessionID pgtype.UUID) (ChatMessage, error) {
|
|
row := q.db.QueryRow(ctx, getLatestAssistantChatMessageForSession, chatSessionID)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getMostRecentUserChatMessage = `-- name: GetMostRecentUserChatMessage :one
|
|
SELECT id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions FROM chat_message
|
|
WHERE chat_session_id = $1 AND role = 'user'
|
|
ORDER BY created_at DESC
|
|
LIMIT 1
|
|
`
|
|
|
|
// Returns the most recent role='user' message in a session. Used by the
|
|
// Lark `/issue` command parser: when the user types `/issue` with no
|
|
// title, the spec falls back to "use the previous user message as the
|
|
// title". Bot replies (role='assistant') are excluded — only human
|
|
// input qualifies as a fallback title source.
|
|
func (q *Queries) GetMostRecentUserChatMessage(ctx context.Context, chatSessionID pgtype.UUID) (ChatMessage, error) {
|
|
row := q.db.QueryRow(ctx, getMostRecentUserChatMessage, chatSessionID)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getPendingChatTask = `-- name: GetPendingChatTask :one
|
|
SELECT id, status, created_at FROM agent_task_queue
|
|
WHERE chat_session_id = $1 AND status IN ('queued', 'dispatched', 'running', 'waiting_local_directory')
|
|
-- Background quick-actions regeneration passes are invisible to the chat UI:
|
|
-- they own no assistant turn and must not raise the StatusPill or disable the
|
|
-- composer (MUL-5149 refresh follow-up).
|
|
AND regenerate_quick_actions_for IS NULL
|
|
ORDER BY created_at DESC
|
|
LIMIT 1
|
|
`
|
|
|
|
type GetPendingChatTaskRow struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
Status string `json:"status"`
|
|
CreatedAt pgtype.Timestamptz `json:"created_at"`
|
|
}
|
|
|
|
// Returns the most recent in-flight task for a chat session, if any.
|
|
// Used by the frontend to recover pending state after refresh / reopen.
|
|
// created_at is the anchor for the chat StatusPill timer (it computes
|
|
// elapsed = now - task.created_at), so the pill survives refresh / reopen
|
|
// without "resetting to 0s".
|
|
func (q *Queries) GetPendingChatTask(ctx context.Context, chatSessionID pgtype.UUID) (GetPendingChatTaskRow, error) {
|
|
row := q.db.QueryRow(ctx, getPendingChatTask, chatSessionID)
|
|
var i GetPendingChatTaskRow
|
|
err := row.Scan(&i.ID, &i.Status, &i.CreatedAt)
|
|
return i, err
|
|
}
|
|
|
|
const hasActiveChatTaskForSession = `-- name: HasActiveChatTaskForSession :one
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM agent_task_queue
|
|
WHERE chat_session_id = $1
|
|
AND status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
) AS has_active
|
|
`
|
|
|
|
// True while ANY task — a normal user turn OR a background quick-actions
|
|
// regenerate — is in flight for the session (contrast GetPendingChatTask, which
|
|
// hides regenerate passes from the UI). A quick-actions refresh is refused when
|
|
// this is true: a running turn is about to change the latest reply, so the
|
|
// target we'd resume is already stale even before its assistant row lands; and a
|
|
// second concurrent regenerate would double-spend quota on the same turn. Read
|
|
// inside the same session lock as the enqueue so it cannot race a sibling insert
|
|
// (MUL-5149 review §1/§2).
|
|
//
|
|
// 'deferred' is included: an auto-retry armed with a backoff fire_at is inserted
|
|
// deferred (CreateRetryTask), and provider_network's final chat attempt waits
|
|
// ~5s that way. During that window the failed turn has written no assistant row,
|
|
// so the latest-persisted check still points at the OLD turn — omitting deferred
|
|
// would let a refresh resume a session the retry is about to advance and attach
|
|
// the new turn's suggestions to the old one (MUL-5149 re-review §1).
|
|
func (q *Queries) HasActiveChatTaskForSession(ctx context.Context, chatSessionID pgtype.UUID) (bool, error) {
|
|
row := q.db.QueryRow(ctx, hasActiveChatTaskForSession, chatSessionID)
|
|
var has_active bool
|
|
err := row.Scan(&has_active)
|
|
return has_active, err
|
|
}
|
|
|
|
const hasPendingChatTasksByCreator = `-- name: HasPendingChatTasksByCreator :one
|
|
SELECT EXISTS (
|
|
SELECT 1
|
|
FROM agent_task_queue atq
|
|
JOIN chat_session cs ON cs.id = atq.chat_session_id
|
|
WHERE atq.chat_session_id IS NOT NULL
|
|
AND atq.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
-- Background quick-actions regeneration passes own no visible turn and must
|
|
-- never light the FAB "running" indicator (MUL-5149 refresh follow-up).
|
|
AND atq.regenerate_quick_actions_for IS NULL
|
|
AND cs.workspace_id = $1
|
|
AND cs.creator_id = $2
|
|
AND cs.agent_id = ANY($3::uuid[])
|
|
) AS has_pending
|
|
`
|
|
|
|
type HasPendingChatTasksByCreatorParams struct {
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
AgentIds []pgtype.UUID `json:"agent_ids"`
|
|
}
|
|
|
|
// Boolean fast-path for the FAB's "running" indicator. Returns a single
|
|
// EXISTS row instead of the full task list, so the planner can stop at the
|
|
// first matching in-flight task (LIMIT 1 semantics via EXISTS).
|
|
//
|
|
// Permission filtering is baked into the query: agent_id = ANY($3) restricts
|
|
// the result to the agents the caller may currently see, so a member who lost
|
|
// access to a private agent never gets a true from a task they can no longer
|
|
// reach. The handler must pass its resolved accessible-agent id set as $3;
|
|
// an empty array yields false.
|
|
func (q *Queries) HasPendingChatTasksByCreator(ctx context.Context, arg HasPendingChatTasksByCreatorParams) (bool, error) {
|
|
row := q.db.QueryRow(ctx, hasPendingChatTasksByCreator, arg.WorkspaceID, arg.CreatorID, arg.AgentIds)
|
|
var has_pending bool
|
|
err := row.Scan(&has_pending)
|
|
return has_pending, err
|
|
}
|
|
|
|
const hasPendingChatTurnForSession = `-- name: HasPendingChatTurnForSession :one
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM agent_task_queue
|
|
WHERE chat_session_id = $1
|
|
AND status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
AND regenerate_quick_actions_for IS NULL
|
|
) AS has_pending
|
|
`
|
|
|
|
// Position-only check for a direct send. Unlike GetPendingChatTask this is an
|
|
// EXISTS query and includes deferred retries: a new user turn must remain a
|
|
// follow-up while an older retry waits for its backoff, otherwise promotion of
|
|
// that retry would make the new message disappear from the visible transcript.
|
|
// Background quick-action regeneration owns no visible turn and is excluded.
|
|
func (q *Queries) HasPendingChatTurnForSession(ctx context.Context, chatSessionID pgtype.UUID) (bool, error) {
|
|
row := q.db.QueryRow(ctx, hasPendingChatTurnForSession, chatSessionID)
|
|
var has_pending bool
|
|
err := row.Scan(&has_pending)
|
|
return has_pending, err
|
|
}
|
|
|
|
const linkChatMessageToTask = `-- name: LinkChatMessageToTask :exec
|
|
UPDATE chat_message
|
|
SET task_id = $2
|
|
WHERE id = $1 AND role = 'user'
|
|
`
|
|
|
|
type LinkChatMessageToTaskParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
}
|
|
|
|
func (q *Queries) LinkChatMessageToTask(ctx context.Context, arg LinkChatMessageToTaskParams) error {
|
|
_, err := q.db.Exec(ctx, linkChatMessageToTask, arg.ID, arg.TaskID)
|
|
return err
|
|
}
|
|
|
|
const linkUnownedChannelChatMessagesToTask = `-- name: LinkUnownedChannelChatMessagesToTask :exec
|
|
UPDATE chat_message AS message
|
|
SET task_id = $1
|
|
WHERE message.chat_session_id = $2
|
|
AND message.role = 'user'
|
|
AND message.task_id IS NULL
|
|
AND message.message_kind != 'channel_command'
|
|
AND NOT EXISTS (
|
|
SELECT 1
|
|
FROM chat_message AS prior
|
|
WHERE prior.chat_session_id = $2
|
|
AND prior.role != 'user'
|
|
AND (prior.created_at, prior.id) > (message.created_at, message.id)
|
|
)
|
|
`
|
|
|
|
type LinkUnownedChannelChatMessagesToTaskParams struct {
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
}
|
|
|
|
// Seals the trailing channel-message batch to its task. The task row and these
|
|
// links are committed together, so an older in-flight task cannot absorb a
|
|
// newer media message and a later assistant row cannot hide that message.
|
|
// channel_command turns were already handled synchronously by Router; keeping
|
|
// them visible but unowned prevents both immediate and delayed re-execution.
|
|
func (q *Queries) LinkUnownedChannelChatMessagesToTask(ctx context.Context, arg LinkUnownedChannelChatMessagesToTaskParams) error {
|
|
_, err := q.db.Exec(ctx, linkUnownedChannelChatMessagesToTask, arg.TaskID, arg.ChatSessionID)
|
|
return err
|
|
}
|
|
|
|
const listAgentBuilderSessionsByCreator = `-- name: ListAgentBuilderSessionsByCreator :many
|
|
SELECT cs.id,
|
|
cs.title,
|
|
cs.created_at,
|
|
cs.updated_at,
|
|
a.runtime_id,
|
|
COALESCE(lm.content, '') AS last_message_content,
|
|
COALESCE(lm.role, '') AS last_message_role,
|
|
lm.created_at AS last_message_at,
|
|
d.draft AS stored_draft
|
|
FROM chat_session cs
|
|
JOIN agent a ON a.id = cs.agent_id
|
|
LEFT JOIN agent_builder_draft d ON d.chat_session_id = cs.id
|
|
LEFT JOIN LATERAL (
|
|
SELECT content, role, created_at
|
|
FROM chat_message m
|
|
WHERE m.chat_session_id = cs.id
|
|
ORDER BY m.created_at DESC
|
|
LIMIT 1
|
|
) lm ON true
|
|
WHERE cs.workspace_id = $1
|
|
AND cs.creator_id = $2
|
|
AND cs.status = 'active'
|
|
AND a.kind = 'system'
|
|
AND a.system_key LIKE 'agent_builder:%'
|
|
AND (lm.created_at IS NOT NULL OR d.chat_session_id IS NOT NULL)
|
|
ORDER BY COALESCE(lm.created_at, d.updated_at, cs.updated_at) DESC
|
|
`
|
|
|
|
type ListAgentBuilderSessionsByCreatorParams struct {
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
}
|
|
|
|
type ListAgentBuilderSessionsByCreatorRow struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
Title string `json:"title"`
|
|
CreatedAt pgtype.Timestamptz `json:"created_at"`
|
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
|
RuntimeID pgtype.UUID `json:"runtime_id"`
|
|
LastMessageContent string `json:"last_message_content"`
|
|
LastMessageRole string `json:"last_message_role"`
|
|
LastMessageAt pgtype.Timestamptz `json:"last_message_at"`
|
|
StoredDraft []byte `json:"stored_draft"`
|
|
}
|
|
|
|
// The caller's unfinished agent-creation conversations.
|
|
//
|
|
// These never appear in ListChatSessionsByCreator: that list is filtered
|
|
// against ListAllAgents, which is `kind = 'user'` only, so a builder session —
|
|
// whose agent is the hidden `kind = 'system'` carrier — is invisible to every
|
|
// chat surface by construction. This statement is the only way back to one,
|
|
// which is why the studio may stop deleting them on navigation.
|
|
//
|
|
// `a.runtime_id` is the whole point of the join. The carrier is what
|
|
// SendDirectChatMessage reads to stamp a chat task's runtime, so it is the only
|
|
// truthful answer to "where does this conversation run". Deliberately NOT
|
|
// cs.runtime_id: that is the daemon's resume pointer, left stale on purpose
|
|
// after a runtime switch (see RebindAgentBuilderRuntime), so resuming from it
|
|
// would put the picker on a runtime that no longer executes anything — the
|
|
// exact split MUL-5163 removed.
|
|
//
|
|
// A conversation qualifies once it holds something the user would miss: a
|
|
// message, or a saved configuration. Requiring a message alone was wrong — the
|
|
// form on the right is editable from the moment the session exists and
|
|
// autosaves, so someone can open the builder, type a name, and leave before the
|
|
// first turn. That session has real work in it and was unreachable. Requiring
|
|
// neither is also wrong: a session opened and abandoned untouched is not a
|
|
// draft, and would put an empty row in front of the user on every accidental
|
|
// entry into the flow.
|
|
//
|
|
// The stored draft rides along instead of needing its own fetch: the studio
|
|
// renders this list beside the conversation it is switching between, so the
|
|
// configuration for the row the user picks has to be in hand at click time.
|
|
// LEFT JOIN because a conversation that has only ever been driven by the AI has
|
|
// no saved draft — the client replays the last <agent_draft> block in that case.
|
|
// A draft-only session has no message to sort by; fall back to when its
|
|
// configuration was last written so it still lands in activity order.
|
|
func (q *Queries) ListAgentBuilderSessionsByCreator(ctx context.Context, arg ListAgentBuilderSessionsByCreatorParams) ([]ListAgentBuilderSessionsByCreatorRow, error) {
|
|
rows, err := q.db.Query(ctx, listAgentBuilderSessionsByCreator, arg.WorkspaceID, arg.CreatorID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ListAgentBuilderSessionsByCreatorRow{}
|
|
for rows.Next() {
|
|
var i ListAgentBuilderSessionsByCreatorRow
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.Title,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.RuntimeID,
|
|
&i.LastMessageContent,
|
|
&i.LastMessageRole,
|
|
&i.LastMessageAt,
|
|
&i.StoredDraft,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listAllChatSessionsByCreator = `-- name: ListAllChatSessionsByCreator :many
|
|
SELECT cs.id, cs.workspace_id, cs.agent_id, cs.creator_id, cs.title, cs.session_id, cs.work_dir, cs.status, cs.created_at, cs.updated_at, cs.unread_since, cs.runtime_id, cs.last_read_at, cs.is_agent_intro, cs.pinned_at, cs.project_id,
|
|
CASE WHEN cs.status = 'archived' THEN 0
|
|
ELSE (SELECT count(*) FROM chat_message m
|
|
WHERE m.chat_session_id = cs.id
|
|
AND m.role = 'assistant'
|
|
AND m.created_at > cs.last_read_at)
|
|
END::int AS unread_count,
|
|
COALESCE(lm.content, '') AS last_message_content,
|
|
COALESCE(lm.role, '') AS last_message_role,
|
|
lm.created_at AS last_message_at,
|
|
lm.failure_reason AS last_message_failure_reason,
|
|
COALESCE(lm.message_kind, '') AS last_message_kind
|
|
FROM chat_session cs
|
|
LEFT JOIN LATERAL (
|
|
SELECT content, role, created_at, failure_reason, message_kind
|
|
FROM chat_message m
|
|
WHERE m.chat_session_id = cs.id
|
|
ORDER BY m.created_at DESC
|
|
LIMIT 1
|
|
) lm ON true
|
|
WHERE cs.workspace_id = $1 AND cs.creator_id = $2
|
|
ORDER BY (cs.pinned_at IS NOT NULL) DESC, cs.pinned_at DESC, COALESCE(lm.created_at, cs.updated_at) DESC
|
|
`
|
|
|
|
type ListAllChatSessionsByCreatorParams struct {
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
}
|
|
|
|
type ListAllChatSessionsByCreatorRow struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
AgentID pgtype.UUID `json:"agent_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
Title string `json:"title"`
|
|
SessionID pgtype.Text `json:"session_id"`
|
|
WorkDir pgtype.Text `json:"work_dir"`
|
|
Status string `json:"status"`
|
|
CreatedAt pgtype.Timestamptz `json:"created_at"`
|
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
|
UnreadSince pgtype.Timestamptz `json:"unread_since"`
|
|
RuntimeID pgtype.UUID `json:"runtime_id"`
|
|
LastReadAt pgtype.Timestamptz `json:"last_read_at"`
|
|
IsAgentIntro bool `json:"is_agent_intro"`
|
|
PinnedAt pgtype.Timestamptz `json:"pinned_at"`
|
|
ProjectID pgtype.UUID `json:"project_id"`
|
|
UnreadCount int32 `json:"unread_count"`
|
|
LastMessageContent string `json:"last_message_content"`
|
|
LastMessageRole string `json:"last_message_role"`
|
|
LastMessageAt pgtype.Timestamptz `json:"last_message_at"`
|
|
LastMessageFailureReason pgtype.Text `json:"last_message_failure_reason"`
|
|
LastMessageKind string `json:"last_message_kind"`
|
|
}
|
|
|
|
// Unlike ListChatSessionsByCreator this returns archived sessions too (for the
|
|
// "Archived" view), so unread must be forced to 0 for archived rows: archiving
|
|
// deliberately does NOT advance last_read_at (so unarchive can restore the true
|
|
// unread state), but an archived session is read-only and hidden from history,
|
|
// so any residual unread is uncleanable and must not light up any badge. Gating
|
|
// on status here is the single source of truth for all unread surfaces (FAB,
|
|
// sidebar Chat tab, chat-window header) — see MUL-4360.
|
|
func (q *Queries) ListAllChatSessionsByCreator(ctx context.Context, arg ListAllChatSessionsByCreatorParams) ([]ListAllChatSessionsByCreatorRow, error) {
|
|
rows, err := q.db.Query(ctx, listAllChatSessionsByCreator, arg.WorkspaceID, arg.CreatorID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ListAllChatSessionsByCreatorRow{}
|
|
for rows.Next() {
|
|
var i ListAllChatSessionsByCreatorRow
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
&i.UnreadCount,
|
|
&i.LastMessageContent,
|
|
&i.LastMessageRole,
|
|
&i.LastMessageAt,
|
|
&i.LastMessageFailureReason,
|
|
&i.LastMessageKind,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listChatDraftRestoresBySession = `-- name: ListChatDraftRestoresBySession :many
|
|
SELECT id, chat_session_id, task_id, content, attachment_ids, created_at FROM chat_draft_restore
|
|
WHERE chat_session_id = $1
|
|
ORDER BY created_at ASC
|
|
`
|
|
|
|
func (q *Queries) ListChatDraftRestoresBySession(ctx context.Context, chatSessionID pgtype.UUID) ([]ChatDraftRestore, error) {
|
|
rows, err := q.db.Query(ctx, listChatDraftRestoresBySession, chatSessionID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ChatDraftRestore{}
|
|
for rows.Next() {
|
|
var i ChatDraftRestore
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.TaskID,
|
|
&i.Content,
|
|
&i.AttachmentIds,
|
|
&i.CreatedAt,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listChatInputMessages = `-- name: ListChatInputMessages :many
|
|
SELECT id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions FROM chat_message
|
|
WHERE task_id = $1 AND role = 'user'
|
|
ORDER BY created_at ASC, id ASC
|
|
`
|
|
|
|
// Loads the immutable user-message input batch owned by a direct-chat task.
|
|
// The caller passes the task's chat_input_task_id (itself for an original send,
|
|
// the root task for an auto-retry child), so a claim reads exactly the messages
|
|
// the user sent for this turn — and never absorbs a message that arrived after
|
|
// the batch was sealed, no matter what the assistant wrote or when. Only used
|
|
// for new task-owned direct-chat tasks; legacy/channel (chat_input_task_id
|
|
// NULL) tasks keep using ListChatMessagesForLegacyTask + trailingUserMessages.
|
|
func (q *Queries) ListChatInputMessages(ctx context.Context, taskID pgtype.UUID) ([]ChatMessage, error) {
|
|
rows, err := q.db.Query(ctx, listChatInputMessages, taskID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ChatMessage{}
|
|
for rows.Next() {
|
|
var i ChatMessage
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listChatMessages = `-- name: ListChatMessages :many
|
|
SELECT message.id, message.chat_session_id, message.role, message.content, message.task_id, message.created_at, message.failure_reason, message.elapsed_ms, message.message_kind, message.channel_media_pending_until, message.channel_ingested, message.quick_actions FROM chat_message AS message
|
|
WHERE message.chat_session_id = $1
|
|
AND NOT (
|
|
message.role = 'user'
|
|
AND EXISTS (
|
|
SELECT 1
|
|
FROM agent_task_queue AS task
|
|
WHERE task.chat_session_id = message.chat_session_id
|
|
AND task.status = 'queued'
|
|
AND task.id = message.task_id
|
|
-- "Queued follow-up" is positional, not the row's transient status:
|
|
-- the first pending task is the current turn even before claim.
|
|
AND task.id <> (
|
|
SELECT head.id
|
|
FROM agent_task_queue AS head
|
|
WHERE head.chat_session_id = $1
|
|
AND head.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
AND head.regenerate_quick_actions_for IS NULL
|
|
ORDER BY
|
|
CASE
|
|
WHEN head.status IN ('dispatched', 'running', 'waiting_local_directory') THEN 0
|
|
WHEN head.status = 'deferred' THEN 1
|
|
ELSE 2
|
|
END,
|
|
head.priority DESC,
|
|
head.created_at ASC,
|
|
head.id ASC
|
|
LIMIT 1
|
|
)
|
|
)
|
|
)
|
|
ORDER BY message.created_at ASC, message.id ASC
|
|
`
|
|
|
|
// IMPORTANT: the visible-head selector below is also used by
|
|
// ListChatMessagesForLegacyTask, ListChatMessagesPage,
|
|
// ReanchorClaimedDirectChatInput, ReanchorNextQueuedDirectChatInput,
|
|
// ListPendingChatTasksForSession, and CancelQueuedAgentTasksForSession in
|
|
// agent.sql. Keep the eligible statuses and ordering identical: a claimed task
|
|
// is current, a deferred retry precedes still-queued work, and queued peers use
|
|
// claim priority/FIFO order. Background quick-action regeneration is invisible.
|
|
// This is a presentation/visibility order, not a scheduling guarantee:
|
|
// deferred rows are not claimable before promotion, so a queued row may be
|
|
// claimed during the backoff and then becomes the visible claimed head.
|
|
func (q *Queries) ListChatMessages(ctx context.Context, chatSessionID pgtype.UUID) ([]ChatMessage, error) {
|
|
rows, err := q.db.Query(ctx, listChatMessages, chatSessionID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ChatMessage{}
|
|
for rows.Next() {
|
|
var i ChatMessage
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listChatMessagesForLegacyTask = `-- name: ListChatMessagesForLegacyTask :many
|
|
SELECT message.id, message.chat_session_id, message.role, message.content, message.task_id, message.created_at, message.failure_reason, message.elapsed_ms, message.message_kind, message.channel_media_pending_until, message.channel_ingested, message.quick_actions FROM chat_message AS message
|
|
WHERE message.chat_session_id = $1
|
|
AND NOT (
|
|
message.role = 'user'
|
|
AND EXISTS (
|
|
SELECT 1
|
|
FROM agent_task_queue AS task
|
|
WHERE task.chat_session_id = message.chat_session_id
|
|
AND task.status = 'queued'
|
|
AND task.id = message.task_id
|
|
AND task.id <> (
|
|
SELECT head.id
|
|
FROM agent_task_queue AS head
|
|
WHERE head.chat_session_id = $1
|
|
AND head.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
AND head.regenerate_quick_actions_for IS NULL
|
|
ORDER BY
|
|
CASE
|
|
WHEN head.status IN ('dispatched', 'running', 'waiting_local_directory') THEN 0
|
|
WHEN head.status = 'deferred' THEN 1
|
|
ELSE 2
|
|
END,
|
|
head.priority DESC,
|
|
head.created_at ASC,
|
|
head.id ASC
|
|
LIMIT 1
|
|
)
|
|
)
|
|
)
|
|
ORDER BY message.created_at ASC, message.id ASC
|
|
`
|
|
|
|
// Legacy/reclaimed daemon tasks use trailing history, but must not absorb a
|
|
// newer queued successor that is already bound to its own user message.
|
|
func (q *Queries) ListChatMessagesForLegacyTask(ctx context.Context, chatSessionID pgtype.UUID) ([]ChatMessage, error) {
|
|
rows, err := q.db.Query(ctx, listChatMessagesForLegacyTask, chatSessionID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ChatMessage{}
|
|
for rows.Next() {
|
|
var i ChatMessage
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listChatMessagesPage = `-- name: ListChatMessagesPage :many
|
|
SELECT message.id, message.chat_session_id, message.role, message.content, message.task_id, message.created_at, message.failure_reason, message.elapsed_ms, message.message_kind, message.channel_media_pending_until, message.channel_ingested, message.quick_actions FROM chat_message AS message
|
|
WHERE message.chat_session_id = $1
|
|
AND NOT (
|
|
message.role = 'user'
|
|
AND EXISTS (
|
|
SELECT 1
|
|
FROM agent_task_queue AS task
|
|
WHERE task.chat_session_id = message.chat_session_id
|
|
AND task.status = 'queued'
|
|
AND task.id = message.task_id
|
|
AND task.id <> (
|
|
SELECT head.id
|
|
FROM agent_task_queue AS head
|
|
WHERE head.chat_session_id = $1
|
|
AND head.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
AND head.regenerate_quick_actions_for IS NULL
|
|
ORDER BY
|
|
CASE
|
|
WHEN head.status IN ('dispatched', 'running', 'waiting_local_directory') THEN 0
|
|
WHEN head.status = 'deferred' THEN 1
|
|
ELSE 2
|
|
END,
|
|
head.priority DESC,
|
|
head.created_at ASC,
|
|
head.id ASC
|
|
LIMIT 1
|
|
)
|
|
)
|
|
)
|
|
AND (
|
|
$3::timestamptz IS NULL
|
|
OR (message.created_at, message.id) < ($3::timestamptz, $4::uuid)
|
|
)
|
|
ORDER BY message.created_at DESC, message.id DESC
|
|
LIMIT $2
|
|
`
|
|
|
|
type ListChatMessagesPageParams struct {
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
Limit int32 `json:"limit"`
|
|
BeforeCreatedAt pgtype.Timestamptz `json:"before_created_at"`
|
|
BeforeID pgtype.UUID `json:"before_id"`
|
|
}
|
|
|
|
func (q *Queries) ListChatMessagesPage(ctx context.Context, arg ListChatMessagesPageParams) ([]ChatMessage, error) {
|
|
rows, err := q.db.Query(ctx, listChatMessagesPage,
|
|
arg.ChatSessionID,
|
|
arg.Limit,
|
|
arg.BeforeCreatedAt,
|
|
arg.BeforeID,
|
|
)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ChatMessage{}
|
|
for rows.Next() {
|
|
var i ChatMessage
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listChatSessionsByCreator = `-- name: ListChatSessionsByCreator :many
|
|
SELECT cs.id, cs.workspace_id, cs.agent_id, cs.creator_id, cs.title, cs.session_id, cs.work_dir, cs.status, cs.created_at, cs.updated_at, cs.unread_since, cs.runtime_id, cs.last_read_at, cs.is_agent_intro, cs.pinned_at, cs.project_id,
|
|
(SELECT count(*) FROM chat_message m
|
|
WHERE m.chat_session_id = cs.id
|
|
AND m.role = 'assistant'
|
|
AND m.created_at > cs.last_read_at)::int AS unread_count,
|
|
COALESCE(lm.content, '') AS last_message_content,
|
|
COALESCE(lm.role, '') AS last_message_role,
|
|
lm.created_at AS last_message_at,
|
|
lm.failure_reason AS last_message_failure_reason,
|
|
COALESCE(lm.message_kind, '') AS last_message_kind
|
|
FROM chat_session cs
|
|
LEFT JOIN LATERAL (
|
|
SELECT content, role, created_at, failure_reason, message_kind
|
|
FROM chat_message m
|
|
WHERE m.chat_session_id = cs.id
|
|
ORDER BY m.created_at DESC
|
|
LIMIT 1
|
|
) lm ON true
|
|
WHERE cs.workspace_id = $1 AND cs.creator_id = $2 AND cs.status = 'active'
|
|
ORDER BY (cs.pinned_at IS NOT NULL) DESC, cs.pinned_at DESC, COALESCE(lm.created_at, cs.updated_at) DESC
|
|
`
|
|
|
|
type ListChatSessionsByCreatorParams struct {
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
}
|
|
|
|
type ListChatSessionsByCreatorRow struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
AgentID pgtype.UUID `json:"agent_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
Title string `json:"title"`
|
|
SessionID pgtype.Text `json:"session_id"`
|
|
WorkDir pgtype.Text `json:"work_dir"`
|
|
Status string `json:"status"`
|
|
CreatedAt pgtype.Timestamptz `json:"created_at"`
|
|
UpdatedAt pgtype.Timestamptz `json:"updated_at"`
|
|
UnreadSince pgtype.Timestamptz `json:"unread_since"`
|
|
RuntimeID pgtype.UUID `json:"runtime_id"`
|
|
LastReadAt pgtype.Timestamptz `json:"last_read_at"`
|
|
IsAgentIntro bool `json:"is_agent_intro"`
|
|
PinnedAt pgtype.Timestamptz `json:"pinned_at"`
|
|
ProjectID pgtype.UUID `json:"project_id"`
|
|
UnreadCount int32 `json:"unread_count"`
|
|
LastMessageContent string `json:"last_message_content"`
|
|
LastMessageRole string `json:"last_message_role"`
|
|
LastMessageAt pgtype.Timestamptz `json:"last_message_at"`
|
|
LastMessageFailureReason pgtype.Text `json:"last_message_failure_reason"`
|
|
LastMessageKind string `json:"last_message_kind"`
|
|
}
|
|
|
|
// IM-style list: each active session with its unread *count* (assistant
|
|
// messages after the read cursor), a preview of the latest message, and
|
|
// ordered by most-recent activity so a new reply bumps a session to the top.
|
|
func (q *Queries) ListChatSessionsByCreator(ctx context.Context, arg ListChatSessionsByCreatorParams) ([]ListChatSessionsByCreatorRow, error) {
|
|
rows, err := q.db.Query(ctx, listChatSessionsByCreator, arg.WorkspaceID, arg.CreatorID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ListChatSessionsByCreatorRow{}
|
|
for rows.Next() {
|
|
var i ListChatSessionsByCreatorRow
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
&i.UnreadCount,
|
|
&i.LastMessageContent,
|
|
&i.LastMessageRole,
|
|
&i.LastMessageAt,
|
|
&i.LastMessageFailureReason,
|
|
&i.LastMessageKind,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listPendingChatTasksByCreator = `-- name: ListPendingChatTasksByCreator :many
|
|
SELECT atq.id AS task_id, atq.status, atq.chat_session_id, cs.agent_id
|
|
FROM agent_task_queue atq
|
|
JOIN chat_session cs ON cs.id = atq.chat_session_id
|
|
WHERE atq.chat_session_id IS NOT NULL
|
|
AND atq.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
-- Exclude background quick-actions regeneration passes: they own no assistant
|
|
-- turn and must not surface as "running" chat work (MUL-5149 refresh follow-up).
|
|
AND atq.regenerate_quick_actions_for IS NULL
|
|
AND cs.workspace_id = $1
|
|
AND cs.creator_id = $2
|
|
ORDER BY atq.created_at DESC
|
|
`
|
|
|
|
type ListPendingChatTasksByCreatorParams struct {
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
CreatorID pgtype.UUID `json:"creator_id"`
|
|
}
|
|
|
|
type ListPendingChatTasksByCreatorRow struct {
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
Status string `json:"status"`
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
AgentID pgtype.UUID `json:"agent_id"`
|
|
}
|
|
|
|
// Aggregate view of all in-flight chat tasks owned by a given creator in a
|
|
// workspace. Drives the FAB's "running" indicator when the chat window is
|
|
// closed and no single session's query is active.
|
|
//
|
|
// Returns cs.agent_id so the handler can filter tasks belonging to private
|
|
// agents the caller has lost access to using the already-loaded `allowed`
|
|
// set — no second ListAllChatSessionsByCreator scan on the hot path.
|
|
//
|
|
// atq.chat_session_id IS NOT NULL is redundant given the JOIN, but stated
|
|
// explicitly so the planner can prove the query predicate is a subset of the
|
|
// idx_agent_task_queue_chat_pending_v3 partial-index predicate and use it.
|
|
func (q *Queries) ListPendingChatTasksByCreator(ctx context.Context, arg ListPendingChatTasksByCreatorParams) ([]ListPendingChatTasksByCreatorRow, error) {
|
|
rows, err := q.db.Query(ctx, listPendingChatTasksByCreator, arg.WorkspaceID, arg.CreatorID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ListPendingChatTasksByCreatorRow{}
|
|
for rows.Next() {
|
|
var i ListPendingChatTasksByCreatorRow
|
|
if err := rows.Scan(
|
|
&i.TaskID,
|
|
&i.Status,
|
|
&i.ChatSessionID,
|
|
&i.AgentID,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const listPendingChatTasksForSession = `-- name: ListPendingChatTasksForSession :many
|
|
SELECT
|
|
task.id,
|
|
task.status,
|
|
task.created_at,
|
|
message.id AS message_id,
|
|
COALESCE(message.content, '')::text AS content
|
|
FROM agent_task_queue AS task
|
|
LEFT JOIN LATERAL (
|
|
SELECT input.id, input.content
|
|
FROM chat_message AS input
|
|
WHERE input.task_id = COALESCE(task.chat_input_task_id, task.id)
|
|
AND input.role = 'user'
|
|
ORDER BY input.created_at ASC, input.id ASC
|
|
LIMIT 1
|
|
) AS message ON TRUE
|
|
WHERE task.chat_session_id = $1
|
|
AND task.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
AND task.regenerate_quick_actions_for IS NULL
|
|
ORDER BY
|
|
CASE
|
|
WHEN task.status IN ('dispatched', 'running', 'waiting_local_directory') THEN 0
|
|
WHEN task.status = 'deferred' THEN 1
|
|
ELSE 2
|
|
END,
|
|
task.priority DESC,
|
|
task.created_at ASC,
|
|
task.id ASC
|
|
`
|
|
|
|
type ListPendingChatTasksForSessionRow struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
Status string `json:"status"`
|
|
CreatedAt pgtype.Timestamptz `json:"created_at"`
|
|
MessageID pgtype.UUID `json:"message_id"`
|
|
Content string `json:"content"`
|
|
}
|
|
|
|
// Returns a claimed task first, then a deferred retry, followed by prioritized
|
|
// then FIFO queued work. See the shared visible-head invariant above
|
|
// ListChatMessages; changing this order requires changing every selector named
|
|
// there in the same patch.
|
|
// The message lateral join reads only the immutable input owned by each task;
|
|
// it avoids loading the session's complete message history just to render a
|
|
// one-line queue preview. GetPendingChatTask remains for legacy callers that
|
|
// only need an existence check.
|
|
func (q *Queries) ListPendingChatTasksForSession(ctx context.Context, chatSessionID pgtype.UUID) ([]ListPendingChatTasksForSessionRow, error) {
|
|
rows, err := q.db.Query(ctx, listPendingChatTasksForSession, chatSessionID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []ListPendingChatTasksForSessionRow{}
|
|
for rows.Next() {
|
|
var i ListPendingChatTasksForSessionRow
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.MessageID,
|
|
&i.Content,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const lockChatSessionForDelete = `-- name: LockChatSessionForDelete :one
|
|
SELECT id FROM chat_session
|
|
WHERE id = $1
|
|
FOR UPDATE
|
|
`
|
|
|
|
// Acquires an exclusive (FOR UPDATE) row lock on chat_session(id). Used by
|
|
// the delete path so that a concurrent SendChatMessage cannot enqueue a new
|
|
// agent_task_queue row referencing this session between our cancel and
|
|
// delete steps. The FK from agent_task_queue.chat_session_id takes a
|
|
// KEY SHARE lock on the parent row during INSERT validation, which
|
|
// conflicts with FOR UPDATE — concurrent inserts block here and then fail
|
|
// their FK check after we commit the delete.
|
|
func (q *Queries) LockChatSessionForDelete(ctx context.Context, id pgtype.UUID) (pgtype.UUID, error) {
|
|
row := q.db.QueryRow(ctx, lockChatSessionForDelete, id)
|
|
var id_2 pgtype.UUID
|
|
err := row.Scan(&id_2)
|
|
return id_2, err
|
|
}
|
|
|
|
const lockChatSessionForDraftWrite = `-- name: LockChatSessionForDraftWrite :one
|
|
SELECT id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id FROM chat_session
|
|
WHERE id = $1
|
|
FOR UPDATE
|
|
`
|
|
|
|
// The autosave half of the agent_builder_draft protocol, and the writer's
|
|
// answer to LockChatSessionForDelete.
|
|
//
|
|
// agent_builder_draft carries no chat_session FK (repo rule), so an INSERT into
|
|
// it takes no lock on the parent row and nothing stops a draft from landing
|
|
// after its conversation is gone: the save path read the session, the delete
|
|
// transaction then committed, and the upsert still succeeded — leaving a row
|
|
// the user explicitly discarded, invisible to the UI and unreachable by every
|
|
// prune except the workspace teardown. The FK that would have rejected that
|
|
// INSERT is the one we are not allowed to have, so the lock replaces it.
|
|
//
|
|
// Returns the whole row, not just the id: the caller must re-check creator and
|
|
// status INSIDE the transaction, because a save blocked here resumes holding
|
|
// values it read before blocking (the same reason the runtime-bind path re-reads
|
|
// the agent under its lock).
|
|
//
|
|
// Same row and same lock mode as LockChatSessionForDelete and
|
|
// LockChatSessionForRuntimeBind, taken as the transaction's first statement, so
|
|
// no ordering against the repo-wide chat_session -> agent_task_queue sequence is
|
|
// introduced and none of the three can deadlock against each other.
|
|
func (q *Queries) LockChatSessionForDraftWrite(ctx context.Context, id pgtype.UUID) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, lockChatSessionForDraftWrite, id)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const lockChatSessionForRuntimeBind = `-- name: LockChatSessionForRuntimeBind :one
|
|
SELECT id FROM chat_session
|
|
WHERE id = $1
|
|
FOR UPDATE
|
|
`
|
|
|
|
// Acquires an exclusive (FOR UPDATE) row lock on chat_session(id), serialising
|
|
// "which runtime does this session execute on" against "enqueue the next task".
|
|
//
|
|
// Both SendDirectChatMessage and the agent-builder runtime switch take this lock
|
|
// for their whole transaction. Without it the two are a read-then-write race: a
|
|
// send reads the carrier agent's runtime_id, the switch then passes its
|
|
// pending-task check and rebinds the carrier, and the send finally inserts a task
|
|
// still stamped with the pre-switch runtime — so the user is told the switch
|
|
// succeeded while their message runs on the old runtime (MUL-5163).
|
|
//
|
|
// The lock alone is not sufficient: the send path must also re-read the agent
|
|
// INSIDE the locked transaction, because a send blocked at INSERT would otherwise
|
|
// resume and write the runtime_id it read before blocking.
|
|
//
|
|
// Same row and same lock mode as LockChatSessionForDelete, and both take it as
|
|
// their first statement, so the delete path and this one cannot deadlock.
|
|
func (q *Queries) LockChatSessionForRuntimeBind(ctx context.Context, id pgtype.UUID) (pgtype.UUID, error) {
|
|
row := q.db.QueryRow(ctx, lockChatSessionForRuntimeBind, id)
|
|
var id_2 pgtype.UUID
|
|
err := row.Scan(&id_2)
|
|
return id_2, err
|
|
}
|
|
|
|
const lockChatSessionForTask = `-- name: LockChatSessionForTask :one
|
|
|
|
SELECT cs.id
|
|
FROM agent_task_queue t
|
|
JOIN chat_session cs ON cs.id = t.chat_session_id
|
|
WHERE t.id = $1
|
|
FOR UPDATE OF cs
|
|
`
|
|
|
|
// The chat_session row lock is the mutual-exclusion protocol between the
|
|
// draft-restore writer (FinalizeDeferredCancelledChat) and every deleter of a
|
|
// chat_session or one of its cascade parents. Without an FK, an INSERT into
|
|
// chat_draft_restore takes no lock on its session, so a prune-then-delete would
|
|
// otherwise miss a restore committed after its snapshot and strand it — with the
|
|
// user's prompt in it — forever (#5219).
|
|
//
|
|
// The contract, held by all five paths below:
|
|
//
|
|
// deleter: lock the sessions FOR UPDATE -> prune restores -> delete the parent
|
|
// finalizer: lock the session FOR UPDATE -> insert the restore
|
|
//
|
|
// Whoever locks first wins: a finalizer that got there first commits its row
|
|
// before the deleter's prune statement takes its snapshot, so the prune sweeps
|
|
// it; a deleter that got there first leaves no session for the finalizer to
|
|
// lock, so it never inserts.
|
|
//
|
|
// Lock order is chat_session -> agent_task_queue everywhere (the finalizer locks
|
|
// the session before claiming its task) so the deleters' cascade into
|
|
// agent_task_queue cannot deadlock against it.
|
|
//
|
|
// The single-session delete path needs no new query: it already holds
|
|
// LockChatSessionForDelete (same FOR UPDATE row lock) across its prune.
|
|
// The finalizer's half of the protocol. No rows means the session is already
|
|
// gone (its cascade NULLs agent_task_queue.chat_session_id), so there is nothing
|
|
// to lock and nothing to restore into.
|
|
func (q *Queries) LockChatSessionForTask(ctx context.Context, id pgtype.UUID) (pgtype.UUID, error) {
|
|
row := q.db.QueryRow(ctx, lockChatSessionForTask, id)
|
|
var id_2 pgtype.UUID
|
|
err := row.Scan(&id_2)
|
|
return id_2, err
|
|
}
|
|
|
|
const lockChatSessionsBySystemRuntimeAgents = `-- name: LockChatSessionsBySystemRuntimeAgents :many
|
|
SELECT cs.id FROM chat_session cs
|
|
JOIN agent a ON a.id = cs.agent_id
|
|
WHERE a.runtime_id = $1 AND a.kind = 'system'
|
|
ORDER BY cs.id
|
|
FOR UPDATE OF cs
|
|
`
|
|
|
|
func (q *Queries) LockChatSessionsBySystemRuntimeAgents(ctx context.Context, runtimeID pgtype.UUID) ([]pgtype.UUID, error) {
|
|
rows, err := q.db.Query(ctx, lockChatSessionsBySystemRuntimeAgents, runtimeID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []pgtype.UUID{}
|
|
for rows.Next() {
|
|
var id pgtype.UUID
|
|
if err := rows.Scan(&id); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, id)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const lockChatSessionsByWorkspace = `-- name: LockChatSessionsByWorkspace :many
|
|
SELECT id FROM chat_session
|
|
WHERE workspace_id = $1
|
|
ORDER BY id
|
|
FOR UPDATE
|
|
`
|
|
|
|
// ORDER BY id: a stable lock order keeps two concurrent deleters from
|
|
// deadlocking against each other.
|
|
func (q *Queries) LockChatSessionsByWorkspace(ctx context.Context, workspaceID pgtype.UUID) ([]pgtype.UUID, error) {
|
|
rows, err := q.db.Query(ctx, lockChatSessionsByWorkspace, workspaceID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []pgtype.UUID{}
|
|
for rows.Next() {
|
|
var id pgtype.UUID
|
|
if err := rows.Scan(&id); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, id)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const markChatSessionRead = `-- name: MarkChatSessionRead :exec
|
|
UPDATE chat_session SET last_read_at = now()
|
|
WHERE id = $1
|
|
`
|
|
|
|
// Advances the read cursor to now, dropping the session's unread_count to 0.
|
|
func (q *Queries) MarkChatSessionRead(ctx context.Context, id pgtype.UUID) error {
|
|
_, err := q.db.Exec(ctx, markChatSessionRead, id)
|
|
return err
|
|
}
|
|
|
|
const prioritizeQueuedChatTask = `-- name: PrioritizeQueuedChatTask :one
|
|
WITH target AS MATERIALIZED (
|
|
SELECT candidate.id
|
|
FROM agent_task_queue AS candidate
|
|
WHERE candidate.id = $2
|
|
AND candidate.chat_session_id = $1
|
|
AND candidate.status = 'queued'
|
|
-- "Send now" is valid only while there is a visible claimed task for the
|
|
-- client to cancel. If the visible head is still queued (or deferred), the
|
|
-- selected row would otherwise replace it without any active_task_id.
|
|
AND EXISTS (
|
|
SELECT 1
|
|
FROM agent_task_queue AS active
|
|
WHERE active.chat_session_id = $1
|
|
AND active.status IN ('dispatched', 'running', 'waiting_local_directory')
|
|
AND active.regenerate_quick_actions_for IS NULL
|
|
)
|
|
FOR UPDATE
|
|
), demoted AS (
|
|
UPDATE agent_task_queue AS queued
|
|
SET priority = 3
|
|
WHERE queued.chat_session_id = $1
|
|
AND queued.id <> $2
|
|
AND queued.status = 'queued'
|
|
AND queued.priority >= 4
|
|
AND EXISTS (SELECT 1 FROM target)
|
|
), prioritized AS (
|
|
UPDATE agent_task_queue AS selected
|
|
SET priority = 4
|
|
FROM target
|
|
WHERE selected.id = target.id
|
|
RETURNING selected.id
|
|
)
|
|
SELECT
|
|
prioritized.id AS task_id,
|
|
(
|
|
SELECT active.id
|
|
FROM agent_task_queue AS active
|
|
WHERE active.chat_session_id = $1
|
|
AND active.status IN ('dispatched', 'running', 'waiting_local_directory')
|
|
AND active.regenerate_quick_actions_for IS NULL
|
|
ORDER BY active.created_at ASC, active.id ASC
|
|
LIMIT 1
|
|
)::uuid AS active_task_id
|
|
FROM prioritized
|
|
`
|
|
|
|
type PrioritizeQueuedChatTaskParams struct {
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
ID pgtype.UUID `json:"id"`
|
|
}
|
|
|
|
type PrioritizeQueuedChatTaskRow struct {
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
ActiveTaskID pgtype.UUID `json:"active_task_id"`
|
|
}
|
|
|
|
func (q *Queries) PrioritizeQueuedChatTask(ctx context.Context, arg PrioritizeQueuedChatTaskParams) (PrioritizeQueuedChatTaskRow, error) {
|
|
row := q.db.QueryRow(ctx, prioritizeQueuedChatTask, arg.ChatSessionID, arg.ID)
|
|
var i PrioritizeQueuedChatTaskRow
|
|
err := row.Scan(&i.TaskID, &i.ActiveTaskID)
|
|
return i, err
|
|
}
|
|
|
|
const promoteChannelChatTasksIfMediaReady = `-- name: PromoteChannelChatTasksIfMediaReady :many
|
|
UPDATE agent_task_queue AS task
|
|
SET status = 'queued', fire_at = NULL
|
|
WHERE task.chat_session_id = $1
|
|
AND task.status = 'deferred'
|
|
AND task.issue_id IS NULL
|
|
AND task.parent_task_id IS NULL
|
|
AND task.escalation_for_task_id IS NULL
|
|
AND NOT EXISTS (
|
|
SELECT 1
|
|
FROM chat_message AS message
|
|
WHERE message.chat_session_id = $1
|
|
AND message.role = 'user'
|
|
AND message.message_kind != 'channel_command'
|
|
AND message.channel_media_pending_until > now()
|
|
)
|
|
RETURNING task.id, task.agent_id, task.issue_id, task.status, task.priority, task.dispatched_at, task.started_at, task.completed_at, task.result, task.error, task.created_at, task.context, task.runtime_id, task.session_id, task.work_dir, task.trigger_comment_id, task.chat_session_id, task.autopilot_run_id, task.attempt, task.max_attempts, task.parent_task_id, task.failure_reason, task.trigger_summary, task.force_fresh_session, task.is_leader_task, task.wait_reason, task.initiator_user_id, task.handoff_note, task.prepare_lease_expires_at, task.squad_id, task.runtime_mcp_overlay, task.escalation_for_task_id, task.fire_at, task.originator_user_id, task.runtime_connected_apps, task.coalesced_comment_ids, task.delivered_comment_ids, task.chat_input_task_id, task.chat_finalize_deferred_at, task.originator_source, task.delegated_from_task_id, task.retry_of_task_id, task.rerun_of_task_id, task.rule_version_id, task.trigger_evidence_kind, task.trigger_evidence_ref_id, task.accountable_user_id, task.session_rollout_missing, task.retired_session_id, task.quick_actions_disabled, task.regenerate_quick_actions_for
|
|
`
|
|
|
|
// Media completion may race with the 3s run batcher. Promote every original
|
|
// channel task waiting for this session only after all unexpired media markers
|
|
// are gone; retry/escalation/direct-chat deferred tasks are excluded. The
|
|
// marker scan uses the same population as GetChannelMediaPendingUntil and the
|
|
// seal: a channel_command turn belongs to no batch, so it must not hold a
|
|
// deferred task back either.
|
|
func (q *Queries) PromoteChannelChatTasksIfMediaReady(ctx context.Context, chatSessionID pgtype.UUID) ([]AgentTaskQueue, error) {
|
|
rows, err := q.db.Query(ctx, promoteChannelChatTasksIfMediaReady, chatSessionID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
items := []AgentTaskQueue{}
|
|
for rows.Next() {
|
|
var i AgentTaskQueue
|
|
if err := rows.Scan(
|
|
&i.ID,
|
|
&i.AgentID,
|
|
&i.IssueID,
|
|
&i.Status,
|
|
&i.Priority,
|
|
&i.DispatchedAt,
|
|
&i.StartedAt,
|
|
&i.CompletedAt,
|
|
&i.Result,
|
|
&i.Error,
|
|
&i.CreatedAt,
|
|
&i.Context,
|
|
&i.RuntimeID,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.TriggerCommentID,
|
|
&i.ChatSessionID,
|
|
&i.AutopilotRunID,
|
|
&i.Attempt,
|
|
&i.MaxAttempts,
|
|
&i.ParentTaskID,
|
|
&i.FailureReason,
|
|
&i.TriggerSummary,
|
|
&i.ForceFreshSession,
|
|
&i.IsLeaderTask,
|
|
&i.WaitReason,
|
|
&i.InitiatorUserID,
|
|
&i.HandoffNote,
|
|
&i.PrepareLeaseExpiresAt,
|
|
&i.SquadID,
|
|
&i.RuntimeMcpOverlay,
|
|
&i.EscalationForTaskID,
|
|
&i.FireAt,
|
|
&i.OriginatorUserID,
|
|
&i.RuntimeConnectedApps,
|
|
&i.CoalescedCommentIds,
|
|
&i.DeliveredCommentIds,
|
|
&i.ChatInputTaskID,
|
|
&i.ChatFinalizeDeferredAt,
|
|
&i.OriginatorSource,
|
|
&i.DelegatedFromTaskID,
|
|
&i.RetryOfTaskID,
|
|
&i.RerunOfTaskID,
|
|
&i.RuleVersionID,
|
|
&i.TriggerEvidenceKind,
|
|
&i.TriggerEvidenceRefID,
|
|
&i.AccountableUserID,
|
|
&i.SessionRolloutMissing,
|
|
&i.RetiredSessionID,
|
|
&i.QuickActionsDisabled,
|
|
&i.RegenerateQuickActionsFor,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
const reanchorClaimedDirectChatInput = `-- name: ReanchorClaimedDirectChatInput :exec
|
|
WITH latest_visible AS (
|
|
SELECT
|
|
claimed_input.id AS claimed_input_id,
|
|
claimed_input.created_at AS input_created_at,
|
|
prior.created_at AS prior_created_at
|
|
FROM chat_message AS claimed_input
|
|
CROSS JOIN LATERAL (
|
|
SELECT prior.created_at
|
|
FROM chat_message AS prior
|
|
WHERE prior.chat_session_id = claimed_input.chat_session_id
|
|
AND prior.id != claimed_input.id
|
|
AND NOT (
|
|
prior.role = 'user'
|
|
AND EXISTS (
|
|
SELECT 1
|
|
FROM agent_task_queue AS queued_task
|
|
WHERE queued_task.chat_session_id = prior.chat_session_id
|
|
AND queued_task.status = 'queued'
|
|
AND queued_task.id = prior.task_id
|
|
AND queued_task.id <> (
|
|
SELECT head.id
|
|
FROM agent_task_queue AS head
|
|
WHERE head.chat_session_id = claimed_input.chat_session_id
|
|
AND head.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
AND head.regenerate_quick_actions_for IS NULL
|
|
ORDER BY
|
|
CASE
|
|
WHEN head.status IN ('dispatched', 'running', 'waiting_local_directory') THEN 0
|
|
WHEN head.status = 'deferred' THEN 1
|
|
ELSE 2
|
|
END,
|
|
head.priority DESC,
|
|
head.created_at ASC,
|
|
head.id ASC
|
|
LIMIT 1
|
|
)
|
|
)
|
|
)
|
|
ORDER BY prior.created_at DESC, prior.id DESC
|
|
LIMIT 1
|
|
) AS prior
|
|
WHERE claimed_input.task_id = $2
|
|
AND claimed_input.role = 'user'
|
|
AND NOT claimed_input.channel_ingested
|
|
)
|
|
UPDATE chat_message AS claimed_input
|
|
SET created_at = GREATEST(
|
|
$1::timestamptz,
|
|
latest_visible.prior_created_at + interval '1 microsecond'
|
|
)
|
|
FROM latest_visible
|
|
WHERE claimed_input.id = latest_visible.claimed_input_id
|
|
AND latest_visible.prior_created_at >= latest_visible.input_created_at
|
|
`
|
|
|
|
type ReanchorClaimedDirectChatInputParams struct {
|
|
DispatchedAt pgtype.Timestamptz `json:"dispatched_at"`
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
}
|
|
|
|
// An idle direct send is visible while it is the positional queue head. A
|
|
// older/pre-deploy follow-up whose enqueue timestamp puts it before a newer
|
|
// settled reply can still reach claim; this is the fallback that moves only
|
|
// that out-of-order user row to the new turn boundary. In-order visible heads
|
|
// keep their original timestamp, so an idle send already returned to a cursor
|
|
// client never moves merely because a daemon claimed it.
|
|
//
|
|
// The task's own created_at remains immutable and continues to drive queue
|
|
// FIFO, wait duration, and elapsed time. Only the message timestamp changes,
|
|
// preserving the existing public ordering/cursor contract for older clients.
|
|
// Add one microsecond when the DB clock ties the latest visible row so UUID
|
|
// ordering can never put the new user turn before the previous assistant row.
|
|
//
|
|
// channel_ingested excludes sealed Slack/Lark batches, which can own multiple
|
|
// user rows and must preserve their original provider order. A direct send
|
|
// owns exactly one non-channel user row today; if direct batching is added,
|
|
// this query must assign a stable per-row offset instead of one timestamp.
|
|
// Retry children are excluded by the caller because their chat_input_task_id
|
|
// names the root task rather than themselves. Stale dispatched reclaim bypasses
|
|
// this query, so re-delivery does not move an already-visible turn.
|
|
//
|
|
// Lock order is task row first (ClaimAgentTask), then this non-FK message-only
|
|
// update. The visibility subqueries take no row locks, avoiding an inverse
|
|
// message -> task lock edge with send/completion transactions.
|
|
func (q *Queries) ReanchorClaimedDirectChatInput(ctx context.Context, arg ReanchorClaimedDirectChatInputParams) error {
|
|
_, err := q.db.Exec(ctx, reanchorClaimedDirectChatInput, arg.DispatchedAt, arg.TaskID)
|
|
return err
|
|
}
|
|
|
|
const reanchorNextQueuedDirectChatInput = `-- name: ReanchorNextQueuedDirectChatInput :exec
|
|
UPDATE chat_message AS queued_input
|
|
SET created_at = $2::timestamptz + interval '1 microsecond'
|
|
WHERE queued_input.chat_session_id = $1
|
|
AND queued_input.role = 'user'
|
|
AND NOT queued_input.channel_ingested
|
|
AND queued_input.created_at <= $2::timestamptz
|
|
AND EXISTS (
|
|
SELECT 1
|
|
FROM agent_task_queue AS queued_task
|
|
WHERE queued_task.id = queued_input.task_id
|
|
AND queued_task.chat_session_id = queued_input.chat_session_id
|
|
AND queued_task.status = 'queued'
|
|
AND queued_task.chat_input_task_id = queued_task.id
|
|
AND queued_task.regenerate_quick_actions_for IS NULL
|
|
AND queued_task.id = (
|
|
SELECT head.id
|
|
FROM agent_task_queue AS head
|
|
WHERE head.chat_session_id = queued_input.chat_session_id
|
|
AND head.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory', 'deferred')
|
|
AND head.regenerate_quick_actions_for IS NULL
|
|
ORDER BY
|
|
CASE
|
|
WHEN head.status IN ('dispatched', 'running', 'waiting_local_directory') THEN 0
|
|
WHEN head.status = 'deferred' THEN 1
|
|
ELSE 2
|
|
END,
|
|
head.priority DESC,
|
|
head.created_at ASC,
|
|
head.id ASC
|
|
LIMIT 1
|
|
)
|
|
)
|
|
`
|
|
|
|
type ReanchorNextQueuedDirectChatInputParams struct {
|
|
ChatSessionID pgtype.UUID `json:"chat_session_id"`
|
|
AssistantCreatedAt pgtype.Timestamptz `json:"assistant_created_at"`
|
|
}
|
|
|
|
// An assistant outcome can make the next queued direct task the positional
|
|
// head before a daemon claims it (MUL-5750). Its user row was persisted before
|
|
// this reply, so move that still-hidden single-row input just after the reply.
|
|
// Callers run this immediately after CreateChatMessage in the same transaction:
|
|
// readers see either the old active head without the reply, or the settled
|
|
// reply followed by the newly-visible head — never user B before assistant A.
|
|
// Channel batches are excluded because their immutable provider ordering may
|
|
// contain multiple user rows. Direct chat currently owns exactly one row; if
|
|
// direct batching is added, assign stable per-row offsets here.
|
|
func (q *Queries) ReanchorNextQueuedDirectChatInput(ctx context.Context, arg ReanchorNextQueuedDirectChatInputParams) error {
|
|
_, err := q.db.Exec(ctx, reanchorNextQueuedDirectChatInput, arg.ChatSessionID, arg.AssistantCreatedAt)
|
|
return err
|
|
}
|
|
|
|
const setChatMessageQuickActionsByTask = `-- name: SetChatMessageQuickActionsByTask :one
|
|
UPDATE chat_message
|
|
SET quick_actions = $2
|
|
WHERE id = (
|
|
SELECT inner_msg.id FROM chat_message AS inner_msg
|
|
WHERE inner_msg.task_id = $1 AND inner_msg.role = 'assistant'
|
|
ORDER BY inner_msg.created_at DESC
|
|
LIMIT 1
|
|
)
|
|
RETURNING id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested, quick_actions
|
|
`
|
|
|
|
type SetChatMessageQuickActionsByTaskParams struct {
|
|
TaskID pgtype.UUID `json:"task_id"`
|
|
QuickActions []byte `json:"quick_actions"`
|
|
}
|
|
|
|
func (q *Queries) SetChatMessageQuickActionsByTask(ctx context.Context, arg SetChatMessageQuickActionsByTaskParams) (ChatMessage, error) {
|
|
row := q.db.QueryRow(ctx, setChatMessageQuickActionsByTask, arg.TaskID, arg.QuickActions)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
&i.MessageKind,
|
|
&i.ChannelMediaPendingUntil,
|
|
&i.ChannelIngested,
|
|
&i.QuickActions,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const setChatSessionArchived = `-- name: SetChatSessionArchived :one
|
|
UPDATE chat_session
|
|
SET status = CASE WHEN $2::bool THEN 'archived' ELSE 'active' END,
|
|
updated_at = now()
|
|
WHERE id = $1
|
|
RETURNING id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id
|
|
`
|
|
|
|
type SetChatSessionArchivedParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
Archived bool `json:"archived"`
|
|
}
|
|
|
|
// Archive/unarchive a chat session by flipping status between 'active' and
|
|
// 'archived'. Bumps updated_at so the row re-sorts on the receiving list. The
|
|
// send-message path refuses archived sessions (see SendChatMessage), so the
|
|
// conversation is effectively read-only until it is unarchived.
|
|
func (q *Queries) SetChatSessionArchived(ctx context.Context, arg SetChatSessionArchivedParams) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, setChatSessionArchived, arg.ID, arg.Archived)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const setChatSessionPinned = `-- name: SetChatSessionPinned :one
|
|
UPDATE chat_session
|
|
SET pinned_at = CASE WHEN $2::bool THEN COALESCE(pinned_at, now()) ELSE NULL END
|
|
WHERE id = $1
|
|
RETURNING id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id
|
|
`
|
|
|
|
type SetChatSessionPinnedParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
Pinned bool `json:"pinned"`
|
|
}
|
|
|
|
// Pin/unpin a chat. Deliberately does NOT touch updated_at: pinning is a
|
|
// list-ordering preference, not activity, so it must not bump the session's
|
|
// last-activity sort key (which would make an unpinned chat jump the list).
|
|
// pinned = true stamps pinned_at only when it was NULL, so re-pinning keeps
|
|
// the original pin order; pinned = false clears it.
|
|
func (q *Queries) SetChatSessionPinned(ctx context.Context, arg SetChatSessionPinnedParams) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, setChatSessionPinned, arg.ID, arg.Pinned)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const setChatTaskInputOwnerSelf = `-- name: SetChatTaskInputOwnerSelf :one
|
|
UPDATE agent_task_queue
|
|
SET chat_input_task_id = id
|
|
WHERE id = $1
|
|
RETURNING id, agent_id, issue_id, status, priority, dispatched_at, started_at, completed_at, result, error, created_at, context, runtime_id, session_id, work_dir, trigger_comment_id, chat_session_id, autopilot_run_id, attempt, max_attempts, parent_task_id, failure_reason, trigger_summary, force_fresh_session, is_leader_task, wait_reason, initiator_user_id, handoff_note, prepare_lease_expires_at, squad_id, runtime_mcp_overlay, escalation_for_task_id, fire_at, originator_user_id, runtime_connected_apps, coalesced_comment_ids, delivered_comment_ids, chat_input_task_id, chat_finalize_deferred_at, originator_source, delegated_from_task_id, retry_of_task_id, rerun_of_task_id, rule_version_id, trigger_evidence_kind, trigger_evidence_ref_id, accountable_user_id, session_rollout_missing, retired_session_id, quick_actions_disabled, regenerate_quick_actions_for
|
|
`
|
|
|
|
// Stamps a freshly-created direct-chat task as the owner of its own input batch
|
|
// (chat_input_task_id = id), so a later claim loads exactly the user messages
|
|
// tagged with this task id (ListChatInputMessages) rather than scanning trailing
|
|
// history. Runs in the same transaction as CreateChatTask + message ownership:
|
|
// direct-send inserts one owned message, while channel enqueue seals its
|
|
// trailing unowned batch. Legacy tasks keep chat_input_task_id NULL and retain
|
|
// the trailing-history fallback during rolling deploys.
|
|
func (q *Queries) SetChatTaskInputOwnerSelf(ctx context.Context, id pgtype.UUID) (AgentTaskQueue, error) {
|
|
row := q.db.QueryRow(ctx, setChatTaskInputOwnerSelf, id)
|
|
var i AgentTaskQueue
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.AgentID,
|
|
&i.IssueID,
|
|
&i.Status,
|
|
&i.Priority,
|
|
&i.DispatchedAt,
|
|
&i.StartedAt,
|
|
&i.CompletedAt,
|
|
&i.Result,
|
|
&i.Error,
|
|
&i.CreatedAt,
|
|
&i.Context,
|
|
&i.RuntimeID,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.TriggerCommentID,
|
|
&i.ChatSessionID,
|
|
&i.AutopilotRunID,
|
|
&i.Attempt,
|
|
&i.MaxAttempts,
|
|
&i.ParentTaskID,
|
|
&i.FailureReason,
|
|
&i.TriggerSummary,
|
|
&i.ForceFreshSession,
|
|
&i.IsLeaderTask,
|
|
&i.WaitReason,
|
|
&i.InitiatorUserID,
|
|
&i.HandoffNote,
|
|
&i.PrepareLeaseExpiresAt,
|
|
&i.SquadID,
|
|
&i.RuntimeMcpOverlay,
|
|
&i.EscalationForTaskID,
|
|
&i.FireAt,
|
|
&i.OriginatorUserID,
|
|
&i.RuntimeConnectedApps,
|
|
&i.CoalescedCommentIds,
|
|
&i.DeliveredCommentIds,
|
|
&i.ChatInputTaskID,
|
|
&i.ChatFinalizeDeferredAt,
|
|
&i.OriginatorSource,
|
|
&i.DelegatedFromTaskID,
|
|
&i.RetryOfTaskID,
|
|
&i.RerunOfTaskID,
|
|
&i.RuleVersionID,
|
|
&i.TriggerEvidenceKind,
|
|
&i.TriggerEvidenceRefID,
|
|
&i.AccountableUserID,
|
|
&i.SessionRolloutMissing,
|
|
&i.RetiredSessionID,
|
|
&i.QuickActionsDisabled,
|
|
&i.RegenerateQuickActionsFor,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const taskHasChannelIngestedMessages = `-- name: TaskHasChannelIngestedMessages :one
|
|
SELECT EXISTS (
|
|
SELECT 1 FROM chat_message
|
|
WHERE task_id = $1
|
|
AND role = 'user'
|
|
AND channel_ingested
|
|
) AS channel_ingested
|
|
`
|
|
|
|
// Immutable channel provenance for a task's user-message input batch:
|
|
// channel_ingested is stamped inside the channel append transaction and never
|
|
// mutated afterwards, so it survives session archiving and installation
|
|
// rebinds that delete the channel_chat_session_binding row. Callers pass the
|
|
// batch OWNER id (chat_input_task_id, which auto-retry clones inherit), not
|
|
// necessarily the task's own id. The cancel restore-delete and the
|
|
// empty-completion silent-drop both gate on this — a channel sender has no
|
|
// Multica composer for a restored draft, and the no_response fallback body
|
|
// must never be pushed to an external channel.
|
|
func (q *Queries) TaskHasChannelIngestedMessages(ctx context.Context, taskID pgtype.UUID) (bool, error) {
|
|
row := q.db.QueryRow(ctx, taskHasChannelIngestedMessages, taskID)
|
|
var channel_ingested bool
|
|
err := row.Scan(&channel_ingested)
|
|
return channel_ingested, err
|
|
}
|
|
|
|
const touchChatSession = `-- name: TouchChatSession :exec
|
|
UPDATE chat_session SET updated_at = now()
|
|
WHERE id = $1
|
|
`
|
|
|
|
func (q *Queries) TouchChatSession(ctx context.Context, id pgtype.UUID) error {
|
|
_, err := q.db.Exec(ctx, touchChatSession, id)
|
|
return err
|
|
}
|
|
|
|
const updateChatSessionProject = `-- name: UpdateChatSessionProject :one
|
|
UPDATE chat_session
|
|
SET project_id = $1
|
|
WHERE id = $2 AND workspace_id = $3
|
|
RETURNING id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id
|
|
`
|
|
|
|
type UpdateChatSessionProjectParams struct {
|
|
ProjectID pgtype.UUID `json:"project_id"`
|
|
ID pgtype.UUID `json:"id"`
|
|
WorkspaceID pgtype.UUID `json:"workspace_id"`
|
|
}
|
|
|
|
// Project context is user-editable session metadata. Do not touch updated_at:
|
|
// changing context is not conversation activity and must not reorder history.
|
|
func (q *Queries) UpdateChatSessionProject(ctx context.Context, arg UpdateChatSessionProjectParams) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, updateChatSessionProject, arg.ProjectID, arg.ID, arg.WorkspaceID)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const updateChatSessionSession = `-- name: UpdateChatSessionSession :exec
|
|
UPDATE chat_session
|
|
SET session_id = COALESCE($1, session_id),
|
|
work_dir = COALESCE($2, work_dir),
|
|
runtime_id = COALESCE($3, runtime_id),
|
|
updated_at = now()
|
|
WHERE id = $4
|
|
`
|
|
|
|
type UpdateChatSessionSessionParams struct {
|
|
SessionID pgtype.Text `json:"session_id"`
|
|
WorkDir pgtype.Text `json:"work_dir"`
|
|
RuntimeID pgtype.UUID `json:"runtime_id"`
|
|
ID pgtype.UUID `json:"id"`
|
|
}
|
|
|
|
// Updates the resume pointer for a chat session. Empty/NULL inputs are
|
|
// ignored via COALESCE so a task that completes without a session_id (e.g.
|
|
// the agent crashed before establishing one) cannot wipe out a previously
|
|
// recorded resume pointer. This makes the chat memory robust against
|
|
// intermittent agent failures.
|
|
func (q *Queries) UpdateChatSessionSession(ctx context.Context, arg UpdateChatSessionSessionParams) error {
|
|
_, err := q.db.Exec(ctx, updateChatSessionSession,
|
|
arg.SessionID,
|
|
arg.WorkDir,
|
|
arg.RuntimeID,
|
|
arg.ID,
|
|
)
|
|
return err
|
|
}
|
|
|
|
const updateChatSessionTitle = `-- name: UpdateChatSessionTitle :one
|
|
UPDATE chat_session SET title = $2, updated_at = now()
|
|
WHERE id = $1
|
|
RETURNING id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id
|
|
`
|
|
|
|
type UpdateChatSessionTitleParams struct {
|
|
ID pgtype.UUID `json:"id"`
|
|
Title string `json:"title"`
|
|
}
|
|
|
|
func (q *Queries) UpdateChatSessionTitle(ctx context.Context, arg UpdateChatSessionTitleParams) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, updateChatSessionTitle, arg.ID, arg.Title)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const updateChatSessionTitleIfCurrent = `-- name: UpdateChatSessionTitleIfCurrent :one
|
|
UPDATE chat_session SET title = $1, updated_at = now()
|
|
WHERE id = $2 AND title = $3
|
|
RETURNING id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_id, last_read_at, is_agent_intro, pinned_at, project_id
|
|
`
|
|
|
|
type UpdateChatSessionTitleIfCurrentParams struct {
|
|
NewTitle string `json:"new_title"`
|
|
ID pgtype.UUID `json:"id"`
|
|
ExpectedTitle string `json:"expected_title"`
|
|
}
|
|
|
|
// Compare-and-swap the title: only overwrite it when it still equals the
|
|
// value the caller observed (@expected_title). This is the idempotency /
|
|
// no-clobber guard behind LLM auto-titling (MUL-4295): the async generator
|
|
// captures the session's current (default/original) title before calling the
|
|
// model, and this write lands only if a manual rename or a competing writer
|
|
// has not changed the title in the meantime. A mismatch returns pgx.ErrNoRows
|
|
// (zero rows updated), which the caller treats as "someone renamed it — leave
|
|
// it alone", NOT as an error.
|
|
func (q *Queries) UpdateChatSessionTitleIfCurrent(ctx context.Context, arg UpdateChatSessionTitleIfCurrentParams) (ChatSession, error) {
|
|
row := q.db.QueryRow(ctx, updateChatSessionTitleIfCurrent, arg.NewTitle, arg.ID, arg.ExpectedTitle)
|
|
var i ChatSession
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.WorkspaceID,
|
|
&i.AgentID,
|
|
&i.CreatorID,
|
|
&i.Title,
|
|
&i.SessionID,
|
|
&i.WorkDir,
|
|
&i.Status,
|
|
&i.CreatedAt,
|
|
&i.UpdatedAt,
|
|
&i.UnreadSince,
|
|
&i.RuntimeID,
|
|
&i.LastReadAt,
|
|
&i.IsAgentIntro,
|
|
&i.PinnedAt,
|
|
&i.ProjectID,
|
|
)
|
|
return i, err
|
|
}
|