Documentation
¶
Index ¶
- Constants
- Variables
- func LatestInterruptedCheckpoint(n int, classify func(int) (assistant, interrupted bool)) int
- type ACPPublication
- type ActiveMessagesMeasure
- type AgentStep
- type AgentStepPersister
- type AssetRef
- type AtomicRoundPersister
- type DBService
- func (s *DBService) DeleteByBot(ctx context.Context, botID string) error
- func (s *DBService) DeleteByIDs(ctx context.Context, ids []string) error
- func (s *DBService) DeleteBySession(ctx context.Context, sessionID string) error
- func (s *DBService) GetByIDBySession(ctx context.Context, sessionID string, messageID string) (Message, error)
- func (s *DBService) GetLatestVisibleTurnBySession(ctx context.Context, sessionID string) (HistoryTurn, error)
- func (s *DBService) GetVisibleTurnByMessage(ctx context.Context, sessionID string, messageID string) (HistoryTurn, error)
- func (s *DBService) LinkAssets(ctx context.Context, messageID string, assets []AssetRef) error
- func (s *DBService) List(ctx context.Context, botID string) ([]Message, error)
- func (s *DBService) ListActiveSince(ctx context.Context, botID string, since time.Time) ([]Message, error)
- func (s *DBService) ListActiveSinceBySession(ctx context.Context, sessionID string, since time.Time) ([]Message, error)
- func (s *DBService) ListActiveSinceBySessionWithinBytes(ctx context.Context, sessionID string, since time.Time, maxBytes int64) ([]Message, error)
- func (s *DBService) ListActiveSinceWithinBytes(ctx context.Context, botID string, since time.Time, maxBytes int64) ([]Message, error)
- func (s *DBService) ListBefore(ctx context.Context, botID string, before time.Time, limit int32) ([]Message, error)
- func (s *DBService) ListBeforeBySession(ctx context.Context, sessionID string, before time.Time, limit int32) ([]Message, error)
- func (s *DBService) ListBeforeMessageBySession(ctx context.Context, sessionID string, beforeMessageID string, limit int32) ([]Message, error)
- func (s *DBService) ListBySession(ctx context.Context, sessionID string) ([]Message, error)
- func (s *DBService) ListLatest(ctx context.Context, botID string, limit int32) ([]Message, error)
- func (s *DBService) ListLatestBySession(ctx context.Context, sessionID string, limit int32) ([]Message, error)
- func (s *DBService) ListLatestUIBySession(ctx context.Context, sessionID string, limit int32) ([]Message, error)
- func (s *DBService) ListSince(ctx context.Context, botID string, since time.Time) ([]Message, error)
- func (s *DBService) ListSinceBySession(ctx context.Context, sessionID string, since time.Time) ([]Message, error)
- func (s *DBService) ListVisibleFromBySession(ctx context.Context, sessionID string, messageID string) ([]Message, error)
- func (s *DBService) LocateByExternalIDBySession(ctx context.Context, sessionID string, externalMessageID string, ...) (LocateResult, error)
- func (s *DBService) MeasureActiveBySession(ctx context.Context, sessionID string, since time.Time) (ActiveMessagesMeasure, error)
- func (s *DBService) Persist(ctx context.Context, input PersistInput) (Message, error)
- func (s *DBService) PersistAgentStep(ctx context.Context, step AgentStep) ([]Message, error)
- func (s *DBService) PersistRound(ctx context.Context, inputs []PersistInput, options RoundPersistenceOptions) ([]Message, bool, error)
- func (s *DBService) PersistToolTailRound(ctx context.Context, inputs []PersistInput) ([]Message, bool, error)
- func (s *DBService) ReplaceTurn(ctx context.Context, sessionID string, oldTurnID string, ...) (HistoryTurn, error)
- type HistoryTurn
- type LocateResult
- type Message
- type MessageAsset
- type PersistInput
- type ReasoningTimingMetadata
- type ReasoningTimingSegment
- type RoundPersistenceOptions
- type Service
- type ToolTailRoundPersister
- type TurnReplacement
- type Writer
Constants ¶
const ( // ReasoningTimingMetadataKey stores server-observed reasoning segment timing // on the assistant row that owns the corresponding reasoning content. ReasoningTimingMetadataKey = "reasoning_timing" ReasoningTimingVersion = 1 )
const ( AgentStepInterruptedMetadataKey = "agent_step_interrupted" HistoryErrorCodeMetadataKey = "error_code" AgentStepInterruptedReasoningPrefix = "[Previous assistant response was interrupted during reasoning. Continue from this checkpoint:]\n" )
Variables ¶
var ErrAgentStepNotWritable = errors.New("agent step is no longer writable")
Functions ¶
Types ¶
type ACPPublication ¶
ACPPublication moves the session's canonical ACP publication head to the round's run inside the same transaction as the round's messages. A head with CheckpointReset=true is canonical but not resumable.
type ActiveMessagesMeasure ¶
ActiveMessagesMeasure is a database-side aggregate over a session's active messages (CM-ADM-001): counts and raw content bytes, computed without loading any payload into the process.
type AgentStep ¶
type AgentStep struct {
RunID string
Messages []PersistInput
Interrupted bool
}
AgentStep is one native-agent step persisted atomically. Interrupted marks a text/reasoning snapshot of a step the model never finished, which takes a different writability predicate than a complete step.
type AgentStepPersister ¶
type AssetRef ¶
type AssetRef struct {
ContentHash string `json:"content_hash"`
Role string `json:"role"`
Ordinal int `json:"ordinal"`
Mime string `json:"mime,omitempty"`
SizeBytes int64 `json:"size_bytes,omitempty"`
StorageKey string `json:"storage_key,omitempty"`
Name string `json:"name,omitempty"`
Metadata map[string]any `json:"metadata,omitempty"`
}
AssetRef links a media asset to a persisted message. ContentHash is the content-addressed identifier for the media file.
type AtomicRoundPersister ¶
type AtomicRoundPersister interface {
PersistRound(ctx context.Context, inputs []PersistInput, options RoundPersistenceOptions) ([]Message, bool, error)
}
AtomicRoundPersister writes a complete round in one transaction. Implementations must enforce any runtime fence carried by ctx, while still supporting unfenced local replacement transactions.
type DBService ¶
type DBService struct {
// contains filtered or unexported fields
}
DBService persists and reads bot history messages.
func NewService ¶
func NewService(log *slog.Logger, queries dbstore.Queries, publishers ...event.Publisher) *DBService
NewService creates a message service.
func (*DBService) DeleteByBot ¶
DeleteByBot deletes all messages for a bot.
func (*DBService) DeleteByIDs ¶
DeleteByIDs deletes specific messages by id.
func (*DBService) DeleteBySession ¶
DeleteBySession deletes all messages for a session.
func (*DBService) GetByIDBySession ¶
func (*DBService) GetLatestVisibleTurnBySession ¶
func (*DBService) GetVisibleTurnByMessage ¶
func (*DBService) LinkAssets ¶
LinkAssets links asset refs to an existing persisted message.
func (*DBService) ListActiveSince ¶
func (s *DBService) ListActiveSince(ctx context.Context, botID string, since time.Time) ([]Message, error)
ListActiveSince returns bot messages since a given time, excluding passive_sync messages.
func (*DBService) ListActiveSinceBySession ¶
func (s *DBService) ListActiveSinceBySession(ctx context.Context, sessionID string, since time.Time) ([]Message, error)
ListActiveSinceBySession returns session messages since a given time, excluding passive_sync messages.
func (*DBService) ListActiveSinceBySessionWithinBytes ¶
func (s *DBService) ListActiveSinceBySessionWithinBytes(ctx context.Context, sessionID string, since time.Time, maxBytes int64) ([]Message, error)
ListActiveSinceBySessionWithinBytes returns session messages since a given time, admitted newest-first within the content byte budget (CM-ADM-001).
func (*DBService) ListActiveSinceWithinBytes ¶
func (s *DBService) ListActiveSinceWithinBytes(ctx context.Context, botID string, since time.Time, maxBytes int64) ([]Message, error)
ListActiveSinceWithinBytes returns bot messages since a given time, admitted newest-first within the content byte budget (CM-ADM-001).
func (*DBService) ListBefore ¶
func (s *DBService) ListBefore(ctx context.Context, botID string, before time.Time, limit int32) ([]Message, error)
ListBefore returns up to limit messages older than before (created_at < before), ordered oldest-first.
func (*DBService) ListBeforeBySession ¶
func (s *DBService) ListBeforeBySession(ctx context.Context, sessionID string, before time.Time, limit int32) ([]Message, error)
ListBeforeBySession returns up to limit session messages older than before.
func (*DBService) ListBeforeMessageBySession ¶
func (s *DBService) ListBeforeMessageBySession(ctx context.Context, sessionID string, beforeMessageID string, limit int32) ([]Message, error)
ListBeforeMessageBySession returns up to limit session messages before the cursor message in visible turn order.
func (*DBService) ListBySession ¶
ListBySession returns all messages for a session.
func (*DBService) ListLatest ¶
ListLatest returns the latest N bot messages (newest first in DB; caller may reverse for ASC).
func (*DBService) ListLatestBySession ¶
func (s *DBService) ListLatestBySession(ctx context.Context, sessionID string, limit int32) ([]Message, error)
ListLatestBySession returns the latest N session messages.
func (*DBService) ListLatestUIBySession ¶
func (s *DBService) ListLatestUIBySession(ctx context.Context, sessionID string, limit int32) ([]Message, error)
ListLatestUIBySession returns the latest N session messages using the lighter field set needed by format=ui rendering.
func (*DBService) ListSince ¶
func (s *DBService) ListSince(ctx context.Context, botID string, since time.Time) ([]Message, error)
ListSince returns bot messages since a given time.
func (*DBService) ListSinceBySession ¶
func (s *DBService) ListSinceBySession(ctx context.Context, sessionID string, since time.Time) ([]Message, error)
ListSinceBySession returns session messages since a given time.
func (*DBService) ListVisibleFromBySession ¶
func (*DBService) LocateByExternalIDBySession ¶
func (*DBService) MeasureActiveBySession ¶
func (s *DBService) MeasureActiveBySession(ctx context.Context, sessionID string, since time.Time) (ActiveMessagesMeasure, error)
MeasureActiveBySession aggregates a session's active message count and content bytes on the database side (CM-ADM-001).
func (*DBService) PersistAgentStep ¶
PersistAgentStep appends one SDK step in a runtime-fenced transaction. Complete steps precede abort intent; interrupted checkpoints remain writable until terminal finalization for cancellation paths without recorded intent.
func (*DBService) PersistRound ¶
func (s *DBService) PersistRound(ctx context.Context, inputs []PersistInput, options RoundPersistenceOptions) ([]Message, bool, error)
PersistRound writes all messages and history links under one PostgreSQL transaction. Distributed callers additionally validate their runtime fence in that transaction; local replacements use the same atomic write without a distributed ownership token.
func (*DBService) PersistToolTailRound ¶
func (s *DBService) PersistToolTailRound(ctx context.Context, inputs []PersistInput) ([]Message, bool, error)
PersistToolTailRound writes the common user -> assistant(tool-call) -> tool -> assistant(final) round with one PostgreSQL statement. Unsupported stores or non-matching inputs return handled=false so callers can use Persist.
type HistoryTurn ¶
type HistoryTurn struct {
ID string `json:"id"`
BotID string `json:"bot_id"`
SessionID string `json:"session_id"`
Position int64 `json:"position"`
RequestMessageID string `json:"request_message_id,omitempty"`
AssistantMessageID string `json:"assistant_message_id,omitempty"`
SupersededByTurnID string `json:"superseded_by_turn_id,omitempty"`
SupersededAt time.Time `json:"superseded_at,omitempty"`
SupersededReason string `json:"superseded_reason,omitempty"`
CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
}
type LocateResult ¶
type Message ¶
type Message struct {
ID string `json:"id"`
BotID string `json:"bot_id"`
SessionID string `json:"session_id,omitempty"`
SenderChannelIdentityID string `json:"sender_channel_identity_id,omitempty"`
SenderUserID string `json:"sender_user_id,omitempty"`
SenderDisplayName string `json:"sender_display_name,omitempty"`
SenderAvatarURL string `json:"sender_avatar_url,omitempty"`
Platform string `json:"platform,omitempty"`
ExternalMessageID string `json:"external_message_id,omitempty"`
SourceReplyToMessageID string `json:"source_reply_to_message_id,omitempty"`
Role string `json:"role"`
Content json.RawMessage `json:"content"`
Metadata map[string]any `json:"metadata,omitempty"`
RawMetadata json.RawMessage `json:"-"`
TurnID string `json:"turn_id,omitempty"`
// TurnPosition is the immutable turn-level sequence reserved at admission
// (SR-TURN-001). It is loaded on the UI read path so the frontend can
// order and reconcile turns without guessing from text or timestamps.
TurnPosition *int64 `json:"turn_position,omitempty"`
Usage json.RawMessage `json:"usage,omitempty"`
SessionMode string `json:"session_mode,omitempty"`
RuntimeType string `json:"runtime_type,omitempty"`
Assets []MessageAsset `json:"assets,omitempty"`
CompactID string `json:"compact_id,omitempty"`
EventID string `json:"event_id,omitempty"`
DisplayContent string `json:"display_content,omitempty"`
CreatedAt time.Time `json:"created_at"`
}
Message represents a single persisted bot message.
type MessageAsset ¶
type MessageAsset struct {
ContentHash string `json:"content_hash"`
Role string `json:"role"`
Ordinal int `json:"ordinal"`
Mime string `json:"mime"`
SizeBytes int64 `json:"size_bytes"`
StorageKey string `json:"storage_key"`
Name string `json:"name,omitempty"`
Metadata map[string]any `json:"metadata,omitempty"`
}
MessageAsset carries media asset metadata attached to a message. ContentHash is the content-addressed identifier for the media file.
type PersistInput ¶
type PersistInput struct {
BotID string
SessionID string
SenderChannelIdentityID string
SenderUserID string
ExternalMessageID string
SourceReplyToMessageID string
Role string
Content json.RawMessage
Metadata map[string]any
Usage json.RawMessage
SessionMode string
RuntimeType string
Assets []AssetRef
ModelID string
EventID string
DisplayText string
TurnRequestMessageID string
SkipHistoryTurn bool
// RunID links this row to the admitted run that produced it. It is
// traceability only: history is read by turn, never by run.
RunID string
// TurnID and TurnPosition carry a turn this message must be filed under
// because admission already allocated it (SR-TURN-001). Both empty means no
// admission decided the turn and the history layer allocates one, which is
// still the case for channel inbound and schedules.
//
// They travel together: a turn id without its position would file the row
// under the right turn at the wrong place in the session's order, and a
// position without its id would take a slot the client cannot name.
TurnID string
TurnPosition *int64
}
PersistInput is the input for persisting a message.
type ReasoningTimingMetadata ¶
type ReasoningTimingMetadata struct {
Version int `json:"version"`
Segments []ReasoningTimingSegment `json:"segments"`
}
ReasoningTimingMetadata is the versioned persistence envelope stored in a bot_history_messages metadata JSONB value. Keeping it outside model content prevents presentation-only timing from being replayed to providers.
type ReasoningTimingSegment ¶
type ReasoningTimingSegment struct {
Ordinal int `json:"ordinal"`
DurationMS int64 `json:"duration_ms"`
State string `json:"state"`
}
ReasoningTimingSegment describes one visible reasoning block. DurationMS is observed by the server and intentionally excludes presentation-only details.
func ReasoningTimingFromMetadata ¶
func ReasoningTimingFromMetadata(metadata map[string]any) []ReasoningTimingSegment
ReasoningTimingFromMetadata decodes the versioned timing envelope from a message metadata map. JSON round-tripping keeps this tolerant of values loaded from JSONB as map[string]any as well as typed values used in tests.
type RoundPersistenceOptions ¶
type RoundPersistenceOptions struct {
Replacement *TurnReplacement
CleanupACPDecisionProjections bool
ACPPublication *ACPPublication
}
type Service ¶
type Service interface {
Writer
List(ctx context.Context, botID string) ([]Message, error)
ListSince(ctx context.Context, botID string, since time.Time) ([]Message, error)
ListActiveSince(ctx context.Context, botID string, since time.Time) ([]Message, error)
// ListActiveSinceWithinBytes is the byte-budgeted variant of
// ListActiveSince (CM-ADM-001): rows are admitted newest-first until
// their content byte total crosses maxBytes, so process memory is
// bounded by the budget regardless of total history size.
ListActiveSinceWithinBytes(ctx context.Context, botID string, since time.Time, maxBytes int64) ([]Message, error)
ListLatest(ctx context.Context, botID string, limit int32) ([]Message, error)
ListBefore(ctx context.Context, botID string, before time.Time, limit int32) ([]Message, error)
ListBySession(ctx context.Context, sessionID string) ([]Message, error)
ListSinceBySession(ctx context.Context, sessionID string, since time.Time) ([]Message, error)
ListActiveSinceBySession(ctx context.Context, sessionID string, since time.Time) ([]Message, error)
// ListActiveSinceBySessionWithinBytes is the byte-budgeted variant of
// ListActiveSinceBySession (CM-ADM-001), same admission semantics as
// ListActiveSinceWithinBytes scoped to one session.
ListActiveSinceBySessionWithinBytes(ctx context.Context, sessionID string, since time.Time, maxBytes int64) ([]Message, error)
// MeasureActiveBySession aggregates message count and content bytes on
// the database side (CM-ADM-001): admission sizes a session's history
// without shipping any payload into the process.
MeasureActiveBySession(ctx context.Context, sessionID string, since time.Time) (ActiveMessagesMeasure, error)
ListLatestBySession(ctx context.Context, sessionID string, limit int32) ([]Message, error)
ListBeforeBySession(ctx context.Context, sessionID string, before time.Time, limit int32) ([]Message, error)
ListBeforeMessageBySession(ctx context.Context, sessionID string, beforeMessageID string, limit int32) ([]Message, error)
LocateByExternalIDBySession(ctx context.Context, sessionID string, externalMessageID string, beforeLimit int32, afterLimit int32) (LocateResult, error)
GetByIDBySession(ctx context.Context, sessionID string, messageID string) (Message, error)
ListVisibleFromBySession(ctx context.Context, sessionID string, messageID string) ([]Message, error)
GetVisibleTurnByMessage(ctx context.Context, sessionID string, messageID string) (HistoryTurn, error)
GetLatestVisibleTurnBySession(ctx context.Context, sessionID string) (HistoryTurn, error)
ReplaceTurn(ctx context.Context, sessionID string, oldTurnID string, replacementTurnID string, replacementTurnPosition *int64, requestMessageID string, assistantMessageID string, reason string) (HistoryTurn, error)
DeleteByIDs(ctx context.Context, ids []string) error
DeleteByBot(ctx context.Context, botID string) error
DeleteBySession(ctx context.Context, sessionID string) error
LinkAssets(ctx context.Context, messageID string, assets []AssetRef) error
}
Service defines message read/write behavior.
type ToolTailRoundPersister ¶
type ToolTailRoundPersister interface {
PersistToolTailRound(ctx context.Context, inputs []PersistInput) ([]Message, bool, error)
}
ToolTailRoundPersister optionally persists a complete user -> assistant(tool-call) -> tool -> assistant(final) round in one write.