285 lines
8.1 KiB
Go
285 lines
8.1 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.MessageID == "" {
|
|
return Event{}, fmt.Errorf("invalid result Outbox envelope")
|
|
}
|
|
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()
|
|
}
|