mirror of
https://github.com/multica-ai/multica.git
synced 2026-08-03 19:20:07 +02:00
* fix(daemon): retire sessions whose history the provider refuses to replay A run killed mid-reply (machine shutdown, force-quit, SIGKILL) can leave an empty assistant message in the agent CLI's transcript. Every later resume replays it, the provider rejects the request, and the (agent, issue) pair is bricked with no self-healing and no user-facing recovery. Multica already has the mechanism for this — poisoned-session classification — but its detector paired "400" with "invalid_request_error", which is the Anthropic wire shape. The same defect reported by any other provider carried neither token, so it classified as agent_error.unknown: resume-safe by omission. GetLastTaskSession kept handing back the dead session on every follow-up, manual Rerun resolved it through the same predicate, and the in-turn fresh-session retry never fired because ResumeRejected is false here (nothing rejected the resume — the transcript loaded and the provider refused to replay it). Add taskfailure.UnresumableHistory, which recognises the defect by what the provider says is wrong — some content is empty, and here is which message in the history — rather than by status code or provider name. Both signals are required, so a tool reporting "field must not be empty" does not match. Wire it into the four places that decide whether a session survives: - classifyPoisonedError, so the task is written as api_invalid_request - shouldRetryWithFreshSession, so the turn recovers on all 17 backends instead of the subset whose adapter learned to detect it; the tools == 0 gate is unchanged, so a run that already used a tool is never re-run - ResumeUnsafeFailure, covering the manual-Rerun path - both resume queries, as defense-in-depth for hosts whose daemon predates this (self-host daemons upgrade on their own cadence) Fixes #6066. Also covers the daemon half of #5760. Co-authored-by: multica-agent <github@multica.ai> * fix(session): close the Chat and fresh-retry paths that resurrect a poisoned session Review found the previous commit stopped short in two places, both of which put the dead transcript back in play. Chat never consulted the guarded query. The claim handler reads chat_session.session_id first and only falls back to GetLastChatTaskSession when it is empty, so a poisoned pointer there bypasses every filter that query applies. The fail path merely declined to OVERWRITE the pointer, leaving it in place. It now clears it in the same transaction, matched on session and runtime so a concurrent turn's newer pointer survives. The promote guard moves to ResumeUnsafeFailure as well — the reason-only check passed an un-upgraded daemon's agent_error.unknown row and re-pinned what the clear had just removed. GetLastChatTaskSession also kept the row-level filter the issue query dropped in GH #5975: it discarded the newest poisoned row and fell back to an older completed row carrying the same dead session. It now judges each session by its latest terminal state, matching GetLastTaskSession. A recovered turn could not retire anything. A terminal report carried one session_id, and an empty one meant both "nothing to report" and "forget the old session", so a fresh-session retry that SUCCEEDED left the id it retried away from selectable — through an older completed row on the issue, or through the chat pointer. agent_task_queue.retired_session_id records the abandonment itself, reported on every terminal path including completed, and both resume lookups exclude it. This is the contract gap the previous PR deferred; the fresh-retry path now runs on all backends, so deferring it is not safe. Also narrows what the cross-backend test claims: it pins the shared decision, not that all 17 adapters surface the error into Result.Error (#5760 is the counter-example), and says so. Co-authored-by: multica-agent <github@multica.ai> * test(session): require pgx.ErrNoRows in the resume-exclusion assertions The `if err == nil && prior.SessionID.Valid` form these tests shared is false-green: any real fault — undefined column, syntax error, dead connection — makes err non-nil, so the condition is false and the test passes. Run against a database missing this branch's new column, the exclusion tests reported PASS on a SQLSTATE 42703, meaning they could not have caught a broken query. requireSessionExcluded demands pgx.ErrNoRows specifically and fails loudly on anything else, so a green run now means the filter worked rather than the query never ran. Applied to all nine sites, not just the four this branch added: the other five guard the same GetLastTaskSession exclusion behaviour that this branch changes, so leaving them false-green would leave the change under-tested. All nine pass on a correctly migrated database. Co-authored-by: multica-agent <github@multica.ai> --------- Co-authored-by: Bohan-J <bohan@devv.ai> Co-authored-by: multica-agent <github@multica.ai>
1918 lines
68 KiB
Go
1918 lines
68 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 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, channel_media_pending_until, channel_ingested
|
|
)
|
|
VALUES (
|
|
$1, $2, $3, $4, $5, $6,
|
|
COALESCE($7::text, 'message'),
|
|
-- 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 $8::float8 IS NULL THEN NULL
|
|
ELSE now() + make_interval(secs => $8::float8) END,
|
|
COALESCE($9::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
|
|
`
|
|
|
|
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"`
|
|
ChannelMediaPendingSecs pgtype.Float8 `json:"channel_media_pending_secs"`
|
|
ChannelIngested pgtype.Bool `json:"channel_ingested"`
|
|
}
|
|
|
|
// message_kind defaults to 'message' 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.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,
|
|
)
|
|
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
|
|
`
|
|
|
|
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,
|
|
)
|
|
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
|
|
`
|
|
|
|
// 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,
|
|
)
|
|
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 deleteChatDraftRestoresByArchivedRuntimeAgents = `-- name: DeleteChatDraftRestoresByArchivedRuntimeAgents :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.archived_at IS NOT NULL
|
|
)
|
|
`
|
|
|
|
// chat_session cascades from agent, so hard-deleting a runtime's archived 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
|
|
// DeleteChannelInstallationsByArchivedRuntimeAgents.
|
|
func (q *Queries) DeleteChatDraftRestoresByArchivedRuntimeAgents(ctx context.Context, runtimeID pgtype.UUID) error {
|
|
_, err := q.db.Exec(ctx, deleteChatDraftRestoresByArchivedRuntimeAgents, runtimeID)
|
|
return err
|
|
}
|
|
|
|
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'
|
|
)
|
|
`
|
|
|
|
// Same cascade, for the system agents a runtime teardown also hard-deletes
|
|
// (DeleteSystemAgentsByRuntime). Split from the archived-agent prune because the
|
|
// runtime-profile teardown deletes only archived agents: pruning system-agent
|
|
// sessions there would destroy restores whose session survives.
|
|
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
|
|
`
|
|
|
|
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,
|
|
)
|
|
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 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.
|
|
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 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,
|
|
)
|
|
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
|
|
), 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')
|
|
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 = 'completed'
|
|
OR (
|
|
status = 'failed'
|
|
AND COALESCE(failure_reason, '') NOT IN ('iteration_limit', 'agent_fallback_message', 'api_invalid_request', 'codex_semantic_inactivity', 'agent_error.context_overflow')
|
|
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]')
|
|
)
|
|
)
|
|
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 both completed and failed tasks: even a failed task
|
|
// may have established a real agent session before failing, 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.
|
|
//
|
|
// 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 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 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,
|
|
)
|
|
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')
|
|
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 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')
|
|
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 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 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.
|
|
func (q *Queries) LinkUnownedChannelChatMessagesToTask(ctx context.Context, arg LinkUnownedChannelChatMessagesToTaskParams) error {
|
|
_, err := q.db.Exec(ctx, linkUnownedChannelChatMessagesToTask, arg.TaskID, arg.ChatSessionID)
|
|
return err
|
|
}
|
|
|
|
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 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 ListChatMessages + 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,
|
|
); 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 id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested FROM chat_message
|
|
WHERE chat_session_id = $1
|
|
ORDER BY created_at ASC, id ASC
|
|
`
|
|
|
|
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,
|
|
); 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 id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms, message_kind, channel_media_pending_until, channel_ingested FROM chat_message
|
|
WHERE chat_session_id = $1
|
|
AND (
|
|
$3::timestamptz IS NULL
|
|
OR (created_at, id) < ($3::timestamptz, $4::uuid)
|
|
)
|
|
ORDER BY created_at DESC, 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,
|
|
); 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')
|
|
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_v2 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 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 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 lockChatSessionsByArchivedRuntimeAgents = `-- name: LockChatSessionsByArchivedRuntimeAgents :many
|
|
SELECT cs.id FROM chat_session cs
|
|
JOIN agent a ON a.id = cs.agent_id
|
|
WHERE a.runtime_id = $1 AND a.archived_at IS NOT NULL
|
|
ORDER BY cs.id
|
|
FOR UPDATE OF cs
|
|
`
|
|
|
|
func (q *Queries) LockChatSessionsByArchivedRuntimeAgents(ctx context.Context, runtimeID pgtype.UUID) ([]pgtype.UUID, error) {
|
|
rows, err := q.db.Query(ctx, lockChatSessionsByArchivedRuntimeAgents, 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 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 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.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
|
|
`
|
|
|
|
// 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.
|
|
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,
|
|
); err != nil {
|
|
return nil, err
|
|
}
|
|
items = append(items, i)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return nil, err
|
|
}
|
|
return items, nil
|
|
}
|
|
|
|
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
|
|
`
|
|
|
|
// 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,
|
|
)
|
|
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
|
|
}
|