Files
lislgosms/gateway/internal/resultoutbox/outbox.go
T

261 lines
6.5 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
inFlight 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) 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.MessageID == "" || event.SubmitID == "" {
return Event{}, fmt.Errorf("invalid result Outbox envelope")
}
if event.Path != "/gateway/events/submit-result" && event.Path != "/gateway/events/submit-segment-result" {
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() }