167 lines
5.3 KiB
Go
167 lines
5.3 KiB
Go
package bell_connector
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/clause"
|
|
|
|
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
|
|
)
|
|
|
|
func EnqueueEvent(tx *gorm.DB, payload []byte, now time.Time) (outbox.Message, error) {
|
|
identity, err := parseEventIdentity(payload)
|
|
if err != nil {
|
|
return outbox.Message{}, err
|
|
}
|
|
keyDigest := sha256.Sum256([]byte(identity.ProducerID + "\x00" + identity.SourceEventID))
|
|
return outbox.Enqueue(tx, outbox.EnqueueInput{
|
|
InternalType: OutboxType,
|
|
BusinessRef: identity.SourceEventID,
|
|
IdempotencyKey: "bell-event-v1:" + hex.EncodeToString(keyDigest[:]),
|
|
PayloadJSON: append([]byte(nil), payload...),
|
|
}, now)
|
|
}
|
|
|
|
func parseEventIdentity(payload []byte) (eventIdentity, error) {
|
|
decoder := json.NewDecoder(bytes.NewReader(payload))
|
|
var identity eventIdentity
|
|
if err := decoder.Decode(&identity); err != nil || !json.Valid(payload) || identity.SchemaVersion != "yovision.event/v1" ||
|
|
!safeIdentifier(identity.ProducerID) || !safeIdentifier(identity.SourceEventID) {
|
|
return eventIdentity{}, errors.New("invalid yovision.event/v1 payload")
|
|
}
|
|
return identity, nil
|
|
}
|
|
|
|
func safeIdentifier(value string) bool {
|
|
if value == "" || len(value) > 128 || strings.ContainsAny(value, "\\/@\x00\r\n") {
|
|
return false
|
|
}
|
|
for index, r := range value {
|
|
allowed := r >= 'A' && r <= 'Z' || r >= 'a' && r <= 'z' || r >= '0' && r <= '9' || (index > 0 && strings.ContainsRune("._:-", r))
|
|
if !allowed {
|
|
return false
|
|
}
|
|
}
|
|
return true
|
|
}
|
|
|
|
type Relay struct {
|
|
DB *gorm.DB
|
|
Client *Client
|
|
Now func() time.Time
|
|
Backoff func(int) time.Duration
|
|
}
|
|
|
|
func (r Relay) DeliverBatch(ctx context.Context, worker string, limit int) (int, error) {
|
|
if r.DB == nil || r.Client == nil || strings.TrimSpace(worker) == "" || limit < 1 || limit > 100 {
|
|
return 0, errors.New("invalid Bell relay configuration")
|
|
}
|
|
items, err := r.claim(worker, limit)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
delivered := 0
|
|
queueRelay := outbox.NewRelay(r.DB)
|
|
queueRelay.Now = r.now
|
|
if r.Backoff != nil {
|
|
queueRelay.Backoff = r.Backoff
|
|
}
|
|
for _, item := range items {
|
|
_, deliveryErr := r.Client.Send(ctx, []byte(item.PayloadJSON))
|
|
if deliveryErr == nil {
|
|
if err = queueRelay.MarkSuccess(item.ID, worker); err != nil {
|
|
return delivered, err
|
|
}
|
|
delivered++
|
|
continue
|
|
}
|
|
var classified *DeliveryError
|
|
if errors.As(deliveryErr, &classified) && classified.Terminal {
|
|
if err = r.markTerminal(item, worker, classified.Code); err != nil {
|
|
return delivered, err
|
|
}
|
|
continue
|
|
}
|
|
if err = queueRelay.MarkFailure(item.ID, worker, deliveryErr.Error()); err != nil {
|
|
return delivered, err
|
|
}
|
|
}
|
|
return delivered, nil
|
|
}
|
|
|
|
func (r Relay) claim(worker string, limit int) ([]outbox.Message, error) {
|
|
now := r.now()
|
|
leaseUntil := now.Add(30 * time.Second)
|
|
claimed := make([]outbox.Message, 0, limit)
|
|
err := r.DB.Transaction(func(tx *gorm.DB) error {
|
|
var candidates []outbox.Message
|
|
query := tx.Where("internal_type = ? AND (((state IN ?) AND available_at <= ?) OR (state = ? AND lease_until < ?))", OutboxType, []string{outbox.StatePending, outbox.StateRetry}, now, outbox.StateProcessing, now).Order("available_at, created_at").Limit(limit)
|
|
if tx.Dialector.Name() == "postgres" {
|
|
query = query.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"})
|
|
}
|
|
if err := query.Find(&candidates).Error; err != nil {
|
|
return err
|
|
}
|
|
for _, item := range candidates {
|
|
result := tx.Model(&outbox.Message{}).Where("id = ? AND version = ?", item.ID, item.Version).Updates(map[string]any{"state": outbox.StateProcessing, "lease_owner": worker, "lease_until": leaseUntil, "version": gorm.Expr("version + 1"), "updated_at": now})
|
|
if result.Error != nil {
|
|
return result.Error
|
|
}
|
|
if result.RowsAffected == 1 {
|
|
item.State, item.LeaseOwner, item.LeaseUntil, item.Version = outbox.StateProcessing, worker, &leaseUntil, item.Version+1
|
|
claimed = append(claimed, item)
|
|
}
|
|
}
|
|
return nil
|
|
})
|
|
return claimed, err
|
|
}
|
|
|
|
func (r Relay) markTerminal(item outbox.Message, worker, detail string) error {
|
|
now := r.now()
|
|
return r.DB.Transaction(func(tx *gorm.DB) error {
|
|
result := tx.Model(&outbox.Message{}).Where("id = ? AND state = ? AND lease_owner = ?", item.ID, outbox.StateProcessing, worker).Updates(map[string]any{
|
|
"state": outbox.StateDead, "attempt_count": gorm.Expr("attempt_count + 1"), "last_error": detail,
|
|
"lease_owner": "", "lease_until": nil, "version": gorm.Expr("version + 1"), "updated_at": now,
|
|
})
|
|
if result.Error != nil || result.RowsAffected != 1 {
|
|
if result.Error != nil {
|
|
return result.Error
|
|
}
|
|
return errors.New("Bell outbox lease lost")
|
|
}
|
|
return tx.Create(&outbox.Attempt{MessageID: item.ID, Number: item.AttemptCount + 1, Outcome: outbox.StateDead, Detail: detail, Worker: worker, CreatedAt: now}).Error
|
|
})
|
|
}
|
|
|
|
func (r Relay) now() time.Time {
|
|
if r.Now != nil {
|
|
return r.Now().UTC()
|
|
}
|
|
return time.Now().UTC()
|
|
}
|
|
|
|
func PreserveIdentity(before, after []byte) error {
|
|
left, err := parseEventIdentity(before)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
right, err := parseEventIdentity(after)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if left.ProducerID != right.ProducerID || left.SourceEventID != right.SourceEventID {
|
|
return fmt.Errorf("relay changed original event identity")
|
|
}
|
|
return nil
|
|
}
|