173 lines
6.4 KiB
Go
173 lines
6.4 KiB
Go
package bell_alert_lifecycle_test
|
|
|
|
import (
|
|
"context"
|
|
"os"
|
|
"sync"
|
|
"sync/atomic"
|
|
"testing"
|
|
"time"
|
|
|
|
"gorm.io/driver/postgres"
|
|
"gorm.io/gorm"
|
|
|
|
adminmodels "go-admin/app/admin/models"
|
|
"go-admin/app/bell/alert"
|
|
"go-admin/app/bell/alert_lifecycle"
|
|
"go-admin/app/bell/event"
|
|
"go-admin/app/bell/rule"
|
|
)
|
|
|
|
func TestConcurrentLifecycleAndPersistence(t *testing.T) {
|
|
dsn := os.Getenv("BELL_ALERT_LIFECYCLE_TEST_DATABASE_URL")
|
|
if dsn == "" {
|
|
t.Skip("integration database not configured")
|
|
}
|
|
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
ctx := context.Background()
|
|
actors := createActors(t, db)
|
|
service := alert_lifecycle.NewService(db)
|
|
alertID := createAlert(t, db, "main")
|
|
const attempts = 20
|
|
var won atomic.Int32
|
|
results := make(chan alert_lifecycle.Result, attempts)
|
|
var wg sync.WaitGroup
|
|
for i := 0; i < attempts; i++ {
|
|
wg.Add(1)
|
|
actor := actors[i%2]
|
|
go func() {
|
|
defer wg.Done()
|
|
result, _ := service.Ack(ctx, alertID, actor)
|
|
if result.Won {
|
|
won.Add(1)
|
|
}
|
|
results <- result
|
|
}()
|
|
}
|
|
wg.Wait()
|
|
close(results)
|
|
if won.Load() != 1 {
|
|
t.Fatalf("ack winners=%d", won.Load())
|
|
}
|
|
detail, err := service.Get(ctx, alertID)
|
|
if err != nil || detail.Projection.Status != alert_lifecycle.StatusAcknowledged || len(detail.Timeline) != 1 {
|
|
t.Fatalf("ack projection=%#v err=%v", detail, err)
|
|
}
|
|
winner := actors[0]
|
|
loser := actors[1]
|
|
if detail.Projection.AcknowledgedBy == nil || *detail.Projection.AcknowledgedBy != winner.ID {
|
|
winner, loser = loser, winner
|
|
}
|
|
for result := range results {
|
|
if result.Detail.Projection.AcknowledgedBy != nil && *result.Detail.Projection.AcknowledgedBy != winner.ID {
|
|
t.Fatal("later ack did not report the true winner")
|
|
}
|
|
}
|
|
if _, err = service.Close(ctx, alertID, alert_lifecycle.CloseInput{Outcome: "site_normal"}, loser); err == nil {
|
|
t.Fatal("non-owner close succeeded")
|
|
}
|
|
if _, err = service.Close(ctx, alertID, alert_lifecycle.CloseInput{}, winner); err == nil {
|
|
t.Fatal("missing outcome succeeded")
|
|
}
|
|
closed, err := service.Close(ctx, alertID, alert_lifecycle.CloseInput{Outcome: "site_normal", Note: "现场正常"}, winner)
|
|
if err != nil || !closed.Won {
|
|
t.Fatalf("close failed: %#v %v", closed, err)
|
|
}
|
|
replay, err := service.Close(ctx, alertID, alert_lifecycle.CloseInput{Outcome: "site_normal", Note: "现场正常"}, winner)
|
|
if err != nil || !replay.Idempotent {
|
|
t.Fatalf("close replay=%#v %v", replay, err)
|
|
}
|
|
if _, err = service.Close(ctx, alertID, alert_lifecycle.CloseInput{Outcome: "false_positive"}, winner); err == nil {
|
|
t.Fatal("conflicting close replay succeeded")
|
|
}
|
|
if err = db.Model(&alert_lifecycle.Fact{}).Where("alert_id = ?", alertID).Update("actor_name", "tampered").Error; err == nil {
|
|
t.Fatal("lifecycle fact update succeeded")
|
|
}
|
|
var facts, rejects int64
|
|
db.Model(&alert_lifecycle.Fact{}).Where("alert_id = ?", alertID).Count(&facts)
|
|
db.Model(&alert_lifecycle.RejectionAudit{}).Where("alert_id = ?", alertID).Count(&rejects)
|
|
if facts != 2 || rejects < 20 {
|
|
t.Fatalf("facts=%d rejects=%d", facts, rejects)
|
|
}
|
|
sqlDB, _ := db.DB()
|
|
_ = sqlDB.Close()
|
|
reopened, err := gorm.Open(postgres.Open(dsn), &gorm.Config{})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
after, err := alert_lifecycle.NewService(reopened).Get(ctx, alertID)
|
|
if err != nil || after.Projection.Status != alert_lifecycle.StatusClosed || len(after.Timeline) != 2 {
|
|
t.Fatalf("restart detail=%#v err=%v", after, err)
|
|
}
|
|
|
|
rollbackID := createAlert(t, reopened, "rollback")
|
|
if err = reopened.Exec(`CREATE FUNCTION bell_test_reject_lifecycle() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'forced lifecycle failure'; END $$`).Error; err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err = reopened.Exec(`CREATE TRIGGER bell_test_reject_lifecycle BEFORE INSERT ON bell_alert_lifecycle_facts FOR EACH ROW EXECUTE FUNCTION bell_test_reject_lifecycle()`).Error; err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if _, err = alert_lifecycle.NewService(reopened).Ack(ctx, rollbackID, winner); err == nil {
|
|
t.Fatal("forced lifecycle failure succeeded")
|
|
}
|
|
rollback, _ := alert_lifecycle.NewService(reopened).Get(ctx, rollbackID)
|
|
if rollback.Projection.Status != alert_lifecycle.StatusOpen || len(rollback.Timeline) != 0 {
|
|
t.Fatal("failed ack left partial projection")
|
|
}
|
|
if err = reopened.Exec(`DROP TRIGGER bell_test_reject_lifecycle ON bell_alert_lifecycle_facts; DROP FUNCTION bell_test_reject_lifecycle()`).Error; err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
|
|
adminCloseID := createAlert(t, reopened, "admin-close")
|
|
if _, err = alert_lifecycle.NewService(reopened).Ack(ctx, adminCloseID, winner); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
administrator := loser
|
|
administrator.Role = "admin"
|
|
if result, closeErr := alert_lifecycle.NewService(reopened).Close(ctx, adminCloseID, alert_lifecycle.CloseInput{Outcome: "danger_confirmed"}, administrator); closeErr != nil || !result.Won {
|
|
t.Fatalf("administrator close failed: %#v %v", result, closeErr)
|
|
}
|
|
}
|
|
|
|
func createActors(t *testing.T, db *gorm.DB) []alert_lifecycle.Actor {
|
|
t.Helper()
|
|
var roleID int
|
|
db.Table("sys_role").Select("role_id").Where("role_key='operator'").Scan(&roleID)
|
|
result := make([]alert_lifecycle.Actor, 2)
|
|
for i := range result {
|
|
user := adminmodels.SysUser{Username: "bell_133_operator_" + string(rune('a'+i)), Password: "test-password-133", NickName: "处置员" + string(rune('A'+i)), RoleId: roleID, DeptId: 1, PostId: 1, Status: "2"}
|
|
if err := db.Create(&user).Error; err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
result[i] = alert_lifecycle.Actor{ID: user.UserId, Name: user.NickName, Role: "operator"}
|
|
}
|
|
return result
|
|
}
|
|
func createAlert(t *testing.T, db *gorm.DB, suffix string) string {
|
|
t.Helper()
|
|
ctx := context.Background()
|
|
eventType := "lifecycle_" + suffix
|
|
createdRule, err := rule.NewService(db).Create(ctx, rule.WriteInput{Code: "lifecycle-" + suffix, Name: "生命周期规则", EventType: &eventType, MinimumSeverity: "low"}, 1)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
created, err := event.NewService(db).Ingest(ctx, event.Command{ProducerID: "bell.lifecycle-test", SourceEventID: suffix, EventType: eventType, OccurredAt: time.Date(2026, 8, 29, 0, 0, 0, 0, time.UTC), Location: "测试地点" + suffix, Severity: "high", Attributes: map[string]any{}}, 1)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
items, _, err := alert.NewService(db).List(ctx, alert.PageQuery{PageIndex: 1, PageSize: 100})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
for _, item := range items {
|
|
if item.PrimaryRuleID == createdRule.ID {
|
|
return item.ID
|
|
}
|
|
}
|
|
t.Fatal("alert not created")
|
|
return created.Event.ID
|
|
}
|