[SEN] 建立内部事务 Outbox 与可靠 relay 状态机 #122
@@ -0,0 +1,19 @@
|
||||
package router
|
||||
|
||||
import (
|
||||
"github.com/gin-gonic/gin"
|
||||
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/common/actions"
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/common/middleware"
|
||||
)
|
||||
|
||||
func init() { routerCheckRole = append(routerCheckRole, registerSenseOutboxRouter) }
|
||||
func registerSenseOutboxRouter(v1 *gin.RouterGroup, auth *jwt.GinJWTMiddleware) {
|
||||
api := &outbox.API{}
|
||||
r := v1.Group("/outbox").Use(auth.MiddlewareFunc()).Use(middleware.AuthCheckRole()).Use(actions.PermissionAction())
|
||||
r.GET("", api.List)
|
||||
r.GET("/:id", api.Get)
|
||||
r.POST("/:id/requeue", api.Requeue)
|
||||
}
|
||||
@@ -0,0 +1,26 @@
|
||||
package router
|
||||
|
||||
import (
|
||||
"github.com/gin-gonic/gin"
|
||||
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
|
||||
"net/http"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestSenseOutboxRoutes(t *testing.T) {
|
||||
gin.SetMode(gin.TestMode)
|
||||
engine := gin.New()
|
||||
registerSenseOutboxRouter(engine.Group("/api/v1"), &jwt.GinJWTMiddleware{})
|
||||
wanted := map[string]bool{http.MethodGet + " /api/v1/outbox": false, http.MethodGet + " /api/v1/outbox/:id": false, http.MethodPost + " /api/v1/outbox/:id/requeue": false}
|
||||
for _, route := range engine.Routes() {
|
||||
key := route.Method + " " + route.Path
|
||||
if _, ok := wanted[key]; ok {
|
||||
wanted[key] = true
|
||||
}
|
||||
}
|
||||
for route, found := range wanted {
|
||||
if !found {
|
||||
t.Fatalf("route not registered: %s", route)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
package local_event
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
|
||||
)
|
||||
|
||||
// CreateWithOutbox commits the local candidate and its internal delivery
|
||||
// record atomically. The payload remains Sense-internal and is not a Bell or
|
||||
// Brain contract.
|
||||
func CreateWithOutbox(ctx context.Context, db *gorm.DB, candidate EventCandidate, payload map[string]interface{}, now time.Time) error {
|
||||
encoded, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode local event outbox payload: %w", err)
|
||||
}
|
||||
return db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Create(&candidate).Error; err != nil {
|
||||
return fmt.Errorf("create local event candidate: %w", err)
|
||||
}
|
||||
_, err = outbox.Enqueue(tx, outbox.EnqueueInput{InternalType: "local_event_candidate", BusinessRef: candidate.ID, IdempotencyKey: "local-event:" + candidate.ID + ":v1", PayloadJSON: encoded}, now)
|
||||
return err
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
package local_event
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
|
||||
)
|
||||
|
||||
func TestCreateWithOutboxCommitsAndRollsBackAtomically(t *testing.T) {
|
||||
db, err := gorm.Open(sqlite.Open("file:"+uuid.NewString()+"?mode=memory&cache=shared"), &gorm.Config{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB, _ := db.DB()
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
if err = db.AutoMigrate(&EventCandidate{}, &outbox.Message{}, &outbox.DeliveryRecord{}, &outbox.Attempt{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
now := time.Date(2026, 8, 28, 10, 0, 0, 0, time.UTC)
|
||||
candidate := EventCandidate{ID: uuid.NewString(), OccurredAt: now, SourceRef: "SEN-CAM-01", RuleRef: "rule-1", RuleName: "区域闯入", CandidateState: CandidateStateCandidate, EvidenceState: EvidenceStatePending, RetainUntil: now.Add(24 * time.Hour)}
|
||||
if err = CreateWithOutbox(context.Background(), db, candidate, map[string]interface{}{"eventId": candidate.ID}, now); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var candidates, messages int64
|
||||
db.Model(&EventCandidate{}).Count(&candidates)
|
||||
db.Model(&outbox.Message{}).Count(&messages)
|
||||
if candidates != 1 || messages != 1 {
|
||||
t.Fatalf("candidates=%d messages=%d", candidates, messages)
|
||||
}
|
||||
duplicate := candidate
|
||||
duplicate.ID = candidate.ID
|
||||
if err = CreateWithOutbox(context.Background(), db, duplicate, map[string]interface{}{"eventId": duplicate.ID}, now); err == nil {
|
||||
t.Fatal("expected duplicate transaction failure")
|
||||
}
|
||||
db.Model(&EventCandidate{}).Count(&candidates)
|
||||
db.Model(&outbox.Message{}).Count(&messages)
|
||||
if candidates != 1 || messages != 1 {
|
||||
t.Fatalf("atomic rollback failed candidates=%d messages=%d", candidates, messages)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,103 @@
|
||||
package outbox
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/gin-gonic/gin/binding"
|
||||
"github.com/go-admin-team/go-admin-core/sdk/api"
|
||||
"github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth/user"
|
||||
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/common"
|
||||
)
|
||||
|
||||
type API struct{ api.Api }
|
||||
|
||||
func (e *API) service(c *gin.Context) (*Service, error) {
|
||||
base := coreService.Service{}
|
||||
if err := e.MakeContext(c).MakeOrm().MakeService(&base).Errors; err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return NewService(base.Orm), nil
|
||||
}
|
||||
|
||||
func (e *API) List(c *gin.Context) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
request := PageRequest{}
|
||||
if err = e.MakeContext(c).Bind(&request).Errors; err != nil {
|
||||
e.audit(c, service, "List", auditFailure, "可靠投递查询条件格式不正确")
|
||||
e.Error(http.StatusBadRequest, err, "查询条件格式不正确")
|
||||
return
|
||||
}
|
||||
response, err := service.List(request)
|
||||
if err != nil {
|
||||
e.audit(c, service, "List", auditFailure, "可靠投递列表查询失败")
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
e.audit(c, service, "List", auditSuccess, "读取可靠投递列表")
|
||||
e.OK(response, "查询成功")
|
||||
}
|
||||
|
||||
func (e *API) Get(c *gin.Context) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
response, err := service.Get(c.Param("id"))
|
||||
if err != nil {
|
||||
e.audit(c, service, "Get", auditFailure, "可靠投递详情查询失败")
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
e.audit(c, service, "Get", auditSuccess, "读取可靠投递详情 "+response.Message.ID)
|
||||
e.OK(response, "查询成功")
|
||||
}
|
||||
|
||||
func (e *API) Requeue(c *gin.Context) {
|
||||
service, err := e.service(c)
|
||||
if err != nil {
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
request := RequeueRequest{}
|
||||
if err = e.MakeContext(c).Bind(&request, binding.JSON).Errors; err != nil {
|
||||
e.audit(c, service, "Requeue", auditFailure, "死信重新排队请求格式不正确")
|
||||
e.Error(http.StatusBadRequest, err, "请求格式不正确")
|
||||
return
|
||||
}
|
||||
message, err := service.Requeue(c.Param("id"), request, user.GetUserId(c))
|
||||
if err != nil {
|
||||
e.audit(c, service, "Requeue", auditFailure, "死信重新排队被拒绝")
|
||||
e.writeError(err)
|
||||
return
|
||||
}
|
||||
e.audit(c, service, "Requeue", auditSuccess, "死信已重新排队 "+message.ID)
|
||||
e.OK(message, "已重新排队")
|
||||
}
|
||||
|
||||
func (e *API) audit(c *gin.Context, service *Service, action, status, remark string) {
|
||||
if err := WriteAudit(service.Orm, Audit{Action: action, Method: c.Request.Method, Status: status, Username: user.GetUserName(c), UserID: user.GetUserId(c), ClientIP: common.GetClientIP(c), Route: c.FullPath(), Remark: remark, At: time.Now()}); err != nil {
|
||||
api.GetRequestLogger(c).Errorf("outbox audit failed: %s", err.Error())
|
||||
}
|
||||
}
|
||||
func (e *API) writeError(err error) {
|
||||
switch {
|
||||
case errors.Is(err, ErrInvalidInput):
|
||||
e.Error(http.StatusBadRequest, err, err.Error())
|
||||
case errors.Is(err, ErrNotFound):
|
||||
e.Error(http.StatusNotFound, err, err.Error())
|
||||
case errors.Is(err, ErrNotDead), errors.Is(err, ErrVersionConflict):
|
||||
e.Error(http.StatusConflict, err, err.Error())
|
||||
default:
|
||||
e.Error(http.StatusInternalServerError, err, "可靠投递操作失败")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,23 @@
|
||||
package outbox
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
|
||||
adminModels "git.ilapage.cn/ila/yovision/Sense/server/app/admin/models"
|
||||
)
|
||||
|
||||
const auditSuccess, auditFailure = "1", "2"
|
||||
|
||||
type Audit struct {
|
||||
Action, Method, Status, Username, ClientIP, Route, Remark string
|
||||
UserID int
|
||||
At time.Time
|
||||
}
|
||||
|
||||
func WriteAudit(db *gorm.DB, input Audit) error {
|
||||
model := adminModels.SysOperaLog{Title: "可靠投递", BusinessType: "other", Method: "outbox.API." + input.Action, RequestMethod: input.Method, OperatorType: "1", OperName: input.Username, OperUrl: input.Route, OperIp: input.ClientIP, Status: input.Status, OperTime: input.At.UTC(), Remark: input.Remark, CreatedAt: input.At.UTC(), UpdatedAt: input.At.UTC()}
|
||||
model.CreateBy, model.UpdateBy = input.UserID, input.UserID
|
||||
return db.Create(&model).Error
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package outbox
|
||||
|
||||
import commonDTO "git.ilapage.cn/ila/yovision/Sense/server/common/dto"
|
||||
|
||||
type PageRequest struct {
|
||||
commonDTO.Pagination `search:"-"`
|
||||
State string `form:"state"`
|
||||
InternalType string `form:"internalType"`
|
||||
Keyword string `form:"keyword"`
|
||||
}
|
||||
|
||||
type Summary struct {
|
||||
Pending int64 `json:"pending"`
|
||||
Retry int64 `json:"retry"`
|
||||
Processing int64 `json:"processing"`
|
||||
Dead int64 `json:"dead"`
|
||||
}
|
||||
|
||||
type PageResponse struct {
|
||||
List []Message `json:"list"`
|
||||
Count int64 `json:"count"`
|
||||
Summary Summary `json:"summary"`
|
||||
}
|
||||
|
||||
type DetailResponse struct {
|
||||
Message Message `json:"message"`
|
||||
Attempts []Attempt `json:"attempts"`
|
||||
}
|
||||
|
||||
type RequeueRequest struct {
|
||||
ExpectedVersion int64 `json:"expectedVersion" binding:"required,min=1"`
|
||||
Reason string `json:"reason" binding:"required"`
|
||||
}
|
||||
|
||||
type EnqueueInput struct {
|
||||
InternalType string
|
||||
BusinessRef string
|
||||
IdempotencyKey string
|
||||
PayloadJSON []byte
|
||||
MaxAttempts int
|
||||
}
|
||||
@@ -0,0 +1,58 @@
|
||||
package outbox
|
||||
|
||||
import "time"
|
||||
|
||||
const (
|
||||
StatePending = "pending"
|
||||
StateProcessing = "processing"
|
||||
StateRetry = "retry"
|
||||
StateDead = "dead"
|
||||
StateDelivered = "delivered"
|
||||
)
|
||||
|
||||
// Message is a Sense-internal delivery record. PayloadJSON is deliberately
|
||||
// excluded from management APIs and is not a cross-product contract.
|
||||
type Message struct {
|
||||
ID string `gorm:"size:36;primaryKey" json:"id"`
|
||||
InternalType string `gorm:"size:64;not null;index" json:"internalType"`
|
||||
BusinessRef string `gorm:"size:128;not null;index" json:"businessRef"`
|
||||
IdempotencyKey string `gorm:"size:191;not null;uniqueIndex" json:"idempotencyKey"`
|
||||
PayloadJSON string `gorm:"column:payload;type:jsonb;not null" json:"-"`
|
||||
State string `gorm:"size:24;not null;index" json:"state"`
|
||||
AttemptCount int `gorm:"not null;default:0" json:"attemptCount"`
|
||||
MaxAttempts int `gorm:"not null;default:12" json:"maxAttempts"`
|
||||
AvailableAt time.Time `gorm:"not null;index" json:"availableAt"`
|
||||
LeaseOwner string `gorm:"size:128;not null;default:''" json:"leaseOwner,omitempty"`
|
||||
LeaseUntil *time.Time `gorm:"index" json:"leaseUntil,omitempty"`
|
||||
LastError string `gorm:"size:512;not null;default:''" json:"lastError,omitempty"`
|
||||
DeliveredAt *time.Time `json:"deliveredAt,omitempty"`
|
||||
Version int64 `gorm:"not null;default:1" json:"version"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
UpdatedAt time.Time `json:"updatedAt"`
|
||||
}
|
||||
|
||||
func (Message) TableName() string { return "sense_outbox_messages" }
|
||||
|
||||
// DeliveryRecord is permanent idempotency evidence. It is never deleted by
|
||||
// queue cleanup and prevents a delivered key from being processed again.
|
||||
type DeliveryRecord struct {
|
||||
IdempotencyKey string `gorm:"size:191;primaryKey" json:"idempotencyKey"`
|
||||
MessageID string `gorm:"size:36;not null;uniqueIndex" json:"messageId"`
|
||||
DeliveredAt time.Time `gorm:"not null" json:"deliveredAt"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
}
|
||||
|
||||
func (DeliveryRecord) TableName() string { return "sense_outbox_deliveries" }
|
||||
|
||||
type Attempt struct {
|
||||
ID uint `gorm:"primaryKey;autoIncrement" json:"id"`
|
||||
MessageID string `gorm:"size:36;not null;index" json:"messageId"`
|
||||
Number int `gorm:"not null" json:"number"`
|
||||
Outcome string `gorm:"size:32;not null" json:"outcome"`
|
||||
Detail string `gorm:"size:512;not null;default:''" json:"detail"`
|
||||
Worker string `gorm:"size:128;not null;default:''" json:"worker,omitempty"`
|
||||
ActorUserID int `gorm:"not null;default:0" json:"actorUserId,omitempty"`
|
||||
CreatedAt time.Time `json:"createdAt"`
|
||||
}
|
||||
|
||||
func (Attempt) TableName() string { return "sense_outbox_attempts" }
|
||||
@@ -0,0 +1,92 @@
|
||||
package outbox
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"gorm.io/driver/postgres"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func TestPostgresConcurrentWorkersDoNotClaimSameMessage(t *testing.T) {
|
||||
dsn := os.Getenv("SENSE_OUTBOX_TEST_DATABASE_URL")
|
||||
if dsn == "" {
|
||||
t.Skip("set SENSE_OUTBOX_TEST_DATABASE_URL to run the PostgreSQL multi-worker test")
|
||||
}
|
||||
base, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
schema := fmt.Sprintf("sense_outbox_78_%d", time.Now().UnixNano())
|
||||
if err = base.Exec("CREATE SCHEMA " + schema).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Cleanup(func() { _ = base.Exec("DROP SCHEMA IF EXISTS " + schema + " CASCADE").Error })
|
||||
scoped, err := gorm.Open(postgres.Open(withSearchPath(dsn, schema)), &gorm.Config{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err = scoped.AutoMigrate(&Message{}, &DeliveryRecord{}, &Attempt{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
now := time.Date(2026, 8, 28, 12, 0, 0, 0, time.UTC)
|
||||
for index := 0; index < 20; index++ {
|
||||
if _, err = Enqueue(scoped, EnqueueInput{InternalType: "local_event_candidate", BusinessRef: fmt.Sprintf("event-%d", index), IdempotencyKey: fmt.Sprintf("event:%d:v1", index), PayloadJSON: []byte(`{}`)}, now); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
workers := []string{"worker-a", "worker-b"}
|
||||
results := make(chan []Message, len(workers))
|
||||
errorsCh := make(chan error, len(workers))
|
||||
var group sync.WaitGroup
|
||||
for _, worker := range workers {
|
||||
group.Add(1)
|
||||
go func(name string) {
|
||||
defer group.Done()
|
||||
relay := NewRelay(scoped)
|
||||
relay.Now = func() time.Time { return now }
|
||||
items, claimErr := relay.Claim(name, 20)
|
||||
if claimErr != nil {
|
||||
errorsCh <- claimErr
|
||||
return
|
||||
}
|
||||
results <- items
|
||||
}(worker)
|
||||
}
|
||||
group.Wait()
|
||||
close(results)
|
||||
close(errorsCh)
|
||||
for claimErr := range errorsCh {
|
||||
t.Fatal(claimErr)
|
||||
}
|
||||
seen := map[string]string{}
|
||||
for batch := range results {
|
||||
for _, item := range batch {
|
||||
if owner, exists := seen[item.ID]; exists {
|
||||
t.Fatalf("message %s claimed by %s and %s", item.ID, owner, item.LeaseOwner)
|
||||
}
|
||||
seen[item.ID] = item.LeaseOwner
|
||||
}
|
||||
}
|
||||
if len(seen) != 20 {
|
||||
t.Fatalf("claimed=%d want=20", len(seen))
|
||||
}
|
||||
}
|
||||
|
||||
func withSearchPath(dsn, schema string) string {
|
||||
if strings.Contains(dsn, "://") {
|
||||
parsed, err := url.Parse(dsn)
|
||||
if err == nil {
|
||||
query := parsed.Query()
|
||||
query.Set("search_path", schema)
|
||||
parsed.RawQuery = query.Encode()
|
||||
return parsed.String()
|
||||
}
|
||||
}
|
||||
return strings.TrimSpace(dsn) + " search_path=" + schema
|
||||
}
|
||||
@@ -0,0 +1,177 @@
|
||||
package outbox
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrDuplicateIdempotency = errors.New("可靠投递幂等键已存在")
|
||||
ErrLeaseLost = errors.New("可靠投递租约已失效")
|
||||
)
|
||||
|
||||
func Enqueue(tx *gorm.DB, input EnqueueInput, now time.Time) (Message, error) {
|
||||
input.InternalType, input.BusinessRef, input.IdempotencyKey = strings.TrimSpace(input.InternalType), strings.TrimSpace(input.BusinessRef), strings.TrimSpace(input.IdempotencyKey)
|
||||
if input.InternalType == "" || len(input.InternalType) > 64 || input.BusinessRef == "" || len(input.BusinessRef) > 128 || input.IdempotencyKey == "" || len(input.IdempotencyKey) > 191 || !json.Valid(input.PayloadJSON) {
|
||||
return Message{}, ErrInvalidInput
|
||||
}
|
||||
if input.MaxAttempts == 0 {
|
||||
input.MaxAttempts = 12
|
||||
}
|
||||
if input.MaxAttempts < 1 || input.MaxAttempts > 100 {
|
||||
return Message{}, ErrInvalidInput
|
||||
}
|
||||
now = now.UTC()
|
||||
message := Message{ID: uuid.NewString(), InternalType: input.InternalType, BusinessRef: input.BusinessRef, IdempotencyKey: input.IdempotencyKey, PayloadJSON: string(input.PayloadJSON), State: StatePending, MaxAttempts: input.MaxAttempts, AvailableAt: now, Version: 1, CreatedAt: now, UpdatedAt: now}
|
||||
var existing int64
|
||||
if err := tx.Model(&Message{}).Where("idempotency_key = ?", input.IdempotencyKey).Count(&existing).Error; err != nil {
|
||||
return Message{}, fmt.Errorf("check outbox idempotency: %w", err)
|
||||
}
|
||||
if existing > 0 {
|
||||
return Message{}, ErrDuplicateIdempotency
|
||||
}
|
||||
if err := tx.Create(&message).Error; err != nil {
|
||||
return Message{}, fmt.Errorf("enqueue outbox message: %w", err)
|
||||
}
|
||||
return message, nil
|
||||
}
|
||||
|
||||
type Relay struct {
|
||||
DB *gorm.DB
|
||||
Now func() time.Time
|
||||
LeaseDuration time.Duration
|
||||
Backoff func(int) time.Duration
|
||||
}
|
||||
|
||||
func NewRelay(db *gorm.DB) *Relay {
|
||||
return &Relay{DB: db, Now: time.Now, LeaseDuration: 30 * time.Second, Backoff: defaultBackoff}
|
||||
}
|
||||
|
||||
func (r *Relay) Claim(worker string, limit int) ([]Message, error) {
|
||||
worker = strings.TrimSpace(worker)
|
||||
if worker == "" || len(worker) > 128 || limit < 1 || limit > 100 {
|
||||
return nil, ErrInvalidInput
|
||||
}
|
||||
now, leaseUntil := r.now(), r.now().Add(r.leaseDuration())
|
||||
claimed := make([]Message, 0, limit)
|
||||
err := r.DB.Transaction(func(tx *gorm.DB) error {
|
||||
var candidates []Message
|
||||
query := tx.Where("((state IN ?) AND available_at <= ?) OR (state = ? AND lease_until < ?)", []string{StatePending, StateRetry}, now, StateProcessing, now).Order("available_at ASC, created_at ASC").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(&Message{}).Where("id = ? AND version = ?", item.ID, item.Version).Updates(map[string]interface{}{"state": 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 = StateProcessing, worker, &leaseUntil, item.Version+1
|
||||
claimed = append(claimed, item)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("claim outbox messages: %w", err)
|
||||
}
|
||||
return claimed, nil
|
||||
}
|
||||
|
||||
func (r *Relay) MarkSuccess(id, worker string) error {
|
||||
now := r.now()
|
||||
return r.DB.Transaction(func(tx *gorm.DB) error {
|
||||
var message Message
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&message, "id = ?", id).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
var delivered int64
|
||||
if err := tx.Model(&DeliveryRecord{}).Where("idempotency_key = ?", message.IdempotencyKey).Count(&delivered).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if delivered > 0 {
|
||||
return nil
|
||||
}
|
||||
if message.State != StateProcessing || message.LeaseOwner != worker || message.LeaseUntil == nil || !message.LeaseUntil.After(now) {
|
||||
return ErrLeaseLost
|
||||
}
|
||||
if err := tx.Create(&DeliveryRecord{IdempotencyKey: message.IdempotencyKey, MessageID: message.ID, DeliveredAt: now, CreatedAt: now}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
number := message.AttemptCount + 1
|
||||
if err := tx.Create(&Attempt{MessageID: message.ID, Number: number, Outcome: StateDelivered, Detail: "投递成功", Worker: worker, CreatedAt: now}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Model(&message).Updates(map[string]interface{}{"state": StateDelivered, "attempt_count": number, "delivered_at": now, "lease_owner": "", "lease_until": nil, "last_error": "", "version": gorm.Expr("version + 1"), "updated_at": now}).Error
|
||||
})
|
||||
}
|
||||
|
||||
func (r *Relay) MarkFailure(id, worker, detail string) error {
|
||||
now := r.now()
|
||||
detail = strings.TrimSpace(detail)
|
||||
if len(detail) > 512 {
|
||||
detail = detail[:512]
|
||||
}
|
||||
if detail == "" {
|
||||
detail = "投递失败"
|
||||
}
|
||||
return r.DB.Transaction(func(tx *gorm.DB) error {
|
||||
var message Message
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&message, "id = ?", id).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if message.State != StateProcessing || message.LeaseOwner != worker || message.LeaseUntil == nil || !message.LeaseUntil.After(now) {
|
||||
return ErrLeaseLost
|
||||
}
|
||||
nextAttempt := message.AttemptCount + 1
|
||||
state := StateRetry
|
||||
available := now.Add(r.backoff(nextAttempt))
|
||||
if nextAttempt >= message.MaxAttempts {
|
||||
state = StateDead
|
||||
available = now
|
||||
}
|
||||
if err := tx.Model(&message).Updates(map[string]interface{}{"state": state, "attempt_count": nextAttempt, "available_at": available, "lease_owner": "", "lease_until": nil, "last_error": detail, "version": gorm.Expr("version + 1"), "updated_at": now}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(&Attempt{MessageID: message.ID, Number: nextAttempt, Outcome: state, 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 (r *Relay) leaseDuration() time.Duration {
|
||||
if r.LeaseDuration <= 0 {
|
||||
return 30 * time.Second
|
||||
}
|
||||
return r.LeaseDuration
|
||||
}
|
||||
func (r *Relay) backoff(attempt int) time.Duration {
|
||||
if r.Backoff != nil {
|
||||
return r.Backoff(attempt)
|
||||
}
|
||||
return defaultBackoff(attempt)
|
||||
}
|
||||
func defaultBackoff(attempt int) time.Duration {
|
||||
if attempt < 1 {
|
||||
attempt = 1
|
||||
}
|
||||
delay := time.Second * time.Duration(1<<min(attempt-1, 8))
|
||||
if delay > 5*time.Minute {
|
||||
return 5 * time.Minute
|
||||
}
|
||||
return delay
|
||||
}
|
||||
@@ -0,0 +1,122 @@
|
||||
package outbox
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func testDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
db, err := gorm.Open(sqlite.Open("file:"+uuid.NewString()+"?mode=memory&cache=shared"), &gorm.Config{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
sqlDB, _ := db.DB()
|
||||
sqlDB.SetMaxOpenConns(1)
|
||||
if err = db.AutoMigrate(&Message{}, &DeliveryRecord{}, &Attempt{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return db
|
||||
}
|
||||
|
||||
func TestEnqueueClaimLeaseRecoveryBackoffAndIdempotentSuccess(t *testing.T) {
|
||||
db := testDB(t)
|
||||
now := time.Date(2026, 8, 28, 8, 0, 0, 0, time.UTC)
|
||||
message, err := Enqueue(db, EnqueueInput{InternalType: "local_event_candidate", BusinessRef: "event-1", IdempotencyKey: "local-event:event-1:v1", PayloadJSON: []byte(`{"event":"event-1"}`), MaxAttempts: 3}, now)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err = Enqueue(db, EnqueueInput{InternalType: "local_event_candidate", BusinessRef: "event-1", IdempotencyKey: message.IdempotencyKey, PayloadJSON: []byte(`{}`)}, now); !errors.Is(err, ErrDuplicateIdempotency) {
|
||||
t.Fatalf("expected duplicate error, got %v", err)
|
||||
}
|
||||
relay := NewRelay(db)
|
||||
relay.Now = func() time.Time { return now }
|
||||
relay.LeaseDuration = 10 * time.Second
|
||||
relay.Backoff = func(int) time.Duration { return 5 * time.Second }
|
||||
first, err := relay.Claim("worker-1", 1)
|
||||
if err != nil || len(first) != 1 {
|
||||
t.Fatalf("first claim=%#v err=%v", first, err)
|
||||
}
|
||||
second, err := relay.Claim("worker-2", 1)
|
||||
if err != nil || len(second) != 0 {
|
||||
t.Fatalf("concurrent claim=%#v err=%v", second, err)
|
||||
}
|
||||
now = now.Add(11 * time.Second)
|
||||
recovered, err := relay.Claim("worker-2", 1)
|
||||
if err != nil || len(recovered) != 1 || recovered[0].LeaseOwner != "worker-2" {
|
||||
t.Fatalf("recovered=%#v err=%v", recovered, err)
|
||||
}
|
||||
if err = relay.MarkFailure(message.ID, "worker-2", "temporary outage"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
now = now.Add(4 * time.Second)
|
||||
waiting, _ := relay.Claim("worker-3", 1)
|
||||
if len(waiting) != 0 {
|
||||
t.Fatalf("claimed before backoff elapsed: %#v", waiting)
|
||||
}
|
||||
now = now.Add(2 * time.Second)
|
||||
retry, err := relay.Claim("worker-3", 1)
|
||||
if err != nil || len(retry) != 1 {
|
||||
t.Fatalf("retry=%#v err=%v", retry, err)
|
||||
}
|
||||
if err = relay.MarkSuccess(message.ID, "worker-3"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err = relay.MarkSuccess(message.ID, "worker-3"); err != nil {
|
||||
t.Fatalf("idempotent success failed: %v", err)
|
||||
}
|
||||
var deliveries int64
|
||||
db.Model(&DeliveryRecord{}).Where("idempotency_key = ?", message.IdempotencyKey).Count(&deliveries)
|
||||
if deliveries != 1 {
|
||||
t.Fatalf("deliveries=%d", deliveries)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFailureBecomesDeadAndManualRequeuePreservesHistory(t *testing.T) {
|
||||
db := testDB(t)
|
||||
now := time.Date(2026, 8, 28, 9, 0, 0, 0, time.UTC)
|
||||
message, err := Enqueue(db, EnqueueInput{InternalType: "audit_projection", BusinessRef: "audit-1", IdempotencyKey: "audit:audit-1:v1", PayloadJSON: []byte(`{}`), MaxAttempts: 1}, now)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
relay := NewRelay(db)
|
||||
relay.Now = func() time.Time { return now }
|
||||
claimed, _ := relay.Claim("worker", 1)
|
||||
if len(claimed) != 1 {
|
||||
t.Fatal("message not claimed")
|
||||
}
|
||||
if err = relay.MarkFailure(message.ID, "worker", "permanent failure"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
service := NewService(db)
|
||||
service.Now = func() time.Time { return now.Add(time.Minute) }
|
||||
detail, err := service.Get(message.ID)
|
||||
if err != nil || detail.Message.State != StateDead || len(detail.Attempts) != 1 {
|
||||
t.Fatalf("detail=%#v err=%v", detail, err)
|
||||
}
|
||||
requeued, err := service.Requeue(message.ID, RequeueRequest{ExpectedVersion: detail.Message.Version, Reason: "出口故障已排除"}, 7)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if requeued.State != StatePending {
|
||||
t.Fatalf("state=%s", requeued.State)
|
||||
}
|
||||
detail, _ = service.Get(message.ID)
|
||||
if len(detail.Attempts) != 2 || detail.Attempts[0].Outcome != "manual_requeue" || detail.Attempts[0].ActorUserID != 7 {
|
||||
t.Fatalf("attempt history=%#v", detail.Attempts)
|
||||
}
|
||||
}
|
||||
|
||||
func TestProductionCannotCreateTestSink(t *testing.T) {
|
||||
if _, err := NewTestSink("prod"); !errors.Is(err, ErrTestSinkForbidden) {
|
||||
t.Fatalf("err=%v", err)
|
||||
}
|
||||
if _, err := NewTestSink("test"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,124 @@
|
||||
package outbox
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
coreService "github.com/go-admin-team/go-admin-core/sdk/service"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrInvalidInput = errors.New("可靠投递请求不符合要求")
|
||||
ErrNotFound = errors.New("可靠投递记录不存在")
|
||||
ErrNotDead = errors.New("仅死信记录可以重新排队")
|
||||
ErrVersionConflict = errors.New("记录已变化,请刷新后重试")
|
||||
)
|
||||
|
||||
type Service struct {
|
||||
coreService.Service
|
||||
Now func() time.Time
|
||||
}
|
||||
|
||||
func NewService(db *gorm.DB) *Service {
|
||||
return &Service{Service: coreService.Service{Orm: db}, Now: time.Now}
|
||||
}
|
||||
|
||||
func (s *Service) List(request PageRequest) (PageResponse, error) {
|
||||
request.State = strings.TrimSpace(request.State)
|
||||
request.InternalType = strings.TrimSpace(request.InternalType)
|
||||
request.Keyword = strings.TrimSpace(request.Keyword)
|
||||
if request.GetPageSize() > 100 || utf8.RuneCountInString(request.Keyword) > 128 || len(request.InternalType) > 64 || (request.State != "" && !validState(request.State)) {
|
||||
return PageResponse{}, ErrInvalidInput
|
||||
}
|
||||
query := s.Orm.Model(&Message{})
|
||||
if request.State != "" {
|
||||
query = query.Where("state = ?", request.State)
|
||||
}
|
||||
if request.InternalType != "" {
|
||||
query = query.Where("internal_type = ?", request.InternalType)
|
||||
}
|
||||
if request.Keyword != "" {
|
||||
pattern := "%" + strings.ToLower(request.Keyword) + "%"
|
||||
query = query.Where("LOWER(id) LIKE ? OR LOWER(idempotency_key) LIKE ? OR LOWER(business_ref) LIKE ?", pattern, pattern, pattern)
|
||||
}
|
||||
var response PageResponse
|
||||
if err := query.Count(&response.Count).Error; err != nil {
|
||||
return PageResponse{}, fmt.Errorf("count outbox messages: %w", err)
|
||||
}
|
||||
if err := query.Order("created_at DESC, id DESC").Limit(request.GetPageSize()).Offset((request.GetPageIndex() - 1) * request.GetPageSize()).Find(&response.List).Error; err != nil {
|
||||
return PageResponse{}, fmt.Errorf("list outbox messages: %w", err)
|
||||
}
|
||||
for state, target := range map[string]*int64{StatePending: &response.Summary.Pending, StateRetry: &response.Summary.Retry, StateProcessing: &response.Summary.Processing, StateDead: &response.Summary.Dead} {
|
||||
if err := s.Orm.Model(&Message{}).Where("state = ?", state).Count(target).Error; err != nil {
|
||||
return PageResponse{}, fmt.Errorf("summarize outbox: %w", err)
|
||||
}
|
||||
}
|
||||
return response, nil
|
||||
}
|
||||
|
||||
func (s *Service) Get(id string) (DetailResponse, error) {
|
||||
id = strings.TrimSpace(id)
|
||||
if id == "" || len(id) > 64 {
|
||||
return DetailResponse{}, ErrInvalidInput
|
||||
}
|
||||
var message Message
|
||||
if err := s.Orm.First(&message, "id = ?", id).Error; err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return DetailResponse{}, ErrNotFound
|
||||
}
|
||||
return DetailResponse{}, fmt.Errorf("get outbox message: %w", err)
|
||||
}
|
||||
var attempts []Attempt
|
||||
if err := s.Orm.Where("message_id = ?", id).Order("created_at DESC, id DESC").Limit(50).Find(&attempts).Error; err != nil {
|
||||
return DetailResponse{}, fmt.Errorf("list outbox attempts: %w", err)
|
||||
}
|
||||
return DetailResponse{Message: message, Attempts: attempts}, nil
|
||||
}
|
||||
|
||||
func (s *Service) Requeue(id string, request RequeueRequest, actorUserID int) (Message, error) {
|
||||
reason := strings.TrimSpace(request.Reason)
|
||||
if utf8.RuneCountInString(reason) < 6 || utf8.RuneCountInString(reason) > 256 || request.ExpectedVersion < 1 {
|
||||
return Message{}, ErrInvalidInput
|
||||
}
|
||||
now := s.now()
|
||||
var result Message
|
||||
err := s.Orm.Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&result, "id = ?", strings.TrimSpace(id)).Error; err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return ErrNotFound
|
||||
}
|
||||
return err
|
||||
}
|
||||
if result.Version != request.ExpectedVersion {
|
||||
return ErrVersionConflict
|
||||
}
|
||||
if result.State != StateDead {
|
||||
return ErrNotDead
|
||||
}
|
||||
result.State, result.AvailableAt, result.LastError = StatePending, now, ""
|
||||
result.LeaseOwner, result.LeaseUntil, result.Version = "", nil, result.Version+1
|
||||
if err := tx.Save(&result).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(&Attempt{MessageID: result.ID, Number: result.AttemptCount, Outcome: "manual_requeue", Detail: reason, ActorUserID: actorUserID, CreatedAt: now}).Error
|
||||
})
|
||||
if err != nil {
|
||||
return Message{}, err
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
func (s *Service) now() time.Time {
|
||||
if s.Now != nil {
|
||||
return s.Now().UTC()
|
||||
}
|
||||
return time.Now().UTC()
|
||||
}
|
||||
func validState(value string) bool {
|
||||
return value == StatePending || value == StateProcessing || value == StateRetry || value == StateDead || value == StateDelivered
|
||||
}
|
||||
@@ -0,0 +1,24 @@
|
||||
package outbox
|
||||
|
||||
import (
|
||||
"errors"
|
||||
"strings"
|
||||
)
|
||||
|
||||
var ErrTestSinkForbidden = errors.New("production 模式禁止启用测试接收器")
|
||||
|
||||
type Sink interface{ Deliver(Message) error }
|
||||
|
||||
type TestSink struct{ Delivered []string }
|
||||
|
||||
func (s *TestSink) Deliver(message Message) error {
|
||||
s.Delivered = append(s.Delivered, message.IdempotencyKey)
|
||||
return nil
|
||||
}
|
||||
|
||||
func NewTestSink(applicationMode string) (*TestSink, error) {
|
||||
if strings.EqualFold(strings.TrimSpace(applicationMode), "prod") || strings.EqualFold(strings.TrimSpace(applicationMode), "production") {
|
||||
return nil, ErrTestSinkForbidden
|
||||
}
|
||||
return &TestSink{}, nil
|
||||
}
|
||||
@@ -0,0 +1,65 @@
|
||||
package version
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"runtime"
|
||||
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration"
|
||||
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
|
||||
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
|
||||
)
|
||||
|
||||
func init() {
|
||||
_, fileName, _, _ := runtime.Caller(0)
|
||||
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSenseOutbox)
|
||||
}
|
||||
|
||||
func migrateSenseOutbox(db *gorm.DB, version string) error {
|
||||
return db.Transaction(func(tx *gorm.DB) error {
|
||||
if err := tx.AutoMigrate(&outbox.Message{}, &outbox.DeliveryRecord{}, &outbox.Attempt{}); err != nil {
|
||||
return err
|
||||
}
|
||||
root, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: senseLayoutMenuName, Title: "视频感知", Icon: "video-camera", Path: "/sense", MenuType: "M", Component: "Layout", Sort: 5, Visible: "0", IsFrame: "1"})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
page, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOutbox", Title: "可靠投递", Icon: "connection", Path: "outbox", Paths: fmt.Sprintf("/0/%d", root.MenuId), MenuType: "C", Permission: "sense:outbox:list", ParentId: root.MenuId, Component: "/sense/outbox/index", Sort: 11, Visible: "0", IsFrame: "1"})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
detail, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOutboxDetail", Title: "查看投递详情", MenuType: "F", Action: "GET", Permission: "sense:outbox:detail", ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 1, Visible: "1", IsFrame: "1"})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
requeue, err := ensureDeviceMenu(tx, migrationModels.SysMenu{MenuName: "SenseOutboxRequeue", Title: "死信重新排队", MenuType: "F", Action: "POST", Permission: "sense:outbox:requeue", ParentId: page.MenuId, Paths: fmt.Sprintf("/0/%d/%d", root.MenuId, page.MenuId), Sort: 2, Visible: "1", IsFrame: "1"})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
|
||||
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{page, detail}); err != nil {
|
||||
return err
|
||||
}
|
||||
for _, policy := range [][2]string{{"/api/v1/outbox", "GET"}, {"/api/v1/outbox/:id", "GET"}} {
|
||||
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: policy[0], V2: policy[1]}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
for _, role := range []string{"implementation_operator", "site_admin"} {
|
||||
if err = attachDeviceRole(tx, role, []migrationModels.SysMenu{requeue}); err != nil {
|
||||
return err
|
||||
}
|
||||
if err = tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&deviceCasbinRule{Ptype: "p", V0: role, V1: "/api/v1/outbox/:id/requeue", V2: "POST"}).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if err = rebuildSenseMenuPaths(tx, root.MenuId, "/0"); err != nil {
|
||||
return err
|
||||
}
|
||||
return tx.Create(&common.Migration{Version: version}).Error
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,41 @@
|
||||
package version
|
||||
|
||||
import (
|
||||
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/outbox"
|
||||
migrationModels "git.ilapage.cn/ila/yovision/Sense/server/cmd/migrate/migration/models"
|
||||
common "git.ilapage.cn/ila/yovision/Sense/server/common/models"
|
||||
"gorm.io/driver/sqlite"
|
||||
"gorm.io/gorm"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestOutboxMigrationAddsRBACWithoutFixtures(t *testing.T) {
|
||||
db, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{})
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err = db.AutoMigrate(&migrationModels.SysRole{}, &migrationModels.SysMenu{}, &deviceCasbinRule{}, &common.Migration{}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, role := range []string{"implementation_operator", "site_admin", "viewer"} {
|
||||
if err = db.Create(&migrationModels.SysRole{RoleName: role, RoleKey: role, Status: "2"}).Error; err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
const version = "2026082815000_outbox.go"
|
||||
if err = migrateSenseOutbox(db, version); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if !db.Migrator().HasTable(&outbox.Message{}) || !db.Migrator().HasTable(&outbox.DeliveryRecord{}) || !db.Migrator().HasTable(&outbox.Attempt{}) {
|
||||
t.Fatal("outbox tables missing")
|
||||
}
|
||||
var menus, reads, writes, fixtures, applied int64
|
||||
db.Model(&migrationModels.SysMenu{}).Where("menu_name LIKE ?", "SenseOutbox%").Count(&menus)
|
||||
db.Model(&deviceCasbinRule{}).Where("v1 LIKE ? AND v2 = ?", "/api/v1/outbox%", "GET").Count(&reads)
|
||||
db.Model(&deviceCasbinRule{}).Where("v1 = ? AND v2 = ?", "/api/v1/outbox/:id/requeue", "POST").Count(&writes)
|
||||
db.Model(&outbox.Message{}).Count(&fixtures)
|
||||
db.Model(&common.Migration{}).Where("version = ?", version).Count(&applied)
|
||||
if menus != 3 || reads != 6 || writes != 2 || fixtures != 0 || applied != 1 {
|
||||
t.Fatalf("menus=%d reads=%d writes=%d fixtures=%d applied=%d", menus, reads, writes, fixtures, applied)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,20 @@
|
||||
import request from '@/utils/request'
|
||||
|
||||
export function listOutbox(query) {
|
||||
return request({ url: '/api/v1/outbox', method: 'get', params: query })
|
||||
}
|
||||
|
||||
export function getOutbox(id) {
|
||||
return request({
|
||||
url: `/api/v1/outbox/${encodeURIComponent(id)}`,
|
||||
method: 'get'
|
||||
})
|
||||
}
|
||||
|
||||
export function requeueOutbox(id, expectedVersion, reason) {
|
||||
return request({
|
||||
url: `/api/v1/outbox/${encodeURIComponent(id)}/requeue`,
|
||||
method: 'post',
|
||||
data: { expectedVersion, reason }
|
||||
})
|
||||
}
|
||||
@@ -0,0 +1,447 @@
|
||||
<template>
|
||||
<BasicLayout>
|
||||
<template #wrapper>
|
||||
<el-card class="box-card">
|
||||
<div class="page-header">
|
||||
<div>
|
||||
<h3>可靠投递</h3>
|
||||
<p>查看 Sense 内部待投递记录、重试进度和人工恢复。</p>
|
||||
</div>
|
||||
<el-button
|
||||
:icon="Refresh"
|
||||
:loading="loading"
|
||||
@click="load"
|
||||
>刷新状态</el-button>
|
||||
</div>
|
||||
|
||||
<el-alert
|
||||
type="info"
|
||||
:closable="false"
|
||||
show-icon
|
||||
class="boundary-alert"
|
||||
>
|
||||
<template #title>未配置外部投递出口不影响 Sense 核心功能</template>
|
||||
内部记录会安全保留,正式 connector 由后续协调工单提供;测试接收器在
|
||||
production 模式不可启用。
|
||||
</el-alert>
|
||||
|
||||
<el-row :gutter="12" class="summary-row" aria-label="可靠投递队列概览">
|
||||
<el-col
|
||||
v-for="card in summaryCards"
|
||||
:key="card.key"
|
||||
:xs="12"
|
||||
:sm="6"
|
||||
><div class="summary-item">
|
||||
<span>{{ card.label }}</span><strong :class="card.className">{{ summary[card.key] }}</strong>
|
||||
</div></el-col>
|
||||
</el-row>
|
||||
|
||||
<el-form
|
||||
:model="query"
|
||||
label-width="76px"
|
||||
class="filter-form"
|
||||
@submit.prevent="search"
|
||||
>
|
||||
<el-form-item label="关键词"><el-input
|
||||
v-model="query.keyword"
|
||||
clearable
|
||||
placeholder="内部记录、业务引用或幂等键"
|
||||
@keyup.enter="search"
|
||||
/></el-form-item>
|
||||
<el-form-item label="状态"><el-select
|
||||
v-model="query.state"
|
||||
clearable
|
||||
placeholder="全部状态"
|
||||
><el-option
|
||||
v-for="(label, value) in stateLabels"
|
||||
:key="value"
|
||||
:label="label"
|
||||
:value="value"
|
||||
/></el-select></el-form-item>
|
||||
<el-form-item label="内部类型"><el-select
|
||||
v-model="query.internalType"
|
||||
clearable
|
||||
placeholder="全部类型"
|
||||
><el-option
|
||||
v-for="(label, value) in typeLabels"
|
||||
:key="value"
|
||||
:label="label"
|
||||
:value="value"
|
||||
/></el-select></el-form-item>
|
||||
<el-form-item class="filter-actions"><el-button
|
||||
type="primary"
|
||||
:icon="Search"
|
||||
native-type="submit"
|
||||
>查询</el-button><el-button
|
||||
:icon="RefreshLeft"
|
||||
@click="reset"
|
||||
>重置</el-button></el-form-item>
|
||||
</el-form>
|
||||
|
||||
<el-table
|
||||
v-loading="loading"
|
||||
:data="items"
|
||||
border
|
||||
stripe
|
||||
empty-text="当前筛选条件下没有可靠投递记录"
|
||||
>
|
||||
<el-table-column
|
||||
label="内部记录"
|
||||
min-width="255"
|
||||
><template #default="scope"><strong>{{ typeLabel(scope.row.internalType) }}</strong>
|
||||
<div class="code-text">{{ scope.row.id }}</div>
|
||||
<small>幂等键:{{ scope.row.idempotencyKey }}</small></template></el-table-column>
|
||||
<el-table-column
|
||||
label="状态"
|
||||
width="112"
|
||||
align="center"
|
||||
><template #default="scope"><el-tag :type="stateType(scope.row.state)">{{
|
||||
stateLabel(scope.row.state)
|
||||
}}</el-tag></template></el-table-column>
|
||||
<el-table-column
|
||||
label="尝试"
|
||||
width="74"
|
||||
align="center"
|
||||
><template #default="scope">{{ scope.row.attemptCount }} 次</template></el-table-column>
|
||||
<el-table-column
|
||||
label="下次动作"
|
||||
min-width="180"
|
||||
><template #default="scope">{{
|
||||
nextAction(scope.row)
|
||||
}}</template></el-table-column>
|
||||
<el-table-column
|
||||
label="租约"
|
||||
min-width="142"
|
||||
><template #default="scope"><div>{{ scope.row.leaseOwner || "未领取" }}</div>
|
||||
<small v-if="scope.row.leaseUntil">{{
|
||||
formatTime(scope.row.leaseUntil)
|
||||
}}</small></template></el-table-column>
|
||||
<el-table-column
|
||||
label="最近结果"
|
||||
min-width="210"
|
||||
><template #default="scope">{{
|
||||
scope.row.lastError ||
|
||||
(scope.row.state === "delivered" ? "投递成功" : "尚未失败")
|
||||
}}</template></el-table-column>
|
||||
<el-table-column
|
||||
label="操作"
|
||||
width="150"
|
||||
fixed="right"
|
||||
><template #default="scope"><el-button
|
||||
v-permisaction="['sense:outbox:detail']"
|
||||
type="primary"
|
||||
link
|
||||
@click="openDetail(scope.row.id)"
|
||||
>详情</el-button><el-button
|
||||
v-if="scope.row.state === 'dead'"
|
||||
v-permisaction="['sense:outbox:requeue']"
|
||||
type="warning"
|
||||
link
|
||||
@click="openRequeue(scope.row)"
|
||||
>重新排队</el-button></template></el-table-column>
|
||||
</el-table>
|
||||
<pagination
|
||||
v-show="total > 0"
|
||||
v-model:current-page="query.pageIndex"
|
||||
v-model:page-size="query.pageSize"
|
||||
:total="total"
|
||||
@pagination="load"
|
||||
/>
|
||||
<p class="boundary-text">
|
||||
页面不展示内部
|
||||
payload、外部凭据或机器身份;人工恢复保留原业务记录、幂等键和失败历史。
|
||||
</p>
|
||||
</el-card>
|
||||
|
||||
<el-dialog
|
||||
v-model="detailOpen"
|
||||
title="投递记录详情"
|
||||
width="min(780px, calc(100vw - 24px))"
|
||||
:close-on-click-modal="false"
|
||||
>
|
||||
<el-descriptions v-if="selected" :column="2" border>
|
||||
<el-descriptions-item label="内部记录">{{
|
||||
selected.message.id
|
||||
}}</el-descriptions-item><el-descriptions-item label="内部类型">{{
|
||||
typeLabel(selected.message.internalType)
|
||||
}}</el-descriptions-item>
|
||||
<el-descriptions-item
|
||||
label="幂等键"
|
||||
:span="2"
|
||||
><span class="code-text">{{
|
||||
selected.message.idempotencyKey
|
||||
}}</span></el-descriptions-item><el-descriptions-item label="当前状态"><el-tag :type="stateType(selected.message.state)">{{
|
||||
stateLabel(selected.message.state)
|
||||
}}</el-tag></el-descriptions-item><el-descriptions-item label="业务引用">{{
|
||||
selected.message.businessRef
|
||||
}}</el-descriptions-item>
|
||||
<el-descriptions-item label="最近结果" :span="2">{{
|
||||
selected.message.lastError || "尚未失败"
|
||||
}}</el-descriptions-item>
|
||||
</el-descriptions>
|
||||
<h4 class="timeline-title">处理时间线</h4>
|
||||
<el-timeline v-if="selected"><el-timeline-item
|
||||
v-for="attempt in selected.attempts"
|
||||
:key="attempt.id"
|
||||
:timestamp="formatTime(attempt.createdAt)"
|
||||
placement="top"
|
||||
><strong>第 {{ attempt.number }} 次 · {{ attempt.outcome }}</strong>
|
||||
<div>{{ attempt.detail }}</div>
|
||||
<small v-if="attempt.actorUserId">操作人编号:{{ attempt.actorUserId }}</small></el-timeline-item><el-timeline-item
|
||||
v-if="!selected.attempts.length"
|
||||
:timestamp="formatTime(selected.message.createdAt)"
|
||||
>与业务记录在同一事务中创建</el-timeline-item></el-timeline>
|
||||
<template #footer><el-button @click="detailOpen = false">关闭</el-button></template>
|
||||
</el-dialog>
|
||||
|
||||
<el-dialog
|
||||
v-model="requeueOpen"
|
||||
title="将死信重新排队"
|
||||
width="min(560px, calc(100vw - 24px))"
|
||||
:close-on-click-modal="false"
|
||||
@closed="resetRequeue"
|
||||
>
|
||||
<el-alert
|
||||
title="这是受控恢复操作"
|
||||
type="warning"
|
||||
:closable="false"
|
||||
show-icon
|
||||
>只创建新的投递尝试,不修改原业务记录或删除失败历史。执行前请先排除失败原因。</el-alert>
|
||||
<el-form
|
||||
ref="requeueFormRef"
|
||||
:model="requeueForm"
|
||||
:rules="requeueRules"
|
||||
label-position="top"
|
||||
class="requeue-form"
|
||||
><el-form-item
|
||||
label="恢复原因"
|
||||
prop="reason"
|
||||
><el-input
|
||||
v-model="requeueForm.reason"
|
||||
type="textarea"
|
||||
:rows="3"
|
||||
maxlength="256"
|
||||
show-word-limit
|
||||
placeholder="例如:出口配置已恢复,已核对幂等键和目标状态"
|
||||
/></el-form-item></el-form>
|
||||
<p class="boundary-text">
|
||||
原因将与操作人、对象版本和时间一起写入脱敏审计。
|
||||
</p>
|
||||
<template #footer><el-button @click="requeueOpen = false">取消</el-button><el-button
|
||||
type="warning"
|
||||
:loading="requeueLoading"
|
||||
@click="confirmRequeue"
|
||||
>确认重新排队</el-button></template>
|
||||
</el-dialog>
|
||||
</template>
|
||||
</BasicLayout>
|
||||
</template>
|
||||
|
||||
<script setup>
|
||||
import { computed, onMounted, reactive, ref } from 'vue'
|
||||
import { Refresh, RefreshLeft, Search } from '@element-plus/icons-vue'
|
||||
import { ElMessage } from 'element-plus'
|
||||
import { getOutbox, listOutbox, requeueOutbox } from '@/api/sense/outbox'
|
||||
import {
|
||||
buildOutboxQuery,
|
||||
formatTime,
|
||||
nextAction,
|
||||
stateLabel,
|
||||
stateLabels,
|
||||
stateType,
|
||||
typeLabel,
|
||||
typeLabels
|
||||
} from './outboxState'
|
||||
|
||||
defineOptions({ name: 'SenseOutbox' })
|
||||
const loading = ref(false)
|
||||
const items = ref([])
|
||||
const total = ref(0)
|
||||
const detailOpen = ref(false)
|
||||
const selected = ref(null)
|
||||
const requeueOpen = ref(false)
|
||||
const requeueLoading = ref(false)
|
||||
const requeueTarget = ref(null)
|
||||
const requeueFormRef = ref(null)
|
||||
const summary = reactive({ pending: 0, retry: 0, processing: 0, dead: 0 })
|
||||
const query = reactive({
|
||||
pageIndex: 1,
|
||||
pageSize: 10,
|
||||
state: '',
|
||||
internalType: '',
|
||||
keyword: ''
|
||||
})
|
||||
const requeueForm = reactive({ reason: '' })
|
||||
const requeueRules = {
|
||||
reason: [
|
||||
{ required: true, message: '请填写恢复原因', trigger: 'blur' },
|
||||
{ min: 6, max: 256, message: '请填写 6 至 256 个字符', trigger: 'blur' }
|
||||
]
|
||||
}
|
||||
const summaryCards = computed(() => [
|
||||
{ key: 'pending', label: '等待投递' },
|
||||
{ key: 'retry', label: '重试等待', className: 'warning-number' },
|
||||
{ key: 'processing', label: '处理中 / 租约' },
|
||||
{ key: 'dead', label: '死信待处理', className: 'danger-number' }
|
||||
])
|
||||
function unwrap(response) {
|
||||
return response?.data?.data ?? response?.data ?? response
|
||||
}
|
||||
async function load() {
|
||||
loading.value = true
|
||||
try {
|
||||
const payload = unwrap(await listOutbox(buildOutboxQuery(query))) || {}
|
||||
items.value = payload.list || []
|
||||
total.value = payload.count || 0
|
||||
Object.assign(
|
||||
summary,
|
||||
payload.summary || { pending: 0, retry: 0, processing: 0, dead: 0 }
|
||||
)
|
||||
} catch (error) {
|
||||
ElMessage.error(error.message || '可靠投递状态加载失败')
|
||||
} finally {
|
||||
loading.value = false
|
||||
}
|
||||
}
|
||||
function search() {
|
||||
query.pageIndex = 1
|
||||
load()
|
||||
}
|
||||
function reset() {
|
||||
Object.assign(query, {
|
||||
pageIndex: 1,
|
||||
state: '',
|
||||
internalType: '',
|
||||
keyword: ''
|
||||
})
|
||||
load()
|
||||
}
|
||||
async function openDetail(id) {
|
||||
try {
|
||||
selected.value = unwrap(await getOutbox(id))
|
||||
detailOpen.value = true
|
||||
} catch (error) {
|
||||
ElMessage.error(error.message || '投递详情加载失败')
|
||||
}
|
||||
}
|
||||
function openRequeue(item) {
|
||||
requeueTarget.value = item
|
||||
requeueOpen.value = true
|
||||
}
|
||||
function resetRequeue() {
|
||||
requeueForm.reason = ''
|
||||
requeueTarget.value = null
|
||||
requeueFormRef.value?.clearValidate()
|
||||
}
|
||||
async function confirmRequeue() {
|
||||
if (!requeueFormRef.value || !requeueTarget.value) return
|
||||
const valid = await requeueFormRef.value.validate().catch(() => false)
|
||||
if (!valid) return
|
||||
requeueLoading.value = true
|
||||
try {
|
||||
await requeueOutbox(
|
||||
requeueTarget.value.id,
|
||||
requeueTarget.value.version,
|
||||
requeueForm.reason.trim()
|
||||
)
|
||||
ElMessage.success('死信已重新排队,原失败历史已保留')
|
||||
requeueOpen.value = false
|
||||
detailOpen.value = false
|
||||
await load()
|
||||
} catch (error) {
|
||||
ElMessage.warning(error.message || '重新排队失败,请刷新状态后重试')
|
||||
} finally {
|
||||
requeueLoading.value = false
|
||||
}
|
||||
}
|
||||
onMounted(load)
|
||||
</script>
|
||||
|
||||
<style scoped>
|
||||
.page-header {
|
||||
display: flex;
|
||||
align-items: flex-start;
|
||||
justify-content: space-between;
|
||||
gap: 16px;
|
||||
}
|
||||
.page-header h3 {
|
||||
margin: 0 0 6px;
|
||||
}
|
||||
.page-header p {
|
||||
margin: 0;
|
||||
color: #909399;
|
||||
}
|
||||
.boundary-alert {
|
||||
margin: 16px 0;
|
||||
}
|
||||
.summary-row {
|
||||
margin-bottom: 18px;
|
||||
}
|
||||
.summary-item {
|
||||
display: flex;
|
||||
align-items: center;
|
||||
justify-content: space-between;
|
||||
min-height: 68px;
|
||||
padding: 12px 16px;
|
||||
border: 1px solid #ebeef5;
|
||||
border-radius: 4px;
|
||||
}
|
||||
.summary-item span {
|
||||
color: #606266;
|
||||
}
|
||||
.summary-item strong {
|
||||
font-size: 22px;
|
||||
}
|
||||
.warning-number {
|
||||
color: #e6a23c;
|
||||
}
|
||||
.danger-number {
|
||||
color: #f56c6c;
|
||||
}
|
||||
.filter-form {
|
||||
display: flex;
|
||||
align-items: flex-end;
|
||||
flex-wrap: wrap;
|
||||
gap: 0 12px;
|
||||
margin-bottom: 2px;
|
||||
}
|
||||
.filter-form .el-form-item {
|
||||
width: 260px;
|
||||
}
|
||||
.filter-form .filter-actions {
|
||||
width: auto;
|
||||
}
|
||||
.filter-form :deep(.el-select) {
|
||||
width: 100%;
|
||||
}
|
||||
.code-text {
|
||||
font-family: ui-monospace, SFMono-Regular, Consolas, monospace;
|
||||
overflow-wrap: anywhere;
|
||||
}
|
||||
.boundary-text,
|
||||
small {
|
||||
color: #909399;
|
||||
font-size: 12px;
|
||||
}
|
||||
.boundary-text {
|
||||
margin: 12px 0 0;
|
||||
}
|
||||
.timeline-title {
|
||||
margin: 20px 0 14px;
|
||||
}
|
||||
.requeue-form {
|
||||
margin-top: 16px;
|
||||
}
|
||||
@media (max-width: 768px) {
|
||||
.page-header {
|
||||
align-items: stretch;
|
||||
flex-direction: column;
|
||||
}
|
||||
.filter-form .el-form-item {
|
||||
width: 100%;
|
||||
}
|
||||
.summary-item {
|
||||
margin-bottom: 8px;
|
||||
}
|
||||
}
|
||||
</style>
|
||||
@@ -0,0 +1,54 @@
|
||||
export const stateLabels = {
|
||||
pending: '等待投递',
|
||||
processing: '处理中',
|
||||
retry: '重试等待',
|
||||
dead: '死信',
|
||||
delivered: '已投递'
|
||||
}
|
||||
export const typeLabels = {
|
||||
local_event_candidate: '本地事件候选',
|
||||
audit_projection: '审计投影'
|
||||
}
|
||||
|
||||
export function stateLabel(value) {
|
||||
return stateLabels[value] || value || '—'
|
||||
}
|
||||
export function stateType(value) {
|
||||
return (
|
||||
{
|
||||
pending: 'info',
|
||||
processing: 'primary',
|
||||
retry: 'warning',
|
||||
dead: 'danger',
|
||||
delivered: 'success'
|
||||
}[value] || 'info'
|
||||
)
|
||||
}
|
||||
export function typeLabel(value) {
|
||||
return typeLabels[value] || value || '—'
|
||||
}
|
||||
export function buildOutboxQuery(query) {
|
||||
return {
|
||||
pageIndex: query.pageIndex,
|
||||
pageSize: query.pageSize,
|
||||
state: query.state || undefined,
|
||||
internalType: query.internalType || undefined,
|
||||
keyword: String(query.keyword || '').trim() || undefined
|
||||
}
|
||||
}
|
||||
export function nextAction(item, now = Date.now()) {
|
||||
if (item.state === 'dead') return '等待人工处理'
|
||||
if (item.state === 'processing') { return item.leaseUntil ? `租约至 ${formatTime(item.leaseUntil)}` : '处理中' }
|
||||
if (item.state === 'retry') {
|
||||
return item.availableAt && new Date(item.availableAt).getTime() > now
|
||||
? `自动重试 ${formatTime(item.availableAt)}`
|
||||
: '等待重试领取'
|
||||
}
|
||||
if (item.state === 'delivered') return '已完成'
|
||||
return '等待 worker 领取'
|
||||
}
|
||||
export function formatTime(value) {
|
||||
return value
|
||||
? new Date(value).toLocaleString('zh-CN', { hour12: false })
|
||||
: '—'
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
import request from '@/utils/request'
|
||||
import { getOutbox, listOutbox, requeueOutbox } from '@/api/sense/outbox'
|
||||
|
||||
jest.mock('@/utils/request', () => jest.fn())
|
||||
|
||||
describe('Sense outbox API', () => {
|
||||
beforeEach(() => request.mockReset())
|
||||
test('encodes identifiers and sends only the recovery fields', () => {
|
||||
listOutbox({ pageIndex: 1, pageSize: 10 })
|
||||
getOutbox('outbox/id unsafe')
|
||||
requeueOutbox('outbox/id unsafe', 4, '出口故障已排除')
|
||||
expect(request).toHaveBeenNthCalledWith(1, {
|
||||
url: '/api/v1/outbox',
|
||||
method: 'get',
|
||||
params: { pageIndex: 1, pageSize: 10 }
|
||||
})
|
||||
expect(request).toHaveBeenNthCalledWith(2, {
|
||||
url: '/api/v1/outbox/outbox%2Fid%20unsafe',
|
||||
method: 'get'
|
||||
})
|
||||
expect(request).toHaveBeenNthCalledWith(3, {
|
||||
url: '/api/v1/outbox/outbox%2Fid%20unsafe/requeue',
|
||||
method: 'post',
|
||||
data: { expectedVersion: 4, reason: '出口故障已排除' }
|
||||
})
|
||||
})
|
||||
})
|
||||
@@ -0,0 +1,37 @@
|
||||
import {
|
||||
buildOutboxQuery,
|
||||
nextAction,
|
||||
stateLabel,
|
||||
stateType,
|
||||
typeLabel
|
||||
} from '@/views/sense/outbox/outboxState'
|
||||
|
||||
describe('Sense outbox presentation state', () => {
|
||||
test('uses operator-facing labels and non-color-only states', () => {
|
||||
expect(stateLabel('dead')).toBe('死信')
|
||||
expect(stateType('dead')).toBe('danger')
|
||||
expect(typeLabel('local_event_candidate')).toBe('本地事件候选')
|
||||
})
|
||||
test('builds an allowlisted trimmed query', () => {
|
||||
expect(
|
||||
buildOutboxQuery({
|
||||
pageIndex: 2,
|
||||
pageSize: 20,
|
||||
state: 'retry',
|
||||
internalType: 'local_event_candidate',
|
||||
keyword: ' key ',
|
||||
ignored: 'no'
|
||||
})
|
||||
).toEqual({
|
||||
pageIndex: 2,
|
||||
pageSize: 20,
|
||||
state: 'retry',
|
||||
internalType: 'local_event_candidate',
|
||||
keyword: 'key'
|
||||
})
|
||||
})
|
||||
test('explains the next recovery action', () => {
|
||||
expect(nextAction({ state: 'dead' })).toBe('等待人工处理')
|
||||
expect(nextAction({ state: 'pending' })).toBe('等待 worker 领取')
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user