300 lines
7.6 KiB
Go
300 lines
7.6 KiB
Go
package resultoutbox
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"net/http"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"cmpp-platform/gateway/internal/metrics"
|
|
|
|
"github.com/redis/go-redis/v9"
|
|
)
|
|
|
|
const (
|
|
// Eight callbacks matched the test VM's API/PostgreSQL capacity through
|
|
// 40 rps. A larger default caused callback bursts to contend with inbound
|
|
// persistence; operators can still raise it after measuring both queues.
|
|
defaultConcurrency = 8
|
|
// Keep the reclaim threshold above the callback timeout. Otherwise another
|
|
// Gateway replica could reclaim a still-running callback and execute the same
|
|
// business transition concurrently before the API stores its idempotency ID.
|
|
defaultMinIdle = 30 * time.Second
|
|
defaultBlock = 2 * time.Second
|
|
defaultHTTPTimeout = 10 * time.Second
|
|
)
|
|
|
|
var acknowledgeAndDeleteScript = redis.NewScript(`
|
|
redis.call('XACK', KEYS[1], ARGV[1], ARGV[2])
|
|
redis.call('XDEL', KEYS[1], ARGV[2])
|
|
return 1
|
|
`)
|
|
|
|
func (o *Outbox) Run(ctx context.Context) error {
|
|
if o.Redis == nil {
|
|
return fmt.Errorf("result Outbox Redis client is required")
|
|
}
|
|
if strings.TrimSpace(o.APIBaseURL) == "" {
|
|
return fmt.Errorf("result Outbox API base URL is required")
|
|
}
|
|
if o.BatchEnabled {
|
|
return retryUntilCanceled(ctx, time.Second, func() error {
|
|
if err := o.ensureGroup(ctx); err != nil {
|
|
log.Printf("gateway result Outbox group initialization failed: %v", err)
|
|
return err
|
|
}
|
|
if err := o.runBatches(ctx); err != nil && ctx.Err() == nil {
|
|
// A transient Redis read error must not permanently stop the only
|
|
// consumer for submit results, receipts, and uplink events.
|
|
log.Printf("gateway result Outbox batch consume failed: %v", err)
|
|
return err
|
|
}
|
|
return ctx.Err()
|
|
})
|
|
}
|
|
if err := o.ensureGroup(ctx); err != nil {
|
|
return err
|
|
}
|
|
pool := newCallbackPool(ctx, o, o.concurrency())
|
|
defer pool.wait()
|
|
for {
|
|
if err := o.recoverPending(ctx, pool); err != nil && ctx.Err() == nil {
|
|
log.Printf("gateway result Outbox pending recovery failed: %v", err)
|
|
sleep(ctx, time.Second)
|
|
continue
|
|
}
|
|
if ctx.Err() != nil {
|
|
return ctx.Err()
|
|
}
|
|
if err := o.consumeOnce(ctx, pool); err != nil && ctx.Err() == nil {
|
|
log.Printf("gateway result Outbox consume failed: %v", err)
|
|
sleep(ctx, time.Second)
|
|
continue
|
|
}
|
|
if ctx.Err() != nil {
|
|
return ctx.Err()
|
|
}
|
|
}
|
|
}
|
|
|
|
func (o *Outbox) ensureGroup(ctx context.Context) error {
|
|
err := o.Redis.XGroupCreateMkStream(ctx, o.stream(), o.group(), "0").Err()
|
|
if err == nil || strings.Contains(err.Error(), "BUSYGROUP") {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
func (o *Outbox) consumeOnce(ctx context.Context, pool *callbackPool) error {
|
|
available, err := pool.waitForCapacity(ctx)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
streams, err := o.Redis.XReadGroup(ctx, &redis.XReadGroupArgs{
|
|
Group: o.group(), Consumer: o.consumer(), Streams: []string{o.stream(), ">"},
|
|
Count: int64(available), Block: defaultBlock,
|
|
}).Result()
|
|
if errors.Is(err, redis.Nil) {
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, stream := range streams {
|
|
for _, message := range stream.Messages {
|
|
if !pool.dispatch(message) {
|
|
return fmt.Errorf("gateway result Outbox capacity accounting mismatch")
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (o *Outbox) recoverPending(ctx context.Context, pool *callbackPool) error {
|
|
if pool.available() == 0 {
|
|
return nil
|
|
}
|
|
messages, _, err := o.Redis.XAutoClaim(ctx, &redis.XAutoClaimArgs{
|
|
Stream: o.stream(), Group: o.group(), Consumer: o.consumer(), MinIdle: o.minIdle(),
|
|
Start: "0-0", Count: int64(pool.available()),
|
|
}).Result()
|
|
if errors.Is(err, redis.Nil) {
|
|
return nil
|
|
}
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, message := range messages {
|
|
if !pool.dispatch(message) {
|
|
return fmt.Errorf("gateway result Outbox recovery capacity accounting mismatch")
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (o *Outbox) processMessage(ctx context.Context, message redis.XMessage) error {
|
|
event, err := EventFromStreamValues(message.Values)
|
|
if err != nil {
|
|
// Malformed internal events cannot be delivered. Keep them pending for operator
|
|
// evidence instead of ACKing and silently losing a supplier result.
|
|
return err
|
|
}
|
|
startedAt := time.Now()
|
|
err = o.post(ctx, event)
|
|
metrics.ObserveSubmitStage("api_callback", err == nil, time.Since(startedAt))
|
|
if err != nil {
|
|
return err
|
|
}
|
|
// Result events have a bounded dedupe key, so successful callbacks can be
|
|
// ACKed and deleted atomically instead of turning the Outbox into an archive.
|
|
return acknowledgeAndDeleteScript.Run(
|
|
ctx,
|
|
o.Redis,
|
|
[]string{o.stream()},
|
|
o.group(),
|
|
message.ID,
|
|
).Err()
|
|
}
|
|
|
|
func (o *Outbox) post(ctx context.Context, event Event) error {
|
|
client := &http.Client{Timeout: o.httpTimeout()}
|
|
req, err := http.NewRequestWithContext(
|
|
ctx,
|
|
http.MethodPost,
|
|
strings.TrimRight(o.APIBaseURL, "/")+event.Path,
|
|
bytes.NewReader(event.Payload),
|
|
)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
req.Header.Set("Content-Type", "application/json")
|
|
req.Header.Set("X-CMPP-Result-Event-ID", event.EventID)
|
|
resp, err := client.Do(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close()
|
|
if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusMultipleChoices {
|
|
body, _ := io.ReadAll(io.LimitReader(resp.Body, 512))
|
|
return fmt.Errorf("result callback returned %d: %s", resp.StatusCode, strings.TrimSpace(string(body)))
|
|
}
|
|
return nil
|
|
}
|
|
|
|
type callbackPool struct {
|
|
ctx context.Context
|
|
outbox *Outbox
|
|
slots chan struct{}
|
|
completed chan struct{}
|
|
group sync.WaitGroup
|
|
mu sync.Mutex
|
|
active map[string]struct{}
|
|
}
|
|
|
|
func newCallbackPool(ctx context.Context, outbox *Outbox, concurrency int) *callbackPool {
|
|
return &callbackPool{
|
|
ctx: ctx, outbox: outbox, slots: make(chan struct{}, concurrency), completed: make(chan struct{}, concurrency), active: make(map[string]struct{}),
|
|
}
|
|
}
|
|
|
|
func (p *callbackPool) available() int { return cap(p.slots) - len(p.slots) }
|
|
|
|
func (p *callbackPool) waitForCapacity(ctx context.Context) (int, error) {
|
|
for p.available() == 0 {
|
|
select {
|
|
case <-ctx.Done():
|
|
return 0, ctx.Err()
|
|
case <-p.completed:
|
|
}
|
|
}
|
|
return p.available(), nil
|
|
}
|
|
|
|
func (p *callbackPool) dispatch(message redis.XMessage) bool {
|
|
p.mu.Lock()
|
|
if _, exists := p.active[message.ID]; exists {
|
|
p.mu.Unlock()
|
|
return true
|
|
}
|
|
select {
|
|
case p.slots <- struct{}{}:
|
|
p.active[message.ID] = struct{}{}
|
|
p.outbox.inFlight.Add(1)
|
|
p.group.Add(1)
|
|
p.mu.Unlock()
|
|
case <-p.ctx.Done():
|
|
p.mu.Unlock()
|
|
return false
|
|
default:
|
|
p.mu.Unlock()
|
|
return false
|
|
}
|
|
go func() {
|
|
defer func() {
|
|
p.mu.Lock()
|
|
delete(p.active, message.ID)
|
|
p.mu.Unlock()
|
|
<-p.slots
|
|
p.outbox.inFlight.Add(-1)
|
|
select {
|
|
case p.completed <- struct{}{}:
|
|
default:
|
|
}
|
|
p.group.Done()
|
|
}()
|
|
if err := p.outbox.processMessage(p.ctx, message); err != nil {
|
|
log.Printf("gateway result Outbox event %s failed: %v", message.ID, err)
|
|
}
|
|
}()
|
|
return true
|
|
}
|
|
|
|
func (p *callbackPool) wait() { p.group.Wait() }
|
|
|
|
func (o *Outbox) concurrency() int {
|
|
if o.Concurrency > 0 {
|
|
return min(o.Concurrency, 1024)
|
|
}
|
|
return defaultConcurrency
|
|
}
|
|
|
|
func (o *Outbox) minIdle() time.Duration {
|
|
if o.MinIdle > 0 {
|
|
return o.MinIdle
|
|
}
|
|
return defaultMinIdle
|
|
}
|
|
|
|
func (o *Outbox) httpTimeout() time.Duration {
|
|
if o.HTTPTimeout > 0 {
|
|
return o.HTTPTimeout
|
|
}
|
|
return defaultHTTPTimeout
|
|
}
|
|
|
|
func (o *Outbox) ConfiguredConcurrency() int { return o.concurrency() }
|
|
|
|
func sleep(ctx context.Context, duration time.Duration) {
|
|
timer := time.NewTimer(duration)
|
|
defer timer.Stop()
|
|
select {
|
|
case <-ctx.Done():
|
|
case <-timer.C:
|
|
}
|
|
}
|
|
|
|
func retryUntilCanceled(ctx context.Context, delay time.Duration, operation func() error) error {
|
|
for ctx.Err() == nil {
|
|
if err := operation(); err != nil && ctx.Err() == nil {
|
|
sleep(ctx, delay)
|
|
}
|
|
}
|
|
return ctx.Err()
|
|
}
|