diff --git a/Sense/server/app/admin/router/sense_outbox.go b/Sense/server/app/admin/router/sense_outbox.go new file mode 100644 index 0000000..603ca10 --- /dev/null +++ b/Sense/server/app/admin/router/sense_outbox.go @@ -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) +} diff --git a/Sense/server/app/admin/router/sense_outbox_test.go b/Sense/server/app/admin/router/sense_outbox_test.go new file mode 100644 index 0000000..44b1bf4 --- /dev/null +++ b/Sense/server/app/admin/router/sense_outbox_test.go @@ -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) + } + } +} diff --git a/Sense/server/app/sense/local_event/outbox.go b/Sense/server/app/sense/local_event/outbox.go new file mode 100644 index 0000000..cf94b3b --- /dev/null +++ b/Sense/server/app/sense/local_event/outbox.go @@ -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 + }) +} diff --git a/Sense/server/app/sense/local_event/outbox_test.go b/Sense/server/app/sense/local_event/outbox_test.go new file mode 100644 index 0000000..372aa70 --- /dev/null +++ b/Sense/server/app/sense/local_event/outbox_test.go @@ -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) + } +} diff --git a/Sense/server/app/sense/outbox/apis.go b/Sense/server/app/sense/outbox/apis.go new file mode 100644 index 0000000..455c5da --- /dev/null +++ b/Sense/server/app/sense/outbox/apis.go @@ -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, "可靠投递操作失败") + } +} diff --git a/Sense/server/app/sense/outbox/audit.go b/Sense/server/app/sense/outbox/audit.go new file mode 100644 index 0000000..5bf729b --- /dev/null +++ b/Sense/server/app/sense/outbox/audit.go @@ -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 +} diff --git a/Sense/server/app/sense/outbox/dto.go b/Sense/server/app/sense/outbox/dto.go new file mode 100644 index 0000000..1ffb9c3 --- /dev/null +++ b/Sense/server/app/sense/outbox/dto.go @@ -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 +} diff --git a/Sense/server/app/sense/outbox/models.go b/Sense/server/app/sense/outbox/models.go new file mode 100644 index 0000000..142eddb --- /dev/null +++ b/Sense/server/app/sense/outbox/models.go @@ -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" } diff --git a/Sense/server/app/sense/outbox/postgres_test.go b/Sense/server/app/sense/outbox/postgres_test.go new file mode 100644 index 0000000..ed424a0 --- /dev/null +++ b/Sense/server/app/sense/outbox/postgres_test.go @@ -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 +} diff --git a/Sense/server/app/sense/outbox/relay.go b/Sense/server/app/sense/outbox/relay.go new file mode 100644 index 0000000..b93dadb --- /dev/null +++ b/Sense/server/app/sense/outbox/relay.go @@ -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< 5*time.Minute { + return 5 * time.Minute + } + return delay +} diff --git a/Sense/server/app/sense/outbox/relay_test.go b/Sense/server/app/sense/outbox/relay_test.go new file mode 100644 index 0000000..35994f5 --- /dev/null +++ b/Sense/server/app/sense/outbox/relay_test.go @@ -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) + } +} diff --git a/Sense/server/app/sense/outbox/service.go b/Sense/server/app/sense/outbox/service.go new file mode 100644 index 0000000..3ac76d5 --- /dev/null +++ b/Sense/server/app/sense/outbox/service.go @@ -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 +} diff --git a/Sense/server/app/sense/outbox/sink.go b/Sense/server/app/sense/outbox/sink.go new file mode 100644 index 0000000..eaf7d12 --- /dev/null +++ b/Sense/server/app/sense/outbox/sink.go @@ -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 +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082815000_outbox.go b/Sense/server/cmd/migrate/migration/version/2026082815000_outbox.go new file mode 100644 index 0000000..ec3cb82 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082815000_outbox.go @@ -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 + }) +} diff --git a/Sense/server/cmd/migrate/migration/version/2026082815000_outbox_test.go b/Sense/server/cmd/migrate/migration/version/2026082815000_outbox_test.go new file mode 100644 index 0000000..aa6e7a6 --- /dev/null +++ b/Sense/server/cmd/migrate/migration/version/2026082815000_outbox_test.go @@ -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) + } +} diff --git a/Sense/ui/src/api/sense/outbox.js b/Sense/ui/src/api/sense/outbox.js new file mode 100644 index 0000000..61d9c9e --- /dev/null +++ b/Sense/ui/src/api/sense/outbox.js @@ -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 } + }) +} diff --git a/Sense/ui/src/views/sense/outbox/index.vue b/Sense/ui/src/views/sense/outbox/index.vue new file mode 100644 index 0000000..19884f5 --- /dev/null +++ b/Sense/ui/src/views/sense/outbox/index.vue @@ -0,0 +1,447 @@ + + + + + diff --git a/Sense/ui/src/views/sense/outbox/outboxState.js b/Sense/ui/src/views/sense/outbox/outboxState.js new file mode 100644 index 0000000..4790d89 --- /dev/null +++ b/Sense/ui/src/views/sense/outbox/outboxState.js @@ -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 }) + : '—' +} diff --git a/Sense/ui/tests/unit/sense/outboxApi.spec.js b/Sense/ui/tests/unit/sense/outboxApi.spec.js new file mode 100644 index 0000000..fa0dd52 --- /dev/null +++ b/Sense/ui/tests/unit/sense/outboxApi.spec.js @@ -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: '出口故障已排除' } + }) + }) +}) diff --git a/Sense/ui/tests/unit/sense/outboxState.spec.js b/Sense/ui/tests/unit/sense/outboxState.spec.js new file mode 100644 index 0000000..3cb6a23 --- /dev/null +++ b/Sense/ui/tests/unit/sense/outboxState.spec.js @@ -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 领取') + }) +})