mirror of
https://github.com/multica-ai/multica.git
synced 2026-07-24 02:39:42 +02:00
* feat(issues): unify run-enqueue decision behind WillEnqueueRun + preview endpoint Collapse the issue update/batch enqueue copies into one service predicate service.IssueService.WillEnqueueRun, shared verbatim with a new dry-run endpoint POST /api/issues/preview-trigger so the four entry points stop drifting (squad/self-loop/batch omissions, MUL-3375). The private-agent gate stays at the HTTP boundary: write paths inject allow-all, preview injects the real gate so it never leaks a private agent's readiness. Add suppress_run to issue update/batch: the change applies but no run starts. Remove the now-dead handler mirrors shouldEnqueueSquadLeaderOnAssign / isSquadLeaderReady. service.Create and the comment trigger chain are untouched. Tests: preview behavior, preview<->write-path match, batch aggregation, member no-trigger, suppress_run skip, malformed-body 400. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * feat(issues): inject handoff note into assigned runs via first-class task field Add an optional handoff_note carried by issue assign/promote into the run's opening prompt and issue_context.md, via a dedicated agent_task_queue column (migration 122) and a daemon assignment-handoff render branch — never a fabricated comment, never trigger_comment_id (MUL-3375 §6.1). Thread the note through enqueueIssueTask/enqueueMentionTask + WithHandoff public variants and dispatchIssueRun; suppress_run or a parked write drops it (no run = nothing to inject). Soft version gate: MinHandoffCLIVersion + HandoffSupported, surfaced per-trigger as handoff_supported in the preview so the UI can gray the note box on old daemons; the assignment never hard-fails. Tests: daemon prompt + issue_context render via the assignment branch (not quick-create/comment), version helper matrix, note persists on the task, suppressed assign enqueues nothing. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * feat(issues): leave a display-only handoff record on the timeline When an assign/promote with a handoff note starts a run, write one type='handoff' timeline record via TaskService.RecordHandoff — a direct Queries.CreateComment + timeline event that bypasses Handler.CreateComment, so it never reaches triggerTasksForComment and cannot start a second run (MUL-3375 §6.2, the must-not-retrigger invariant). Author is the actor who handed off; body is the note. Migration 123 admits the 'handoff' comment type. Recorded only on a real run start: suppress_run or a parked write writes nothing. enqueueSquadLeaderTask now reports whether it enqueued so the trace is gated on an actual dispatch. Test: exactly one handoff record on assign-with-note, exactly one task (no re-trigger), and no record when suppressed. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * feat(issues): frontend plumbing for issue-trigger preview + handoff (core) Add api.previewIssueTrigger + IssueTriggerPreviewSchema (zod parseWithFallback), the use-issue-trigger-preview hook, issueKeys.issueTriggerPreview(+All) with WS queue-state invalidation, suppress_run/handoff_note on UpdateIssueRequest, the 'handoff' CommentType, and stripping of the control fields from optimistic update/batch cache patches (MUL-3375 §9). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * fix(issues): exclude handoff records from new-comment counting type='handoff' is a display-only timeline record, not conversation. Exclude it from CountNewCommentsSince so a handoff note never inflates the count of "new comments to catch up on" fed to a claiming agent (MUL-3375 §12). Analytics already excludes it (RecordHandoff is a direct write that emits no analytics event), and the comment-trigger path is already bypassed. Test: a handoff record does not bump the new-comment count; a real comment does. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * feat(issues): pre-trigger preview UI, handoff note, timeline card (web/desktop) Wire the §9 frontend onto the preview endpoint + handoff fields: - Delete the backlog blocking dialog (backlog-agent-hint*) and its modal type; the over-eager nag is gone. Backlog awareness is now a passive label. - RunConfirmModal: single assign + batch assign/status route here. Shows the backend predicate's verdict ("将启动 @X" / "将启动 N 个" / parked), an optional handoff note (assign only, soft-gated by handoff_supported), and 暂不启动 — then applies via update/batch. No frontend guessing. - create modal: passive CreateRunHint ("将启动 @X" / backlog parked). - single status change stays a direct apply (unchanged). - timeline: render type='handoff' as a distinct, non-interactive handoff card. - i18n run_confirm + handoff_card across en/ja/ko/zh-Hans; drop backlog action keys; locale parity green. Tests: use-issue-actions (assign → run-confirm modal, member → direct), create-issue + comment-card suites updated/green; views typecheck + lint clean. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * test(issues): use a valid anchor in the handoff count-exclusion test CountNewCommentsSince filters id <> @anchor_id; SQL id <> NULL is NULL and excludes every row, so an empty anchor made the control assertion read 0. The production caller always passes a real anchor — mirror that with a non-matching sentinel uuid. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * test(issues): RunConfirmModal apply logic (start/suppress/note-gate/batch) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * test(core): preview schema malformed/missing/null fallback coverage Cover IssueTriggerPreviewSchema via parseWithFallback (MUL-3375): well-formed parse, top-level + item default fills (empty/older backend), and fallback to { triggers: [], total_count: 0 } for malformed shapes, a dropped required issue_id, a wrong-typed total_count, and null/non-object bodies — so the four entry points degrade to "nothing will start" instead of throwing. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * refactor(issues): remove display-only handoff timeline record (留痕) The handoff "留痕" timeline record (type='handoff' comment written on run start) was judged superfluous and dropped per product call. This removes only the display-only trace; the handoff NOTE injection into the run's opening prompt + issue_context.md is untouched. - backend: drop RecordHandoff + its call in dispatchIssueRun - db: drop the `type <> 'handoff'` exclusion in CountNewCommentsSince and migration 123 (comment_type_check reverts to the 4-type set from 001); no production data exists for this unreleased feature - frontend: drop the "handoff" CommentType, HandoffCard, and handoff_card i18n (all locales) - tests: drop handoff_count_test.go and the record-write assertions in issue_trigger_preview_test.go (note-injection tests retained) Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * feat(issues): dismissable run-confirm modal + team-handoff copy Two fixes to the pre-trigger confirm modal (MUL-3375). 1. Dismissable: switch RunConfirmModal from AlertDialog to the standard shadcn Dialog so it has the close (X) button + Esc + click-outside. Previously the only choices were "start" / "don't start now" with no way to abort the action entirely; dismissing now cancels with no write. 2. Copy: rework the action-surface wording away from the backend term "run" toward team-handoff voice — 指派 / 开始 / 交接 (run stays only on record surfaces). Unifies the note's three names to "交接说明", and parallels the rewrite across en/ja/ko. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai> * chore(agent): bump handoff note min CLI version to 0.3.28 The daemon release that renders handoff notes ships in 0.3.28 (0.3.27 was the prior tag), so move the soft-gate threshold up. Below this the note is silently dropped and the frontend grays the note box — assignment is never blocked. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * fix(issues): skip run-confirm when batch-moving issues to backlog A move into backlog never starts a run (service/issue_trigger.go), so the pre-trigger confirm modal degenerated to an empty "won't start" box with a single Apply button — pure friction. Apply directly instead, matching the single-issue status path. Other target statuses still route through the modal. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * feat(issues): refine pre-trigger preview hint and copy - Move the create-issue run hint to a reveal band (grid 0fr→1fr) above the property toolbar. It was sharing the footer button row and, lacking a width constraint, reflowed the submit buttons whenever it appeared. Restyle to a borderless, comment-style avatar+caption that is purely a caption (non-interactive avatar). - Distinguish squad from agent in the pre-trigger copy: a squad's leader evaluates and delegates rather than "starting work" itself. Add will_start_named_squad / will_start_squad / create_will_start_squad across en/zh/ja/ko (reusing the squad_leader_* evaluate→arrange vocabulary) and branch run-confirm + the create hint on squad assignees. - Bold the assignee name in the run-confirm headline via a language-safe sentinel split (no per-language prefix/suffix keys). - Align zh "开始处理" → "开始工作" on the single-assign copy. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> * test(issues): stub ActorAvatar in create-issue suite CreateRunHint now renders an ActorAvatar for agent/squad assignees, which pulls in getActorInitials/getActorAvatarUrl + the workspace/presence/navigation hook tree. This form-focused suite only stubbed getActorName, so the squad-forwarding test crashed with "getActorInitials is not a function". Stub the avatar inert — its own behavior is covered elsewhere. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com> --------- Co-authored-by: Walt <walt@multica.ai> Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com> Co-authored-by: multica-agent <github@multica.ai>
746 lines
22 KiB
Go
746 lines
22 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 createChatMessage = `-- name: CreateChatMessage :one
|
|
INSERT INTO chat_message (chat_session_id, role, content, task_id, failure_reason, elapsed_ms)
|
|
VALUES ($1, $2, $3, $4, $5, $6)
|
|
RETURNING id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms
|
|
`
|
|
|
|
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"`
|
|
}
|
|
|
|
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,
|
|
)
|
|
var i ChatMessage
|
|
err := row.Scan(
|
|
&i.ID,
|
|
&i.ChatSessionID,
|
|
&i.Role,
|
|
&i.Content,
|
|
&i.TaskID,
|
|
&i.CreatedAt,
|
|
&i.FailureReason,
|
|
&i.ElapsedMs,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const createChatSession = `-- name: CreateChatSession :one
|
|
INSERT INTO chat_session (workspace_id, agent_id, creator_id, title, runtime_id)
|
|
VALUES ($1, $2, $3, $4, (SELECT runtime_id FROM agent WHERE id = $2))
|
|
RETURNING id, workspace_id, agent_id, creator_id, title, session_id, work_dir, status, created_at, updated_at, unread_since, runtime_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"`
|
|
}
|
|
|
|
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,
|
|
)
|
|
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,
|
|
)
|
|
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)
|
|
VALUES ($1, $2, NULL, 'queued', $3, $4, $5)
|
|
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
|
|
`
|
|
|
|
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"`
|
|
}
|
|
|
|
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,
|
|
)
|
|
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,
|
|
)
|
|
return i, 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
|
|
`
|
|
|
|
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,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getChatMessage = `-- name: GetChatMessage :one
|
|
SELECT id, chat_session_id, role, content, task_id, created_at, failure_reason, elapsed_ms 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,
|
|
)
|
|
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 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,
|
|
)
|
|
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 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,
|
|
)
|
|
return i, err
|
|
}
|
|
|
|
const getLastChatTaskSession = `-- name: GetLastChatTaskSession :one
|
|
SELECT session_id, work_dir, runtime_id FROM agent_task_queue
|
|
WHERE chat_session_id = $1
|
|
AND (
|
|
status = 'completed'
|
|
OR (
|
|
status = 'failed'
|
|
AND COALESCE(failure_reason, '') NOT IN ('iteration_limit', 'agent_fallback_message', 'api_invalid_request', 'codex_semantic_inactivity')
|
|
AND NOT (COALESCE(error, '') ILIKE '%400%' AND COALESCE(error, '') ILIKE '%invalid_request_error%')
|
|
)
|
|
)
|
|
AND session_id IS NOT NULL
|
|
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.
|
|
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 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,
|
|
)
|
|
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 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 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.unread_since IS NOT NULL)::bool AS has_unread
|
|
FROM chat_session cs
|
|
WHERE cs.workspace_id = $1 AND cs.creator_id = $2
|
|
ORDER BY 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"`
|
|
HasUnread bool `json:"has_unread"`
|
|
}
|
|
|
|
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.HasUnread,
|
|
); 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 FROM chat_message
|
|
WHERE chat_session_id = $1
|
|
ORDER BY created_at 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,
|
|
); 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 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,
|
|
); 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.unread_since IS NOT NULL)::bool AS has_unread
|
|
FROM chat_session cs
|
|
WHERE cs.workspace_id = $1 AND cs.creator_id = $2 AND cs.status = 'active'
|
|
ORDER BY 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"`
|
|
HasUnread bool `json:"has_unread"`
|
|
}
|
|
|
|
// Returns active sessions with a boolean unread flag. Unread is strictly
|
|
// per-session: either the user has uncleared assistant replies in this
|
|
// session or they don't. Counting messages would be misleading.
|
|
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.HasUnread,
|
|
); 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
|
|
FROM agent_task_queue atq
|
|
JOIN chat_session cs ON cs.id = atq.chat_session_id
|
|
WHERE cs.workspace_id = $1
|
|
AND cs.creator_id = $2
|
|
AND atq.status IN ('queued', 'dispatched', 'running', 'waiting_local_directory')
|
|
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"`
|
|
}
|
|
|
|
// 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.
|
|
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); 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 markChatSessionRead = `-- name: MarkChatSessionRead :exec
|
|
UPDATE chat_session SET unread_since = NULL
|
|
WHERE id = $1
|
|
`
|
|
|
|
// Clears unread_since, 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 setUnreadSinceIfNull = `-- name: SetUnreadSinceIfNull :exec
|
|
UPDATE chat_session SET unread_since = now()
|
|
WHERE id = $1 AND unread_since IS NULL
|
|
`
|
|
|
|
// Atomically stamps the first unread assistant message's arrival time.
|
|
// No-op if the session is already in "has unread" state — keeps the earliest
|
|
// unread boundary stable across multiple incoming replies.
|
|
func (q *Queries) SetUnreadSinceIfNull(ctx context.Context, id pgtype.UUID) error {
|
|
_, err := q.db.Exec(ctx, setUnreadSinceIfNull, id)
|
|
return 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 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
|
|
`
|
|
|
|
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,
|
|
)
|
|
return i, err
|
|
}
|