Files

288 lines
8.3 KiB
Go

package resultoutbox
import (
"context"
"encoding/json"
"fmt"
"strings"
"sync/atomic"
"time"
"cmpp-platform/gateway/internal/queue"
"github.com/redis/go-redis/v9"
)
const (
defaultStream = "gateway.submit.results"
defaultGroup = "cmpp-api-callback"
defaultConsumer = "gateway-1"
defaultDedupeTTL = 7 * 24 * time.Hour
)
var publishScript = redis.NewScript(`
local inserted = redis.call('SET', KEYS[2], '1', 'NX', 'EX', ARGV[1])
if inserted then
redis.call('XADD', KEYS[1], '*', 'data', ARGV[2])
return 1
end
return 0
`)
var publishAndAckScript = redis.NewScript(`
local inserted = redis.call('SET', KEYS[2], '1', 'NX', 'EX', ARGV[1])
if inserted then
redis.call('XADD', KEYS[1], '*', 'data', ARGV[2])
end
redis.call('XACK', KEYS[3], ARGV[3], ARGV[4])
if inserted then return 1 end
return 0
`)
type Event struct {
SchemaVersion string `json:"schemaVersion"`
EventID string `json:"eventId"`
EventType string `json:"eventType"`
Path string `json:"path"`
TraceID string `json:"traceId,omitempty"`
MessageID string `json:"messageId"`
ChannelID string `json:"channelId"`
SubmitID string `json:"submitId"`
Payload json.RawMessage `json:"payload"`
CreatedAt time.Time `json:"createdAt"`
}
type Outbox struct {
Redis *redis.Client
Stream string
Group string
Consumer string
DedupeTTL time.Duration
APIBaseURL string
HTTPTimeout time.Duration
Concurrency int
MinIdle time.Duration
BatchEnabled bool
BatchSize int
BatchWait time.Duration
GatewayInstanceID string
DeadLetterStream string
inFlight atomic.Int64
batchRequests atomic.Int64
batchEvents atomic.Int64
batchRetries atomic.Int64
deadLetters atomic.Int64
}
func New(client *redis.Client) *Outbox {
return &Outbox{Redis: client}
}
func (o *Outbox) PublishSubmitSegment(ctx context.Context, command queue.SubmitCommand, segment queue.SubmitSegmentResult) error {
payload := struct {
queue.Envelope
SubmitID string `json:"submitId,omitempty"`
queue.SubmitSegmentResult
}{
Envelope: command.Envelope,
SubmitID: command.SubmitID,
SubmitSegmentResult: segment,
}
event, err := newEvent(
fmt.Sprintf("submit:%s:segment:%d", command.SubmitID, segment.SegmentIndex),
"submit_segment_result",
"/gateway/events/submit-segment-result",
command,
payload,
)
if err != nil {
return err
}
return o.publish(ctx, event)
}
func (o *Outbox) PublishSubmitResultAndAck(
ctx context.Context,
commandStream string,
commandGroup string,
commandMessageID string,
command queue.SubmitCommand,
result queue.SubmitResult,
) error {
if o.Redis == nil {
return fmt.Errorf("result Outbox Redis client is required")
}
event, err := newEvent(
fmt.Sprintf("submit:%s:aggregate", command.SubmitID),
"submit_result",
"/gateway/events/submit-result",
command,
result,
)
if err != nil {
return err
}
data, err := json.Marshal(event)
if err != nil {
return err
}
// XADD and command XACK share one Redis script so a process crash cannot leave
// an acknowledged supplier command without its aggregate result in the Outbox.
_, err = publishAndAckScript.Run(
ctx,
o.Redis,
[]string{o.stream(), o.dedupeKey(event.EventID), commandStream},
int64(o.dedupeTTL().Seconds()),
string(data),
commandGroup,
commandMessageID,
).Result()
return err
}
func (o *Outbox) PublishSubmitResult(ctx context.Context, command queue.SubmitCommand, result queue.SubmitResult) error {
event, err := newEvent(
fmt.Sprintf("submit:%s:aggregate", command.SubmitID),
"submit_result",
"/gateway/events/submit-result",
command,
result,
)
if err != nil {
return err
}
return o.publish(ctx, event)
}
func (o *Outbox) PublishReceipt(ctx context.Context, event queue.ReceiptEvent) error {
return o.publishRaw(ctx, Event{SchemaVersion: queue.SchemaVersion, EventID: fmt.Sprintf("receipt:%s:%s:%d", event.GatewayMessageID, event.RawStatus, event.SequenceID), EventType: "receipt_intake", Path: "/gateway/events/receipt/intake", TraceID: event.TraceID, MessageID: event.MessageID, ChannelID: event.ChannelID, Payload: mustMarshal(event), CreatedAt: time.Now().UTC()})
}
func (o *Outbox) PublishUplink(ctx context.Context, event queue.UplinkEvent) error {
return o.publishRaw(ctx, Event{SchemaVersion: queue.SchemaVersion, EventID: fmt.Sprintf("uplink:%s:%d:%d", event.ChannelID, event.SequenceID, event.ReceivedAt.UnixNano()), EventType: "uplink", Path: "/gateway/events/uplink", TraceID: event.TraceID, MessageID: event.MessageID, ChannelID: event.ChannelID, Payload: mustMarshal(event), CreatedAt: time.Now().UTC()})
}
func mustMarshal(value any) json.RawMessage { data, _ := json.Marshal(value); return data }
func (o *Outbox) publishRaw(ctx context.Context, event Event) error { return o.publish(ctx, event) }
func (o *Outbox) publish(ctx context.Context, event Event) error {
if o.Redis == nil {
return fmt.Errorf("result Outbox Redis client is required")
}
data, err := json.Marshal(event)
if err != nil {
return err
}
_, err = publishScript.Run(
ctx,
o.Redis,
[]string{o.stream(), o.dedupeKey(event.EventID)},
int64(o.dedupeTTL().Seconds()),
string(data),
).Result()
return err
}
func newEvent(eventID string, eventType string, path string, command queue.SubmitCommand, payload any) (Event, error) {
data, err := json.Marshal(payload)
if err != nil {
return Event{}, err
}
var object map[string]interface{}
if err := json.Unmarshal(data, &object); err != nil {
return Event{}, err
}
// The API stores the deterministic ID on the submit attempt after successful
// processing, making callback redelivery idempotent across Gateway restarts.
object["eventId"] = eventID
data, err = json.Marshal(object)
if err != nil {
return Event{}, err
}
return Event{
SchemaVersion: queue.SchemaVersion,
EventID: eventID,
EventType: eventType,
Path: path,
TraceID: command.TraceID,
MessageID: command.MessageID,
ChannelID: command.ChannelID,
SubmitID: command.SubmitID,
Payload: data,
CreatedAt: time.Now().UTC(),
}, nil
}
func EventFromStreamValues(values map[string]interface{}) (Event, error) {
raw, ok := values["data"]
if !ok {
return Event{}, fmt.Errorf("result Outbox data field is required")
}
var data string
switch value := raw.(type) {
case string:
data = value
case []byte:
data = string(value)
default:
data = fmt.Sprint(value)
}
var event Event
if err := json.Unmarshal([]byte(data), &event); err != nil {
return Event{}, err
}
if event.SchemaVersion != queue.SchemaVersion || event.EventID == "" || event.ChannelID == "" {
return Event{}, fmt.Errorf("invalid result Outbox envelope")
}
if (event.EventType == "submit_result" || event.EventType == "submit_segment_result") && event.MessageID == "" {
return Event{}, fmt.Errorf("submit result Outbox messageId is required")
}
if event.Path != "/gateway/events/submit-result" && event.Path != "/gateway/events/submit-segment-result" && event.Path != "/gateway/events/receipt/intake" && event.Path != "/gateway/events/uplink" && event.Path != "/gateway/events/dead-letter" {
return Event{}, fmt.Errorf("unsupported result Outbox path %q", event.Path)
}
if len(event.Payload) == 0 {
return Event{}, fmt.Errorf("result Outbox payload is required")
}
return event, nil
}
func (o *Outbox) stream() string {
if strings.TrimSpace(o.Stream) != "" {
return o.Stream
}
return defaultStream
}
func (o *Outbox) group() string {
if strings.TrimSpace(o.Group) != "" {
return o.Group
}
return defaultGroup
}
func (o *Outbox) consumer() string {
if strings.TrimSpace(o.Consumer) != "" {
return o.Consumer
}
return defaultConsumer
}
func (o *Outbox) dedupeTTL() time.Duration {
if o.DedupeTTL > 0 {
return o.DedupeTTL
}
return defaultDedupeTTL
}
func (o *Outbox) dedupeKey(eventID string) string {
return o.stream() + ":dedupe:" + eventID
}
func (o *Outbox) StreamName() string { return o.stream() }
func (o *Outbox) GroupName() string { return o.group() }
func (o *Outbox) InFlight() int64 { return o.inFlight.Load() }
func (o *Outbox) BatchCounts() (int64, int64, int64, int64) {
return o.batchRequests.Load(), o.batchEvents.Load(), o.batchRetries.Load(), o.deadLetters.Load()
}