Files
ai-gateway-go/internal/outbox/store.go
T
superidou 5759c1862e AI Gateway Go 0.10.0 源码快照 + 旗舰版需求规划报告
M0-M7 已完成:核心网关(身份/RBAC/TOTP/OIDC/SAML/Provider/配额/路由/内容策略/审计/定价)+ 资源市场(MCP/Skills/数字员工)。
含 22 个 PostgreSQL 迁移、管理端/门户端前端源码、OpenAPI 契约、部署 compose。

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-12 11:45:54 +08:00

202 lines
7.5 KiB
Go

package outbox
import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
)
var (
ErrStoreUnavailable = errors.New("outbox store unavailable")
ErrEventNotFound = errors.New("outbox event not found")
)
type Event struct {
EventID string `json:"event_id"`
EventType string `json:"event_type"`
EventVersion int `json:"event_version"`
TenantID *string `json:"tenant_id"`
AggregateType string `json:"aggregate_type"`
AggregateID string `json:"aggregate_id"`
Payload json.RawMessage `json:"payload"`
TraceContext json.RawMessage `json:"trace_context"`
OccurredAt time.Time `json:"occurred_at"`
Attempts int `json:"attempts"`
}
type Store struct{ pool *pgxpool.Pool }
func NewStore(pool *pgxpool.Pool) *Store { return &Store{pool: pool} }
func (s *Store) Claim(ctx context.Context, workerID string, limit int, lease time.Duration) ([]Event, error) {
if s == nil || s.pool == nil {
return nil, ErrStoreUnavailable
}
rows, err := s.pool.Query(ctx, `
WITH candidates AS (
SELECT event_id FROM gateway.outbox_events
WHERE processed_at IS NULL AND dead_lettered_at IS NULL AND available_at <= clock_timestamp()
AND (locked_at IS NULL OR locked_at < clock_timestamp()-($1 * interval '1 second'))
ORDER BY available_at,occurred_at
FOR UPDATE SKIP LOCKED LIMIT $2
)
UPDATE gateway.outbox_events e
SET locked_at=clock_timestamp(),locked_by=$3,attempts=e.attempts+1
FROM candidates c WHERE e.event_id=c.event_id
RETURNING e.event_id::text,e.event_type,e.event_version,e.tenant_id::text,e.aggregate_type,
e.aggregate_id,e.payload,e.trace_context,e.occurred_at,e.attempts`, lease.Seconds(), limit, workerID)
if err != nil {
return nil, fmt.Errorf("%w: claim: %v", ErrStoreUnavailable, err)
}
defer rows.Close()
events := make([]Event, 0, limit)
for rows.Next() {
var event Event
if err := rows.Scan(&event.EventID, &event.EventType, &event.EventVersion, &event.TenantID,
&event.AggregateType, &event.AggregateID, &event.Payload, &event.TraceContext, &event.OccurredAt, &event.Attempts); err != nil {
return nil, fmt.Errorf("%w: scan claim: %v", ErrStoreUnavailable, err)
}
events = append(events, event)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("%w: claim rows: %v", ErrStoreUnavailable, err)
}
return events, nil
}
func (s *Store) MarkProcessed(ctx context.Context, eventID, workerID, streamID string) error {
if s == nil || s.pool == nil {
return ErrStoreUnavailable
}
result, err := s.pool.Exec(ctx, `UPDATE gateway.outbox_events SET processed_at=clock_timestamp(),locked_at=NULL,locked_by=NULL,last_error=NULL,published_stream_id=$3 WHERE event_id=$1 AND locked_by=$2 AND processed_at IS NULL`, eventID, workerID, streamID)
if err != nil {
return fmt.Errorf("%w: mark processed: %v", ErrStoreUnavailable, err)
}
if result.RowsAffected() != 1 {
return ErrEventNotFound
}
return nil
}
func (s *Store) MarkFailed(ctx context.Context, event Event, workerID string, deliveryErr error, maxAttempts int, delay time.Duration) error {
if s == nil || s.pool == nil {
return ErrStoreUnavailable
}
message := deliveryErr.Error()
if len(message) > 2048 {
message = message[:2048]
}
dead := event.Attempts >= maxAttempts
result, err := s.pool.Exec(ctx, `
UPDATE gateway.outbox_events SET locked_at=NULL,locked_by=NULL,last_error=$3,
available_at=CASE WHEN $4 THEN available_at ELSE clock_timestamp()+($5 * interval '1 second') END,
dead_lettered_at=CASE WHEN $4 THEN clock_timestamp() ELSE NULL END
WHERE event_id=$1 AND locked_by=$2 AND processed_at IS NULL`, event.EventID, workerID, message, dead, delay.Seconds())
if err != nil {
return fmt.Errorf("%w: mark failed: %v", ErrStoreUnavailable, err)
}
if result.RowsAffected() != 1 {
return ErrEventNotFound
}
return nil
}
type EventView struct {
Event
AvailableAt time.Time `json:"available_at"`
LockedAt *time.Time `json:"locked_at"`
LockedBy *string `json:"locked_by"`
ProcessedAt *time.Time `json:"processed_at"`
DeadLetteredAt *time.Time `json:"dead_lettered_at"`
LastError *string `json:"last_error"`
PublishedStream *string `json:"published_stream_id"`
}
func (s *Store) List(ctx context.Context, status, eventType string, limit int) ([]EventView, error) {
if s == nil || s.pool == nil {
return nil, ErrStoreUnavailable
}
where := "TRUE"
switch status {
case "pending":
where = "processed_at IS NULL AND dead_lettered_at IS NULL"
case "dead":
where = "dead_lettered_at IS NOT NULL"
case "processed":
where = "processed_at IS NOT NULL"
}
args := []any{limit}
if eventType != "" {
args = append(args, eventType)
where += fmt.Sprintf(" AND event_type=$%d", len(args))
}
rows, err := s.pool.Query(ctx, `SELECT event_id::text,event_type,event_version,tenant_id::text,aggregate_type,aggregate_id,payload,trace_context,occurred_at,attempts,available_at,locked_at,locked_by,processed_at,dead_lettered_at,last_error,published_stream_id FROM gateway.outbox_events WHERE `+where+` ORDER BY occurred_at DESC LIMIT $1`, args...)
if err != nil {
return nil, fmt.Errorf("%w: list: %v", ErrStoreUnavailable, err)
}
defer rows.Close()
items := make([]EventView, 0)
for rows.Next() {
var item EventView
if err := rows.Scan(&item.EventID, &item.EventType, &item.EventVersion, &item.TenantID, &item.AggregateType,
&item.AggregateID, &item.Payload, &item.TraceContext, &item.OccurredAt, &item.Attempts, &item.AvailableAt,
&item.LockedAt, &item.LockedBy, &item.ProcessedAt, &item.DeadLetteredAt, &item.LastError, &item.PublishedStream); err != nil {
return nil, fmt.Errorf("%w: scan list: %v", ErrStoreUnavailable, err)
}
items = append(items, item)
}
return items, rows.Err()
}
func (s *Store) Retry(ctx context.Context, eventID string) error {
if s == nil || s.pool == nil {
return ErrStoreUnavailable
}
result, err := s.pool.Exec(ctx, `UPDATE gateway.outbox_events SET attempts=0,available_at=clock_timestamp(),locked_at=NULL,locked_by=NULL,dead_lettered_at=NULL,last_error=NULL WHERE event_id=$1 AND processed_at IS NULL`, eventID)
if err != nil {
return fmt.Errorf("%w: retry: %v", ErrStoreUnavailable, err)
}
if result.RowsAffected() != 1 {
return ErrEventNotFound
}
return nil
}
type ConsumerHandler func(context.Context, pgx.Tx) error
func (s *Store) Consume(ctx context.Context, subscriber, eventID string, handler ConsumerHandler) (bool, error) {
if s == nil || s.pool == nil {
return false, ErrStoreUnavailable
}
subscriber = strings.TrimSpace(subscriber)
if subscriber == "" || eventID == "" || handler == nil {
return false, errors.New("subscriber, event ID, and handler are required")
}
tx, err := s.pool.Begin(ctx)
if err != nil {
return false, fmt.Errorf("%w: begin consumption: %v", ErrStoreUnavailable, err)
}
defer func() { _ = tx.Rollback(ctx) }()
var inserted bool
err = tx.QueryRow(ctx, `WITH inserted AS (INSERT INTO gateway.event_consumptions(subscriber,event_id) VALUES($1,$2) ON CONFLICT DO NOTHING RETURNING 1) SELECT EXISTS(SELECT 1 FROM inserted)`, subscriber, eventID).Scan(&inserted)
if err != nil {
return false, fmt.Errorf("%w: reserve consumption: %v", ErrStoreUnavailable, err)
}
if !inserted {
return false, nil
}
if err := handler(ctx, tx); err != nil {
return false, err
}
if err := tx.Commit(ctx); err != nil {
return false, fmt.Errorf("%w: commit consumption: %v", ErrStoreUnavailable, err)
}
return true, nil
}