5759c1862e
M0-M7 已完成:核心网关(身份/RBAC/TOTP/OIDC/SAML/Provider/配额/路由/内容策略/审计/定价)+ 资源市场(MCP/Skills/数字员工)。 含 22 个 PostgreSQL 迁移、管理端/门户端前端源码、OpenAPI 契约、部署 compose。 Co-Authored-By: Claude <noreply@anthropic.com>
202 lines
7.5 KiB
Go
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
|
|
}
|