Files
Memoh/internal/db/sqlc/messages.sql.go
T
Acbox Liu 7d7d0e4b51 refactor: introduce multi-session chat support (#session) (#267)
* refactor: introduce multi-session chat support (#session)

Replace the single-context-per-bot model with multiple chat sessions.

Database:
- Add bot_sessions table (route_id, channel_type, title, metadata, soft delete)
- Migrate bot_history_messages from (route_id, channel_type) to session_id
- Add active_session_id to bot_channel_routes
- Migration 0036 handles data migration from existing messages

Backend:
- New internal/session service for session CRUD
- Update message service/types to use session_id instead of route_id
- Update conversation flow (resolver, history, store) for session context
- Channel inbound auto-creates/retrieves active session via SessionEnsurer
- New REST endpoints: /bots/:bot_id/sessions (CRUD)
- WebSocket and message handlers accept optional session_id
- Wire session service into FX dependency graph (agent + memoh)

Frontend:
- Refactor chat store: sessions replaces chats, sessionId replaces chatId
- Session-aware message loading, sending, and pagination
- WebSocket sends include session_id
- New session sidebar component with select/delete
- Chat area header shows active session title + new session button
- API layer updated: fetchSessions, createSession, deleteSession
- i18n strings for session management (en + zh)

SDK:
- Regenerated TypeScript SDK and Swagger docs with session endpoints

* fix: update tests for session refactoring (RouteID → SessionID)

Remove references to removed RouteID and Platform fields from
PersistInput/Message in channel_test.go and service_integration_test.go.

* fix: restore accidentally deleted SDK files and guard migration 0032

- Restore packages/sdk/src/container-stream.ts and extra/index.ts that
  were accidentally removed during SDK regeneration
- Wrap migration 0032 route_id index creation in a column existence check
  to avoid failure on fresh databases where 0001_init.up.sql no longer
  has route_id

* fix: guard migration 0036 data steps for fresh databases

Wrap steps 3-7 (which reference route_id/channel_type on
bot_history_messages) in a column existence check so the migration
is safe on fresh databases where 0001_init.up.sql already reflects
the final schema without those columns.

* feat: add title model setting and auto-generate session titles on user input

- Add title_model_id to bots table (migration 0037) and bot settings API
- Implement async title generation triggered at user message time (not after
  assistant response) for faster title availability
- Publish session_title_updated events via SSE event hub for real-time
  frontend updates without page refresh
- Fix SSE message event parsing: use direct JSON.parse instead of
  normalizeStreamEvent which silently dropped non-chat-stream event types
- Add title model selector in bot settings UI with i18n support

* fix: session-scoped message filtering and URL-based chat routing

- Filter realtime SSE messages by session_id to prevent cross-session
  message leakage after page refresh
- Add /chat/:sessionId? route with bidirectional URL ↔ store sync
- Visiting /chat shows a clean state with no bot or session pre-selected
- Visiting /chat/:sessionId loads the specific session directly
- Session switches from sidebar automatically update the URL
- Fix stale RouteID field in dedupe test (removed during session refactor)

* fix: skip cross-channel stream events to prevent session leakage

The bot-level web stream pushes events from all channels (Telegram,
Discord, etc.) without session_id context. Previously these were
rendered inline in the current chat view regardless of session.

Now cross-channel events are ignored in handleLocalStreamEvent;
persisted messages arrive via the SSE message events stream with
proper session_id filtering through appendRealtimeMessage.

* feat: show IM avatars and platform badges on session sidebar

- Add sender_avatar_url to route metadata from identity resolution
- Resolve group avatar and handle via directory adapter for group chats
- JOIN bot_channel_routes in ListSessionsByBot to return route metadata
- Display avatar with ChannelBadge on IM session items (group avatar
  for groups, sender avatar for private chats)
- Show @groupname or @username as session sub-label

* fix: clean up RunConfig unused fields, fix skill system and copy bug

- Remove unused RunConfig fields: Tools, Channels, CurrentChannel,
  ActiveContextTime
- Remove unused SessionContext fields: DisplayName, ConversationType
- Fix EnabledSkillNames copy bug: make([]string, 0, n) + copy copies
  zero elements; changed to make([]string, n)
- Fix prepareRunConfig dead code: remove no-op loop over
  CurrentPlatform runes; compute supportsImageInput from model's
  InputModalities
- Fix EnabledSkills always nil in system prompt: resolve enabled skill
  entries from EnabledSkillNames + Skills
- Fix use_skill tool returning empty response: now returns full skill
  content (description + instructions) so LLM gets it in the same turn
- Skip use_skill tool registration when no skills are available
- Conditionally render Skills section in system prompt (hidden when
  no skills exist)

* feat: add session type field and bind sessions to heartbeat/schedule executions

- Add `type` column to `bot_sessions` (chat | heartbeat | schedule)
- Add `session_id` to `bot_heartbeat_logs` for per-execution session tracking
- Create `schedule_logs` table binding schedule_id + session_id
- Heartbeat and schedule runs now create independent sessions and persist
  agent messages via storeRound, enabling full conversation replay
- Add schedule logs API endpoints (list by bot, list by schedule, delete)
- Update Triggerer interfaces to return TriggerResult with status/usage/model

* refactor: modular system prompts per session type (chat/heartbeat/schedule)

Split the monolithic system.md into three type-specific system prompts
with shared fragments via {{include:_xxx}} syntax, so each session type
gets a focused prompt without irrelevant instructions.

* fix: prevent message duplication after task completion

message_created events from Persist() had an empty platform field because
toMessageFromCreate() didn't extract it from the session. This caused
appendRealtimeMessage to fail the platform === 'web' guard, and
hasMessageWithId to fail because local IDs differ from server UUIDs,
resulting in all messages being appended as duplicates.

- Extract platform from metadata in toMessageFromCreate so published events
  carry the correct value
- Pass channel_type: 'web' when creating sessions from the web frontend so
  List queries return the correct platform via the session JOIN

* fix: use per-message usage from SDK instead of misaligned step-level usages

Previously, token usage was stored via a separate per-step usages array
that didn't align with messages (off-by-one from prepending user message,
step count != message count). This caused:
- User messages incorrectly receiving usage data
- Usage values shifted across messages in multi-step rounds
- Last assistant message getting the accumulated total instead of its own step usage
- InputTokenDetails/OutputTokenDetails lost during manual accumulation

Now each sdk.Message carries its own per-step Usage (set by the SDK in
buildStepMessages), which is extracted in sdkMessagesToModelMessages and
stored directly via ModelMessage.Usage. The storeRound/storeMessages path
no longer needs external usage/usages parameters.

Also fixes the totalUsage accumulation in runStream to include all detail
fields (InputTokenDetails, OutputTokenDetails).

* feat: add /new slash command to create a new active session from IM channels

Users in Telegram/Discord/Feishu can now send /new to start a fresh
conversation, resetting the session context for the current chat thread.
The command resolves the channel route, creates a new session, sets it as
the active session on the route, and replies with a confirmation message.

* feat: distinguish heartbeat and schedule sessions with dedicated icons in sidebar

Heartbeat sessions show a heart-pulse icon (rose), schedule sessions
show a clock icon (amber), and both display a type label beneath the
session title.

* refactor: remove enabledSkills system prompt injection, keep sorted skill listing

use_skill now returns skill content directly as tool output, so there is
no need to inject enabled skill body text into the system prompt. Remove
the entire enabledSkills tracking chain (RunConfig.EnabledSkillNames,
StreamEvent.Skills, GenerateResult.Skills, ChatRequest/Response.Skills,
enableSkill closures in runStream/runGenerate, prepareRunConfig matching).

Keep a lightweight skills listing (name + description only) in the system
prompt so the model knows which skills are available. Sort entries by name
to guarantee deterministic ordering and maximize KV cache reuse.

* refactor: remove inbox system, persist passive messages directly to history

Replace the bot_inbox table and service with direct writes to
bot_history_messages for group conversations where the bot is not
@mentioned. Trigger-path messages continue to be persisted after the
agent responds (unchanged).

- Drop bot_inbox table and max_inbox_items column (migration 0039)
- Delete internal/inbox/, handlers/inbox.go, command/inbox.go,
  agent/tools/inbox.go and the MCP message provider
- Add persistPassiveMessage() in channel inbound to write user
  messages into the active session immediately
- Rewrite ListObservedConversationsByChannelIdentity to query
  bot_history_messages + bot_sessions instead of bot_inbox
- Extract shared send/react logic into internal/messaging/executor.go;
  agent/tools/message.go is now a thin SDK adapter
- Clean up all inbox references from agent prompts, flow resolver,
  email trigger, settings, commands, DI wiring, and frontend
- Regenerate sqlc, swagger, and SDK

* feat: add list_sessions and search_messages agent tools

Provide agents with the ability to query session metadata and search
message history across all sessions. search_messages supports filtering
by time range, keyword (JSONB-aware ILIKE), session, contact, and role,
with a default 7-day lookback when no start_time is given.

* feat: inject last_heartbeat time and improve heartbeat search guidance

Query the previous heartbeat's started_at timestamp and pass it through
TriggerPayload into the heartbeat prompt template. Update system prompt
and HEARTBEAT.md checklist to guide agents to use search_messages with
start_time=last_heartbeat for efficient cross-session message review.

* fix: pass BridgeProvider to FSClient and store full heartbeat prompt

FSClient was always created with nil provider, causing all container
file reads (IDENTITY.md, SOUL.md, MEMORY.md, HEARTBEAT.md, etc.) to
silently return empty strings. Expose Agent.BridgeProvider() and wire
it into Resolver. Also fix heartbeat trigger to store the full prompt
template as the user message instead of the literal "heartbeat" string.

* feat: add line numbers to container file read output

Move line-number formatting from the bridge gRPC server to the agent
tool layer so that the raw content stored and transmitted via gRPC
remains clean, while the read_file tool output includes numbered lines
for easier reference by the agent.

* chore(deps): update twilight-ai to v0.3.2

* fix: lint, test
2026-03-21 15:57:22 +08:00

1149 lines
35 KiB
Go

// Code generated by sqlc. DO NOT EDIT.
// versions:
// sqlc v1.30.0
// source: messages.sql
package sqlc
import (
"context"
"github.com/jackc/pgx/v5/pgtype"
)
const createMessage = `-- name: CreateMessage :one
INSERT INTO bot_history_messages (
bot_id,
session_id,
sender_channel_identity_id,
sender_account_user_id,
source_message_id,
source_reply_to_message_id,
role,
content,
metadata,
usage,
model_id
)
VALUES (
$1,
$2::uuid,
$3::uuid,
$4::uuid,
$5::text,
$6::text,
$7,
$8,
$9,
$10,
$11::uuid
)
RETURNING
id,
bot_id,
session_id,
sender_channel_identity_id,
sender_account_user_id AS sender_user_id,
source_message_id AS external_message_id,
source_reply_to_message_id,
role,
content,
metadata,
usage,
created_at
`
type CreateMessageParams struct {
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
ModelID pgtype.UUID `json:"model_id"`
}
type CreateMessageRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
}
func (q *Queries) CreateMessage(ctx context.Context, arg CreateMessageParams) (CreateMessageRow, error) {
row := q.db.QueryRow(ctx, createMessage,
arg.BotID,
arg.SessionID,
arg.SenderChannelIdentityID,
arg.SenderUserID,
arg.ExternalMessageID,
arg.SourceReplyToMessageID,
arg.Role,
arg.Content,
arg.Metadata,
arg.Usage,
arg.ModelID,
)
var i CreateMessageRow
err := row.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
)
return i, err
}
const deleteMessagesByBot = `-- name: DeleteMessagesByBot :exec
DELETE FROM bot_history_messages
WHERE bot_id = $1
`
func (q *Queries) DeleteMessagesByBot(ctx context.Context, botID pgtype.UUID) error {
_, err := q.db.Exec(ctx, deleteMessagesByBot, botID)
return err
}
const deleteMessagesBySession = `-- name: DeleteMessagesBySession :exec
DELETE FROM bot_history_messages
WHERE session_id = $1
`
func (q *Queries) DeleteMessagesBySession(ctx context.Context, sessionID pgtype.UUID) error {
_, err := q.db.Exec(ctx, deleteMessagesBySession, sessionID)
return err
}
const listActiveMessagesSince = `-- name: ListActiveMessagesSince :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.bot_id = $1
AND m.created_at >= $2
AND (m.metadata->>'trigger_mode' IS NULL OR m.metadata->>'trigger_mode' != 'passive_sync')
ORDER BY m.created_at ASC
`
type ListActiveMessagesSinceParams struct {
BotID pgtype.UUID `json:"bot_id"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
}
type ListActiveMessagesSinceRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListActiveMessagesSince(ctx context.Context, arg ListActiveMessagesSinceParams) ([]ListActiveMessagesSinceRow, error) {
rows, err := q.db.Query(ctx, listActiveMessagesSince, arg.BotID, arg.CreatedAt)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListActiveMessagesSinceRow
for rows.Next() {
var i ListActiveMessagesSinceRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listActiveMessagesSinceBySession = `-- name: ListActiveMessagesSinceBySession :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.session_id = $1
AND m.created_at >= $2
AND (m.metadata->>'trigger_mode' IS NULL OR m.metadata->>'trigger_mode' != 'passive_sync')
ORDER BY m.created_at ASC
`
type ListActiveMessagesSinceBySessionParams struct {
SessionID pgtype.UUID `json:"session_id"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
}
type ListActiveMessagesSinceBySessionRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListActiveMessagesSinceBySession(ctx context.Context, arg ListActiveMessagesSinceBySessionParams) ([]ListActiveMessagesSinceBySessionRow, error) {
rows, err := q.db.Query(ctx, listActiveMessagesSinceBySession, arg.SessionID, arg.CreatedAt)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListActiveMessagesSinceBySessionRow
for rows.Next() {
var i ListActiveMessagesSinceBySessionRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessages = `-- name: ListMessages :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.bot_id = $1
ORDER BY m.created_at ASC
LIMIT 10000
`
type ListMessagesRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessages(ctx context.Context, botID pgtype.UUID) ([]ListMessagesRow, error) {
rows, err := q.db.Query(ctx, listMessages, botID)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesRow
for rows.Next() {
var i ListMessagesRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessagesBefore = `-- name: ListMessagesBefore :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.bot_id = $1
AND m.created_at < $2
ORDER BY m.created_at DESC
LIMIT $3
`
type ListMessagesBeforeParams struct {
BotID pgtype.UUID `json:"bot_id"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
MaxCount int32 `json:"max_count"`
}
type ListMessagesBeforeRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessagesBefore(ctx context.Context, arg ListMessagesBeforeParams) ([]ListMessagesBeforeRow, error) {
rows, err := q.db.Query(ctx, listMessagesBefore, arg.BotID, arg.CreatedAt, arg.MaxCount)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesBeforeRow
for rows.Next() {
var i ListMessagesBeforeRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessagesBeforeBySession = `-- name: ListMessagesBeforeBySession :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.session_id = $1
AND m.created_at < $2
ORDER BY m.created_at DESC
LIMIT $3
`
type ListMessagesBeforeBySessionParams struct {
SessionID pgtype.UUID `json:"session_id"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
MaxCount int32 `json:"max_count"`
}
type ListMessagesBeforeBySessionRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessagesBeforeBySession(ctx context.Context, arg ListMessagesBeforeBySessionParams) ([]ListMessagesBeforeBySessionRow, error) {
rows, err := q.db.Query(ctx, listMessagesBeforeBySession, arg.SessionID, arg.CreatedAt, arg.MaxCount)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesBeforeBySessionRow
for rows.Next() {
var i ListMessagesBeforeBySessionRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessagesBySession = `-- name: ListMessagesBySession :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.session_id = $1
ORDER BY m.created_at ASC
LIMIT 10000
`
type ListMessagesBySessionRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessagesBySession(ctx context.Context, sessionID pgtype.UUID) ([]ListMessagesBySessionRow, error) {
rows, err := q.db.Query(ctx, listMessagesBySession, sessionID)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesBySessionRow
for rows.Next() {
var i ListMessagesBySessionRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessagesLatest = `-- name: ListMessagesLatest :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.bot_id = $1
ORDER BY m.created_at DESC
LIMIT $2
`
type ListMessagesLatestParams struct {
BotID pgtype.UUID `json:"bot_id"`
MaxCount int32 `json:"max_count"`
}
type ListMessagesLatestRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessagesLatest(ctx context.Context, arg ListMessagesLatestParams) ([]ListMessagesLatestRow, error) {
rows, err := q.db.Query(ctx, listMessagesLatest, arg.BotID, arg.MaxCount)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesLatestRow
for rows.Next() {
var i ListMessagesLatestRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessagesLatestBySession = `-- name: ListMessagesLatestBySession :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.session_id = $1
ORDER BY m.created_at DESC
LIMIT $2
`
type ListMessagesLatestBySessionParams struct {
SessionID pgtype.UUID `json:"session_id"`
MaxCount int32 `json:"max_count"`
}
type ListMessagesLatestBySessionRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessagesLatestBySession(ctx context.Context, arg ListMessagesLatestBySessionParams) ([]ListMessagesLatestBySessionRow, error) {
rows, err := q.db.Query(ctx, listMessagesLatestBySession, arg.SessionID, arg.MaxCount)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesLatestBySessionRow
for rows.Next() {
var i ListMessagesLatestBySessionRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessagesSince = `-- name: ListMessagesSince :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.bot_id = $1
AND m.created_at >= $2
ORDER BY m.created_at ASC
`
type ListMessagesSinceParams struct {
BotID pgtype.UUID `json:"bot_id"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
}
type ListMessagesSinceRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessagesSince(ctx context.Context, arg ListMessagesSinceParams) ([]ListMessagesSinceRow, error) {
rows, err := q.db.Query(ctx, listMessagesSince, arg.BotID, arg.CreatedAt)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesSinceRow
for rows.Next() {
var i ListMessagesSinceRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listMessagesSinceBySession = `-- name: ListMessagesSinceBySession :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.sender_account_user_id AS sender_user_id,
m.source_message_id AS external_message_id,
m.source_reply_to_message_id,
m.role,
m.content,
m.metadata,
m.usage,
m.created_at,
ci.display_name AS sender_display_name,
ci.avatar_url AS sender_avatar_url,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.session_id = $1
AND m.created_at >= $2
ORDER BY m.created_at ASC
`
type ListMessagesSinceBySessionParams struct {
SessionID pgtype.UUID `json:"session_id"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
}
type ListMessagesSinceBySessionRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
SenderUserID pgtype.UUID `json:"sender_user_id"`
ExternalMessageID pgtype.Text `json:"external_message_id"`
SourceReplyToMessageID pgtype.Text `json:"source_reply_to_message_id"`
Role string `json:"role"`
Content []byte `json:"content"`
Metadata []byte `json:"metadata"`
Usage []byte `json:"usage"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
SenderAvatarUrl pgtype.Text `json:"sender_avatar_url"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) ListMessagesSinceBySession(ctx context.Context, arg ListMessagesSinceBySessionParams) ([]ListMessagesSinceBySessionRow, error) {
rows, err := q.db.Query(ctx, listMessagesSinceBySession, arg.SessionID, arg.CreatedAt)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListMessagesSinceBySessionRow
for rows.Next() {
var i ListMessagesSinceBySessionRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.SenderUserID,
&i.ExternalMessageID,
&i.SourceReplyToMessageID,
&i.Role,
&i.Content,
&i.Metadata,
&i.Usage,
&i.CreatedAt,
&i.SenderDisplayName,
&i.SenderAvatarUrl,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const listObservedConversationsByChannelIdentity = `-- name: ListObservedConversationsByChannelIdentity :many
WITH observed_routes AS (
SELECT
s.route_id,
MAX(m.created_at)::timestamptz AS last_observed_at
FROM bot_history_messages m
JOIN bot_sessions s ON s.id = m.session_id
WHERE m.bot_id = $1
AND m.sender_channel_identity_id = $2::uuid
AND s.route_id IS NOT NULL
GROUP BY s.route_id
)
SELECT
r.id AS route_id,
r.channel_type AS channel,
CASE
WHEN LOWER(COALESCE(r.conversation_type, '')) IN ('thread', 'topic') THEN 'thread'
ELSE 'group'
END AS conversation_type,
r.external_conversation_id AS conversation_id,
COALESCE(r.external_thread_id, '') AS thread_id,
COALESCE(r.metadata->>'conversation_name', '')::text AS conversation_name,
rr.last_observed_at
FROM observed_routes rr
JOIN bot_channel_routes r ON r.id = rr.route_id
WHERE LOWER(COALESCE(r.conversation_type, '')) NOT IN ('', 'p2p', 'private', 'direct', 'dm')
GROUP BY
r.id,
r.channel_type,
r.conversation_type,
r.external_conversation_id,
r.external_thread_id,
r.metadata,
rr.last_observed_at
ORDER BY rr.last_observed_at DESC
`
type ListObservedConversationsByChannelIdentityParams struct {
BotID pgtype.UUID `json:"bot_id"`
ChannelIdentityID pgtype.UUID `json:"channel_identity_id"`
}
type ListObservedConversationsByChannelIdentityRow struct {
RouteID pgtype.UUID `json:"route_id"`
Channel string `json:"channel"`
ConversationType string `json:"conversation_type"`
ConversationID string `json:"conversation_id"`
ThreadID string `json:"thread_id"`
ConversationName string `json:"conversation_name"`
LastObservedAt pgtype.Timestamptz `json:"last_observed_at"`
}
func (q *Queries) ListObservedConversationsByChannelIdentity(ctx context.Context, arg ListObservedConversationsByChannelIdentityParams) ([]ListObservedConversationsByChannelIdentityRow, error) {
rows, err := q.db.Query(ctx, listObservedConversationsByChannelIdentity, arg.BotID, arg.ChannelIdentityID)
if err != nil {
return nil, err
}
defer rows.Close()
var items []ListObservedConversationsByChannelIdentityRow
for rows.Next() {
var i ListObservedConversationsByChannelIdentityRow
if err := rows.Scan(
&i.RouteID,
&i.Channel,
&i.ConversationType,
&i.ConversationID,
&i.ThreadID,
&i.ConversationName,
&i.LastObservedAt,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}
const searchMessages = `-- name: SearchMessages :many
SELECT
m.id,
m.bot_id,
m.session_id,
m.sender_channel_identity_id,
m.role,
m.content,
m.created_at,
ci.display_name AS sender_display_name,
s.channel_type AS platform
FROM bot_history_messages m
LEFT JOIN channel_identities ci ON ci.id = m.sender_channel_identity_id
LEFT JOIN bot_sessions s ON s.id = m.session_id
WHERE m.bot_id = $1
AND ($2::uuid IS NULL OR m.session_id = $2::uuid)
AND ($3::uuid IS NULL OR m.sender_channel_identity_id = $3::uuid)
AND ($4::timestamptz IS NULL OR m.created_at >= $4::timestamptz)
AND ($5::timestamptz IS NULL OR m.created_at <= $5::timestamptz)
AND ($6::text IS NULL OR m.role = $6::text)
AND ($7::text IS NULL OR (
CASE
WHEN jsonb_typeof(m.content->'content') = 'string'
THEN m.content->>'content'
WHEN jsonb_typeof(m.content->'content') = 'array'
THEN (SELECT COALESCE(string_agg(elem->>'text', ' '), '')
FROM jsonb_array_elements(m.content->'content') AS elem
WHERE elem->>'type' = 'text')
ELSE ''
END
) ILIKE '%' || $7::text || '%')
ORDER BY m.created_at DESC
LIMIT $8
`
type SearchMessagesParams struct {
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
ContactID pgtype.UUID `json:"contact_id"`
StartTime pgtype.Timestamptz `json:"start_time"`
EndTime pgtype.Timestamptz `json:"end_time"`
Role pgtype.Text `json:"role"`
Keyword pgtype.Text `json:"keyword"`
MaxCount int32 `json:"max_count"`
}
type SearchMessagesRow struct {
ID pgtype.UUID `json:"id"`
BotID pgtype.UUID `json:"bot_id"`
SessionID pgtype.UUID `json:"session_id"`
SenderChannelIdentityID pgtype.UUID `json:"sender_channel_identity_id"`
Role string `json:"role"`
Content []byte `json:"content"`
CreatedAt pgtype.Timestamptz `json:"created_at"`
SenderDisplayName pgtype.Text `json:"sender_display_name"`
Platform pgtype.Text `json:"platform"`
}
func (q *Queries) SearchMessages(ctx context.Context, arg SearchMessagesParams) ([]SearchMessagesRow, error) {
rows, err := q.db.Query(ctx, searchMessages,
arg.BotID,
arg.SessionID,
arg.ContactID,
arg.StartTime,
arg.EndTime,
arg.Role,
arg.Keyword,
arg.MaxCount,
)
if err != nil {
return nil, err
}
defer rows.Close()
var items []SearchMessagesRow
for rows.Next() {
var i SearchMessagesRow
if err := rows.Scan(
&i.ID,
&i.BotID,
&i.SessionID,
&i.SenderChannelIdentityID,
&i.Role,
&i.Content,
&i.CreatedAt,
&i.SenderDisplayName,
&i.Platform,
); err != nil {
return nil, err
}
items = append(items, i)
}
if err := rows.Err(); err != nil {
return nil, err
}
return items, nil
}