Merge PR #135: Bell Event、Receipt 与合成事件入口

Implements #131; merge into dev for user acceptance.
This commit was merged in pull request #135.
This commit is contained in:
ila
2026-08-29 09:01:11 +08:00
18 changed files with 870 additions and 1 deletions
+190
View File
@@ -0,0 +1,190 @@
[CmdletBinding()]
param(
[string]$PostgresBin = 'D:\pgsql17\bin'
)
Set-StrictMode -Version 3.0
$ErrorActionPreference = 'Stop'
$server = $null
$pgStarted = $false
$testRoot = Join-Path ([IO.Path]::GetTempPath()) ('yovision-bell-131-' + [guid]::NewGuid().ToString('N'))
$pgData = Join-Path $testRoot 'postgres'
$pgLog = Join-Path $testRoot 'postgres.log'
$pgCtlOut = Join-Path $testRoot 'pg-ctl.out.log'
$pgCtlErr = Join-Path $testRoot 'pg-ctl.err.log'
$serverOut = Join-Path $testRoot 'bell.out.log'
$serverErr = Join-Path $testRoot 'bell.err.log'
$serverExe = Join-Path $testRoot 'bell-server.exe'
$serverRoot = (Resolve-Path (Join-Path $PSScriptRoot '..\server')).Path
function Get-FreeTcpPort {
$listener = [Net.Sockets.TcpListener]::new([Net.IPAddress]::Loopback, 0)
try {
$listener.Start()
return ([Net.IPEndPoint]$listener.LocalEndpoint).Port
} finally {
$listener.Stop()
}
}
function Wait-Tcp([int]$Port, [bool]$Open, [int]$Attempts = 120) {
for ($attempt = 0; $attempt -lt $Attempts; $attempt++) {
$client = [Net.Sockets.TcpClient]::new()
try {
$connected = $client.ConnectAsync('127.0.0.1', $Port).Wait(250) -and $client.Connected
} catch {
$connected = $false
} finally {
$client.Dispose()
}
if ($connected -eq $Open) { return }
Start-Sleep -Milliseconds 250
}
throw "TCP port $Port did not reach open=$Open"
}
function Wait-Health([string]$BaseUrl) {
for ($attempt = 0; $attempt -lt 100; $attempt++) {
try {
$health = Invoke-RestMethod -Uri "$BaseUrl/healthz" -TimeoutSec 2 -NoProxy
if ($health.status -eq 'ok' -and $health.service -eq 'bell') { return }
} catch {}
Start-Sleep -Milliseconds 300
}
throw 'Bell health endpoint did not become ready'
}
function Start-Bell([string]$BaseUrl, [string]$Config = 'config/settings.demo.yml') {
$script:server = Start-Process -FilePath $serverExe -ArgumentList @('server', '-c', $Config) -WorkingDirectory $serverRoot -RedirectStandardOutput $serverOut -RedirectStandardError $serverErr -WindowStyle Hidden -PassThru
Wait-Health $BaseUrl
}
function Stop-Bell {
if ($null -ne $script:server -and -not $script:server.HasExited) {
Stop-Process -Id $script:server.Id -Force
$script:server.WaitForExit(5000) | Out-Null
}
$script:server = $null
}
New-Item -ItemType Directory -Path $testRoot | Out-Null
$pgPort = Get-FreeTcpPort
$bellPort = Get-FreeTcpPort
$baseUrl = "http://127.0.0.1:$bellPort"
$database = 'bell_131'
$adminPassword = [guid]::NewGuid().ToString('N')
try {
foreach ($required in @('initdb.exe', 'pg_ctl.exe', 'createdb.exe', 'psql.exe')) {
$path = Join-Path $PostgresBin $required
if (-not (Test-Path -LiteralPath $path -PathType Leaf)) { throw "Missing PostgreSQL tool: $path" }
}
& (Join-Path $PostgresBin 'initdb.exe') -D $pgData -U postgres -A trust --encoding=UTF8 --no-locale | Out-Null
if ($LASTEXITCODE -ne 0) { throw 'isolated PostgreSQL initdb failed' }
$pgArguments = "-D `"$pgData`" -l `"$pgLog`" -o `"-p $pgPort -h 127.0.0.1`" start"
Start-Process -FilePath (Join-Path $PostgresBin 'pg_ctl.exe') -ArgumentList $pgArguments -RedirectStandardOutput $pgCtlOut -RedirectStandardError $pgCtlErr -WindowStyle Hidden | Out-Null
Wait-Tcp -Port $pgPort -Open $true
$pgStarted = $true
& (Join-Path $PostgresBin 'createdb.exe') -h 127.0.0.1 -p $pgPort -U postgres $database
if ($LASTEXITCODE -ne 0) { throw 'isolated Bell database creation failed' }
$env:GOTOOLCHAIN = 'go1.26.5'
$env:BELL_DATABASE_URL = "host=127.0.0.1 port=$pgPort user=postgres dbname=$database sslmode=disable"
$env:BELL_JWT_SECRET = [guid]::NewGuid().ToString('N') + [guid]::NewGuid().ToString('N')
$env:BELL_BOOTSTRAP_USERNAME = 'bell_131_admin'
$env:BELL_BOOTSTRAP_PASSWORD = $adminPassword
$env:BELL_HOST = '127.0.0.1'
$env:BELL_PORT = $bellPort.ToString()
$env:BELL_SYNTHETIC_EVENTS_ENABLED = 'true'
Push-Location $serverRoot
try {
go run . migrate -c config/settings.demo.yml *> (Join-Path $testRoot 'migrate.log')
if ($LASTEXITCODE -ne 0) { throw 'Bell migration failed' }
go build -o $serverExe .
if ($LASTEXITCODE -ne 0) { throw 'Bell build failed' }
} finally {
Pop-Location
}
Start-Bell $baseUrl
$unauthorized = Invoke-RestMethod -Method Post -Uri "$baseUrl/api/v1/bell/synthetic-events" -ContentType 'application/json' -Body '{}' -TimeoutSec 5 -NoProxy
if ([int]$unauthorized.code -ne 401) { throw "unauthenticated request returned code $($unauthorized.code)" }
$loginBody = @{ username = 'bell_131_admin'; password = $adminPassword; code = '0'; uuid = '0' } | ConvertTo-Json -Compress
$login = Invoke-RestMethod -Method Post -Uri "$baseUrl/api/v1/login" -ContentType 'application/json' -Body $loginBody -TimeoutSec 5 -NoProxy
if ([int]$login.code -ne 200 -or [string]::IsNullOrWhiteSpace($login.token)) { throw 'Bell login failed' }
$headers = @{ Authorization = "Bearer $($login.token)" }
$occurredAt = '2026-08-29T00:00:00Z'
$eventBody = @{
sourceEventId = 'acceptance-001'; eventType = 'danger_area_entered'; occurredAt = $occurredAt
location = '东门'; severity = 'high'; evidenceRef = 'evidence/acceptance-001'
attributes = @{ rule = 'area-01'; target = 'anonymous' }
} | ConvertTo-Json -Depth 6 -Compress
$created = Invoke-RestMethod -Method Post -Uri "$baseUrl/api/v1/bell/synthetic-events" -Headers $headers -ContentType 'application/json; charset=utf-8' -Body $eventBody -TimeoutSec 10 -NoProxy
if ([int]$created.code -ne 200 -or $created.data.duplicate -ne $false) { throw 'first synthetic event was not accepted as new' }
$eventId = [string]$created.data.event.id
$receiptId = [string]$created.data.receipt.id
if ([string]::IsNullOrWhiteSpace($eventId) -or [string]::IsNullOrWhiteSpace($receiptId)) { throw 'event or receipt id missing' }
$replay = Invoke-RestMethod -Method Post -Uri "$baseUrl/api/v1/bell/synthetic-events" -Headers $headers -ContentType 'application/json; charset=utf-8' -Body $eventBody -TimeoutSec 10 -NoProxy
if ([int]$replay.code -ne 200 -or $replay.data.duplicate -ne $true -or $replay.data.event.id -ne $eventId) { throw 'idempotent replay failed' }
$concurrentBody = @{
sourceEventId = 'acceptance-concurrent'; eventType = 'danger_area_entered'; occurredAt = $occurredAt
location = '西门'; severity = 'medium'; evidenceRef = 'evidence/acceptance-concurrent'
attributes = @{ rule = 'area-02'; target = 'anonymous' }
} | ConvertTo-Json -Depth 6 -Compress
$authHeader = [string]$headers.Authorization
$concurrent = 1..12 | ForEach-Object -Parallel {
$requestHeaders = @{ Authorization = $using:authHeader }
Invoke-RestMethod -Method Post -Uri "$using:baseUrl/api/v1/bell/synthetic-events" -Headers $requestHeaders -ContentType 'application/json; charset=utf-8' -Body $using:concurrentBody -TimeoutSec 20 -NoProxy
} -ThrottleLimit 12
$concurrentIds = @($concurrent | ForEach-Object { [string]$_.data.event.id } | Sort-Object -Unique)
$newCount = @($concurrent | Where-Object { $_.data.duplicate -eq $false }).Count
if ($concurrent.Count -ne 12 -or $concurrentIds.Count -ne 1 -or $newCount -ne 1) { throw 'concurrent idempotency did not converge to one Event' }
$conflictObject = $eventBody | ConvertFrom-Json
$conflictObject.severity = 'critical'
$conflict = Invoke-RestMethod -Method Post -Uri "$baseUrl/api/v1/bell/synthetic-events" -Headers $headers -ContentType 'application/json; charset=utf-8' -Body ($conflictObject | ConvertTo-Json -Depth 6 -Compress) -TimeoutSec 10 -NoProxy
if ([int]$conflict.code -ne 409) { throw "idempotency conflict returned code $($conflict.code)" }
$malformed = Invoke-RestMethod -Method Post -Uri "$baseUrl/api/v1/bell/synthetic-events" -Headers $headers -ContentType 'application/json' -Body '{' -TimeoutSec 10 -NoProxy
if ([int]$malformed.code -ne 400) { throw "malformed request returned code $($malformed.code)" }
$oversized = @{ sourceEventId='too-large'; eventType='test'; occurredAt=$occurredAt; location='lab'; severity='low'; attributes=@{ blob=('x' * 70000) } } | ConvertTo-Json -Depth 5 -Compress
$tooLarge = Invoke-RestMethod -Method Post -Uri "$baseUrl/api/v1/bell/synthetic-events" -Headers $headers -ContentType 'application/json' -Body $oversized -TimeoutSec 10 -NoProxy
if ([int]$tooLarge.code -ne 400) { throw "oversized request returned code $($tooLarge.code)" }
$psql = Join-Path $PostgresBin 'psql.exe'
$counts = & $psql -h 127.0.0.1 -p $pgPort -U postgres -d $database -Atc "select count(*) from bell_events; select count(*) from bell_event_receipts; select count(*) from bell_event_ingest_audits where outcome='accepted'; select count(*) from bell_event_ingest_audits where outcome='replay'; select count(*) from bell_event_ingest_audits where outcome='conflict';"
if (($counts -join ',') -ne '2,2,2,12,1') { throw "unexpected Bell fact counts: $($counts -join ',')" }
$payloadAuditCount = & $psql -h 127.0.0.1 -p $pgPort -U postgres -d $database -Atc "select count(*) from sys_opera_log where oper_url='/api/v1/bell/synthetic-events' and (oper_param like '%sourceEventId%' or json_result like '%acceptance-001%');"
if ([int]$payloadAuditCount -ne 0) { throw 'generic GoAdmin audit retained synthetic event payload' }
& $psql -h 127.0.0.1 -p $pgPort -U postgres -d $database -v ON_ERROR_STOP=1 -c "update bell_events set severity='low' where id='$eventId'" 2>$null | Out-Null
if ($LASTEXITCODE -eq 0) { throw 'immutable Event update unexpectedly succeeded' }
& $psql -h 127.0.0.1 -p $pgPort -U postgres -d $database -v ON_ERROR_STOP=1 -c "delete from bell_event_receipts where id='$receiptId'" 2>$null | Out-Null
if ($LASTEXITCODE -eq 0) { throw 'immutable Receipt delete unexpectedly succeeded' }
Stop-Bell
Start-Bell $baseUrl
$afterRestart = Invoke-RestMethod -Uri "$baseUrl/api/v1/bell/events/$eventId" -Headers $headers -TimeoutSec 10 -NoProxy
if ([int]$afterRestart.code -ne 200 -or $afterRestart.data.id -ne $eventId) { throw 'Event was not readable after Bell restart' }
Stop-Bell
Start-Bell $baseUrl 'config/settings.yml'
$productionSynthetic = Invoke-WebRequest -Method Post -Uri "$baseUrl/api/v1/bell/synthetic-events" -ContentType 'application/json' -Body '{}' -TimeoutSec 10 -NoProxy -SkipHttpErrorCheck
if ([int]$productionSynthetic.StatusCode -ne 404) { throw "production synthetic route returned HTTP $($productionSynthetic.StatusCode)" }
Write-Output "BELL_131_SMOKE events=2 receipts=2 concurrent=12 audit=redacted immutable=true restart=true production_synthetic=404"
} finally {
Stop-Bell
if ($pgStarted) {
& (Join-Path $PostgresBin 'pg_ctl.exe') -D $pgData -m fast stop *> (Join-Path $testRoot 'pg-stop.log')
}
foreach ($name in @('BELL_DATABASE_URL','BELL_JWT_SECRET','BELL_BOOTSTRAP_USERNAME','BELL_BOOTSTRAP_PASSWORD','BELL_HOST','BELL_PORT','BELL_SYNTHETIC_EVENTS_ENABLED')) {
Remove-Item "Env:$name" -ErrorAction SilentlyContinue
}
Write-Verbose "Bell #131 temporary artifacts: $testRoot"
}
+9
View File
@@ -0,0 +1,9 @@
package event
import "errors"
var (
ErrInvalid = errors.New("事件字段无效")
ErrConflict = errors.New("幂等键已用于不同事件载荷")
ErrNotFound = errors.New("事件不存在")
)
+24
View File
@@ -0,0 +1,24 @@
package event
import (
"encoding/json"
"time"
)
// Event is an immutable Bell business fact. It intentionally does not embed
// GoAdmin's mutable/soft-delete model fields.
type Event struct {
ID string `json:"id" gorm:"type:uuid;primaryKey"`
ProducerID string `json:"producerId" gorm:"size:128;not null;uniqueIndex:bell_event_key"`
SourceEventID string `json:"sourceEventId" gorm:"size:256;not null;uniqueIndex:bell_event_key"`
EventType string `json:"eventType" gorm:"size:128;not null;index"`
OccurredAt time.Time `json:"occurredAt" gorm:"type:timestamptz;not null;index"`
Location string `json:"location" gorm:"size:256;not null;index"`
Severity string `json:"severity" gorm:"size:16;not null;index"`
EvidenceRef *string `json:"evidenceRef,omitempty" gorm:"size:512"`
NormalizedPayload json.RawMessage `json:"normalizedPayload" gorm:"type:jsonb;not null"`
PayloadSHA256 string `json:"-" gorm:"type:char(64);not null"`
ReceivedAt time.Time `json:"receivedAt" gorm:"type:timestamptz;not null;index"`
}
func (Event) TableName() string { return "bell_events" }
+105
View File
@@ -0,0 +1,105 @@
package event
import (
"bytes"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"fmt"
"strings"
"time"
"unicode/utf8"
)
const (
maxProducerID = 128
maxSourceEventID = 256
maxEventType = 128
maxLocation = 256
maxEvidenceRef = 512
maxAttributes = 48 * 1024
)
type Command struct {
ProducerID string `json:"producerId"`
SourceEventID string `json:"sourceEventId"`
EventType string `json:"eventType"`
OccurredAt time.Time `json:"occurredAt"`
Location string `json:"location"`
Severity string `json:"severity"`
EvidenceRef *string `json:"evidenceRef,omitempty"`
Attributes map[string]any `json:"attributes"`
}
type Normalized struct {
Command Command
Payload []byte
Digest string
}
func Normalize(command Command) (Normalized, error) {
command.ProducerID = strings.TrimSpace(command.ProducerID)
command.SourceEventID = strings.TrimSpace(command.SourceEventID)
command.EventType = strings.TrimSpace(command.EventType)
command.Location = strings.TrimSpace(command.Location)
command.Severity = strings.ToLower(strings.TrimSpace(command.Severity))
command.OccurredAt = command.OccurredAt.UTC()
if !validText(command.ProducerID, maxProducerID) ||
!validText(command.SourceEventID, maxSourceEventID) ||
!validText(command.EventType, maxEventType) ||
!validText(command.Location, maxLocation) || command.OccurredAt.IsZero() {
return Normalized{}, ErrInvalid
}
switch command.Severity {
case "low", "medium", "high", "critical":
default:
return Normalized{}, ErrInvalid
}
if command.EvidenceRef != nil {
value := strings.TrimSpace(*command.EvidenceRef)
lower := strings.ToLower(value)
if !validText(value, maxEvidenceRef) || strings.Contains(value, "@") ||
strings.Contains(value, "\\") || strings.HasPrefix(lower, "file:") {
return Normalized{}, fmt.Errorf("%w: evidenceRef 必须是安全的逻辑引用", ErrInvalid)
}
command.EvidenceRef = &value
}
if command.Attributes == nil {
command.Attributes = map[string]any{}
}
attributes, err := canonicalJSON(command.Attributes)
if err != nil || len(attributes) > maxAttributes {
return Normalized{}, fmt.Errorf("%w: attributes 无效或过大", ErrInvalid)
}
// Round trip the attributes so map ordering and nested values have one
// deterministic representation before the complete command is hashed.
if err = json.Unmarshal(attributes, &command.Attributes); err != nil {
return Normalized{}, fmt.Errorf("%w: attributes 无效", ErrInvalid)
}
payload, err := json.Marshal(command)
if err != nil {
return Normalized{}, fmt.Errorf("normalize event: %w", err)
}
digest := sha256.Sum256(payload)
return Normalized{Command: command, Payload: payload, Digest: hex.EncodeToString(digest[:])}, nil
}
func canonicalJSON(value any) ([]byte, error) {
raw, err := json.Marshal(value)
if err != nil {
return nil, err
}
decoder := json.NewDecoder(bytes.NewReader(raw))
decoder.UseNumber()
var normalized any
if err = decoder.Decode(&normalized); err != nil {
return nil, err
}
return json.Marshal(normalized)
}
func validText(value string, max int) bool {
return value != "" && utf8.ValidString(value) && utf8.RuneCountInString(value) <= max &&
!strings.ContainsAny(value, "\x00\r\n")
}
+110
View File
@@ -0,0 +1,110 @@
package event
import (
"context"
"errors"
"fmt"
"time"
"github.com/google/uuid"
"gorm.io/gorm"
"go-admin/app/bell/receipt"
)
type Result struct {
Event Event `json:"event"`
Receipt receipt.Receipt `json:"receipt"`
Duplicate bool `json:"duplicate"`
}
type Service struct{ DB *gorm.DB }
func NewService(db *gorm.DB) Service { return Service{DB: db} }
func (s Service) Ingest(ctx context.Context, command Command, actorID int) (Result, error) {
if s.DB == nil {
return Result{}, errors.New("event database is unavailable")
}
normalized, err := Normalize(command)
if err != nil {
return Result{}, err
}
var result Result
err = s.DB.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
lockKey := fmt.Sprintf("%d:%s:%s", len(normalized.Command.ProducerID), normalized.Command.ProducerID, normalized.Command.SourceEventID)
if err := tx.Exec("SELECT pg_advisory_xact_lock(hashtextextended(?, 0))", lockKey).Error; err != nil {
return fmt.Errorf("lock event idempotency key: %w", err)
}
var existing receipt.Receipt
err := tx.Where("producer_id = ? AND source_event_id = ?", normalized.Command.ProducerID, normalized.Command.SourceEventID).
First(&existing).Error
if err == nil {
if existing.PayloadSHA256 != normalized.Digest {
return ErrConflict
}
if err = tx.First(&result.Event, "id = ?", existing.EventID).Error; err != nil {
return fmt.Errorf("load replay event: %w", err)
}
result.Receipt = existing
result.Duplicate = true
return tx.Create(&receipt.IngestAudit{
ProducerID: normalized.Command.ProducerID, SourceEventID: normalized.Command.SourceEventID,
PayloadSHA256: normalized.Digest, Outcome: receipt.OutcomeReplay, ActorID: actorID, CreatedAt: time.Now().UTC(),
}).Error
}
if !errors.Is(err, gorm.ErrRecordNotFound) {
return fmt.Errorf("load event receipt: %w", err)
}
now := time.Now().UTC()
result.Event = Event{
ID: uuid.NewString(), ProducerID: normalized.Command.ProducerID,
SourceEventID: normalized.Command.SourceEventID, EventType: normalized.Command.EventType,
OccurredAt: normalized.Command.OccurredAt, Location: normalized.Command.Location,
Severity: normalized.Command.Severity, EvidenceRef: normalized.Command.EvidenceRef,
NormalizedPayload: normalized.Payload, PayloadSHA256: normalized.Digest, ReceivedAt: now,
}
result.Receipt = receipt.Receipt{
ID: uuid.NewString(), EventID: result.Event.ID, ProducerID: result.Event.ProducerID,
SourceEventID: result.Event.SourceEventID, PayloadSHA256: normalized.Digest, AcceptedAt: now,
}
if err = tx.Create(&result.Event).Error; err != nil {
return fmt.Errorf("create event: %w", err)
}
if err = tx.Create(&result.Receipt).Error; err != nil {
return fmt.Errorf("create event receipt: %w", err)
}
return tx.Create(&receipt.IngestAudit{
ProducerID: result.Event.ProducerID, SourceEventID: result.Event.SourceEventID,
PayloadSHA256: normalized.Digest, Outcome: receipt.OutcomeAccepted, ActorID: actorID, CreatedAt: now,
}).Error
})
if errors.Is(err, ErrConflict) {
// The conflict audit must commit independently from the rejected ingest.
auditErr := s.DB.WithContext(ctx).Create(&receipt.IngestAudit{
ProducerID: normalized.Command.ProducerID, SourceEventID: normalized.Command.SourceEventID,
PayloadSHA256: normalized.Digest, Outcome: receipt.OutcomeConflict, ActorID: actorID, CreatedAt: time.Now().UTC(),
}).Error
if auditErr != nil {
return Result{}, fmt.Errorf("record idempotency conflict audit: %w", auditErr)
}
return Result{}, ErrConflict
}
return result, err
}
func (s Service) Get(ctx context.Context, id string) (Event, error) {
if _, err := uuid.Parse(id); err != nil {
return Event{}, ErrNotFound
}
var item Event
if err := s.DB.WithContext(ctx).First(&item, "id = ?", id).Error; err != nil {
if errors.Is(err, gorm.ErrRecordNotFound) {
return Event{}, ErrNotFound
}
return Event{}, err
}
return item, nil
}
+23
View File
@@ -0,0 +1,23 @@
package receipt
import "time"
const (
OutcomeAccepted = "accepted"
OutcomeReplay = "replay"
OutcomeConflict = "conflict"
)
// IngestAudit contains only identifiers, a digest and an outcome. Event payloads,
// credentials and tokens must never be copied into this append-only audit table.
type IngestAudit struct {
ID int64 `json:"id" gorm:"primaryKey;autoIncrement"`
ProducerID string `json:"producerId" gorm:"size:128;not null;index"`
SourceEventID string `json:"sourceEventId" gorm:"size:256;not null;index"`
PayloadSHA256 string `json:"-" gorm:"type:char(64);not null"`
Outcome string `json:"outcome" gorm:"size:16;not null"`
ActorID int `json:"actorId" gorm:"not null"`
CreatedAt time.Time `json:"createdAt" gorm:"type:timestamptz;not null"`
}
func (IngestAudit) TableName() string { return "bell_event_ingest_audits" }
+16
View File
@@ -0,0 +1,16 @@
package receipt
import "time"
// Receipt permanently binds an idempotency key to the accepted Event payload.
// It deliberately has no UpdatedAt or soft-delete fields because it is immutable.
type Receipt struct {
ID string `json:"id" gorm:"type:uuid;primaryKey"`
EventID string `json:"eventId" gorm:"type:uuid;not null;uniqueIndex"`
ProducerID string `json:"producerId" gorm:"size:128;not null;uniqueIndex:bell_receipt_key"`
SourceEventID string `json:"sourceEventId" gorm:"size:256;not null;uniqueIndex:bell_receipt_key"`
PayloadSHA256 string `json:"-" gorm:"type:char(64);not null"`
AcceptedAt time.Time `json:"acceptedAt" gorm:"type:timestamptz;not null"`
}
func (Receipt) TableName() string { return "bell_event_receipts" }
+38
View File
@@ -0,0 +1,38 @@
package router
import (
"errors"
"net/http"
"github.com/gin-gonic/gin"
"github.com/go-admin-team/go-admin-core/sdk/api"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"go-admin/app/bell/event"
)
type eventHandler struct{ api.Api }
func registerEventRouter(v1 *gin.RouterGroup, authMiddleware *jwt.GinJWTMiddleware) {
handler := eventHandler{}
v1.GET("/events/:id", authMiddleware.MiddlewareFunc(), handler.Get)
}
func (h eventHandler) Get(c *gin.Context) {
h.MakeContext(c).MakeOrm()
if h.Errors != nil {
h.Error(http.StatusInternalServerError, errors.New("数据库连接获取失败"), "数据库连接获取失败")
return
}
item, err := event.NewService(h.Orm).Get(c.Request.Context(), c.Param("id"))
if err != nil {
if errors.Is(err, event.ErrNotFound) {
h.Error(http.StatusNotFound, event.ErrNotFound, event.ErrNotFound.Error())
return
}
h.Logger.Errorf("load Bell event failed: %v", err)
h.Error(http.StatusInternalServerError, errors.New("读取事件失败"), "读取事件失败")
return
}
h.OK(item, "查询成功")
}
+38
View File
@@ -0,0 +1,38 @@
package router
import (
"os"
"github.com/gin-gonic/gin"
log "github.com/go-admin-team/go-admin-core/logger"
"github.com/go-admin-team/go-admin-core/sdk"
"github.com/go-admin-team/go-admin-core/sdk/config"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"go-admin/app/bell/synthetic"
"go-admin/common/middleware"
)
type registrar func(*gin.RouterGroup, *jwt.GinJWTMiddleware)
var registrars = []registrar{registerEventRouter}
func InitRouter() {
engine, ok := sdk.Runtime.GetEngine().(*gin.Engine)
if !ok || engine == nil {
log.Error("Bell business router requires Gin engine")
return
}
authMiddleware, err := middleware.AuthInit()
if err != nil {
log.Errorf("Bell business JWT init error: %v", err)
return
}
v1 := engine.Group("/api/v1/bell")
for _, register := range registrars {
register(v1, authMiddleware)
}
if synthetic.Enabled(config.ApplicationConfig.Mode, os.Getenv) {
registerSyntheticRouter(v1, authMiddleware)
}
}
+13
View File
@@ -0,0 +1,13 @@
package router
import (
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"go-admin/app/bell/synthetic"
)
func registerSyntheticRouter(v1 *gin.RouterGroup, authMiddleware *jwt.GinJWTMiddleware) {
handler := synthetic.Handler{}
v1.POST("/synthetic-events", authMiddleware.MiddlewareFunc(), handler.Create)
}
+19
View File
@@ -0,0 +1,19 @@
package synthetic
import (
"os"
"strings"
)
const EnabledEnv = "BELL_SYNTHETIC_EVENTS_ENABLED"
// Enabled requires both a non-production runtime and an explicit opt-in.
func Enabled(mode string, getenv func(string) string) bool {
if strings.EqualFold(strings.TrimSpace(mode), "prod") || strings.EqualFold(strings.TrimSpace(mode), "production") {
return false
}
value := strings.ToLower(strings.TrimSpace(getenv(EnabledEnv)))
return value == "true" || value == "1"
}
func EnabledFromEnvironment(mode string) bool { return Enabled(mode, os.Getenv) }
+58
View File
@@ -0,0 +1,58 @@
package synthetic
import (
"errors"
"net/http"
"github.com/gin-gonic/gin"
"github.com/gin-gonic/gin/binding"
"github.com/go-admin-team/go-admin-core/sdk/api"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth/user"
"go-admin/app/bell/event"
)
const maxRequestBytes = 64 * 1024
type Handler struct{ api.Api }
func IsAdministrator(c *gin.Context) bool {
claims := jwt.ExtractClaims(c)
role, _ := claims[jwt.RoleKey].(string)
return role == "admin"
}
func (h Handler) Create(c *gin.Context) {
if !IsAdministrator(c) {
h.MakeContext(c).Error(http.StatusForbidden, errors.New("无权使用合成事件入口"), "无权使用合成事件入口")
return
}
if err := restoreRequestBody(c); err != nil {
h.MakeContext(c).Error(http.StatusBadRequest, errors.New("请求内容格式不正确"), "请求内容格式不正确")
return
}
var request CreateRequest
h.MakeContext(c).MakeOrm().Bind(&request, binding.JSON)
if h.Errors != nil {
h.Error(http.StatusBadRequest, errors.New("请求内容格式不正确"), "请求内容格式不正确")
return
}
result, err := event.NewService(h.Orm).Ingest(c.Request.Context(), request.Command(), user.GetUserId(c))
if err != nil {
switch {
case errors.Is(err, event.ErrInvalid):
h.Error(http.StatusBadRequest, event.ErrInvalid, event.ErrInvalid.Error())
case errors.Is(err, event.ErrConflict):
h.Error(http.StatusConflict, event.ErrConflict, event.ErrConflict.Error())
default:
h.Logger.Errorf("synthetic event ingest failed: %v", err)
h.Error(http.StatusInternalServerError, errors.New("合成事件写入失败"), "合成事件写入失败")
}
return
}
h.OK(result, "合成事件已接收")
// The client already received the full response. Keep only a fixed marker
// for GoAdmin's generic operation audit, which runs after the handler.
c.Set("result", gin.H{"code": http.StatusOK, "data": "<redacted>"})
}
@@ -0,0 +1,53 @@
package synthetic
import (
"bytes"
"errors"
"io"
"net/http"
"github.com/gin-gonic/gin"
)
const (
requestBodyKey = "bell.synthetic.request-body"
requestBodyErrorKey = "bell.synthetic.request-body-error"
)
var redactedRequestBody = []byte(`{"redacted":true}`)
// RedactRequestBody runs before GoAdmin's operation logger. It keeps the real
// body only in the request context for the handler and exposes a fixed marker
// to the generic audit middleware so event payloads never enter sys_opera_log.
func RedactRequestBody() gin.HandlerFunc {
return func(c *gin.Context) {
if c.Request.Method != http.MethodPost || c.Request.URL.Path != "/api/v1/bell/synthetic-events" {
c.Next()
return
}
body, err := io.ReadAll(io.LimitReader(c.Request.Body, maxRequestBytes+1))
if err != nil {
c.Set(requestBodyErrorKey, err)
} else if len(body) > maxRequestBytes {
c.Set(requestBodyErrorKey, errors.New("request body too large"))
} else {
c.Set(requestBodyKey, body)
}
_ = c.Request.Body.Close()
c.Request.Body = io.NopCloser(bytes.NewReader(redactedRequestBody))
c.Next()
}
}
func restoreRequestBody(c *gin.Context) error {
if value, ok := c.Get(requestBodyErrorKey); ok {
return value.(error)
}
value, ok := c.Get(requestBodyKey)
if !ok {
return errors.New("synthetic request body was not captured")
}
body := value.([]byte)
c.Request.Body = io.NopCloser(bytes.NewReader(body))
return nil
}
+27
View File
@@ -0,0 +1,27 @@
package synthetic
import (
"time"
"go-admin/app/bell/event"
)
const ProducerID = "bell.synthetic"
type CreateRequest struct {
SourceEventID string `json:"sourceEventId" binding:"required"`
EventType string `json:"eventType" binding:"required"`
OccurredAt time.Time `json:"occurredAt" binding:"required"`
Location string `json:"location" binding:"required"`
Severity string `json:"severity" binding:"required"`
EvidenceRef *string `json:"evidenceRef,omitempty"`
Attributes map[string]any `json:"attributes"`
}
func (r CreateRequest) Command() event.Command {
return event.Command{
ProducerID: ProducerID, SourceEventID: r.SourceEventID, EventType: r.EventType,
OccurredAt: r.OccurredAt, Location: r.Location, Severity: r.Severity,
EvidenceRef: r.EvidenceRef, Attributes: r.Attributes,
}
}
+5 -1
View File
@@ -20,6 +20,8 @@ import (
"go-admin/app/admin/models"
"go-admin/app/admin/router"
bellrouter "go-admin/app/bell/router"
"go-admin/app/bell/synthetic"
"go-admin/common/bellconfig"
"go-admin/common/database"
"go-admin/common/global"
@@ -54,6 +56,7 @@ func init() {
//注册路由 fixme 其他应用的路由,在本目录新建文件放在init方法
AppRouters = append(AppRouters, router.InitRouter)
AppRouters = append(AppRouters, bellrouter.InitRouter)
}
func setup() error {
@@ -178,7 +181,8 @@ func initRouter() {
//r.Use(middleware.Metrics())
r.Use(common.Sentinel()).
Use(common.RequestId(pkg.TrafficKey)).
Use(api.SetRequestLogger)
Use(api.SetRequestLogger).
Use(synthetic.RedactRequestBody())
common.InitMiddleware(r)
@@ -0,0 +1,39 @@
package version_local
import (
"runtime"
"gorm.io/gorm"
"go-admin/app/bell/event"
"go-admin/app/bell/receipt"
"go-admin/cmd/migrate/migration"
common "go-admin/common/models"
)
func init() {
_, fileName, _, _ := runtime.Caller(0)
migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateBellEventReceipt)
}
func migrateBellEventReceipt(db *gorm.DB, version string) error {
return db.Transaction(func(tx *gorm.DB) error {
if err := tx.AutoMigrate(new(event.Event), new(receipt.Receipt), new(receipt.IngestAudit)); err != nil {
return err
}
statements := []string{
`ALTER TABLE bell_event_receipts ADD CONSTRAINT bell_event_receipts_event_fk FOREIGN KEY (event_id) REFERENCES bell_events(id) ON UPDATE RESTRICT ON DELETE RESTRICT`,
`ALTER TABLE bell_events ADD CONSTRAINT bell_events_severity_check CHECK (severity IN ('low','medium','high','critical'))`,
`CREATE OR REPLACE FUNCTION bell_reject_immutable_fact() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'Bell immutable fact cannot be changed' USING ERRCODE = '55000'; END $$`,
`CREATE TRIGGER bell_events_immutable BEFORE UPDATE OR DELETE ON bell_events FOR EACH ROW EXECUTE FUNCTION bell_reject_immutable_fact()`,
`CREATE TRIGGER bell_event_receipts_immutable BEFORE UPDATE OR DELETE ON bell_event_receipts FOR EACH ROW EXECUTE FUNCTION bell_reject_immutable_fact()`,
`CREATE TRIGGER bell_event_ingest_audits_immutable BEFORE UPDATE OR DELETE ON bell_event_ingest_audits FOR EACH ROW EXECUTE FUNCTION bell_reject_immutable_fact()`,
}
for _, statement := range statements {
if err := tx.Exec(statement).Error; err != nil {
return err
}
}
return tx.Create(&common.Migration{Version: version}).Error
})
}
@@ -0,0 +1,61 @@
package bell_event_test
import (
"errors"
"testing"
"time"
"go-admin/app/bell/event"
)
func validCommand() event.Command {
return event.Command{
ProducerID: " bell.synthetic ", SourceEventID: " test-001 ", EventType: " danger_area_entered ",
OccurredAt: time.Date(2026, 8, 29, 0, 0, 0, 0, time.FixedZone("CST", 8*60*60)),
Location: " 东门 ", Severity: " HIGH ", Attributes: map[string]any{"z": 2, "a": "first"},
}
}
func TestNormalizeIsDeterministicAndTrimsBusinessFields(t *testing.T) {
first, err := event.Normalize(validCommand())
if err != nil {
t.Fatal(err)
}
secondCommand := validCommand()
secondCommand.Attributes = map[string]any{"a": "first", "z": 2}
second, err := event.Normalize(secondCommand)
if err != nil {
t.Fatal(err)
}
if first.Digest != second.Digest || string(first.Payload) != string(second.Payload) {
t.Fatalf("normalized payload is not deterministic: %s != %s", first.Payload, second.Payload)
}
if first.Command.ProducerID != "bell.synthetic" || first.Command.Severity != "high" {
t.Fatalf("fields were not normalized: %#v", first.Command)
}
if first.Command.OccurredAt.Location() != time.UTC {
t.Fatalf("occurredAt was not converted to UTC: %v", first.Command.OccurredAt)
}
}
func TestNormalizeRejectsUnsafeOrInvalidInput(t *testing.T) {
tests := []struct {
name string
mutate func(*event.Command)
}{
{name: "missing source id", mutate: func(c *event.Command) { c.SourceEventID = " " }},
{name: "unknown severity", mutate: func(c *event.Command) { c.Severity = "urgent" }},
{name: "newline in location", mutate: func(c *event.Command) { c.Location = "东门\nsecret" }},
{name: "unsafe evidence path", mutate: func(c *event.Command) { value := `C:\\secret.jpg`; c.EvidenceRef = &value }},
{name: "oversized attributes", mutate: func(c *event.Command) { c.Attributes = map[string]any{"blob": string(make([]byte, 49*1024))} }},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
command := validCommand()
test.mutate(&command)
if _, err := event.Normalize(command); !errors.Is(err, event.ErrInvalid) {
t.Fatalf("expected ErrInvalid, got %v", err)
}
})
}
}
@@ -0,0 +1,42 @@
package bell_event_test
import (
"testing"
"github.com/gin-gonic/gin"
jwt "github.com/go-admin-team/go-admin-core/sdk/pkg/jwtauth"
"go-admin/app/bell/synthetic"
)
func TestSyntheticEnabledRequiresExplicitNonProductionOptIn(t *testing.T) {
getenv := func(name string) string {
if name == synthetic.EnabledEnv {
return "true"
}
return ""
}
if synthetic.Enabled("prod", getenv) || synthetic.Enabled("production", getenv) {
t.Fatal("synthetic route must stay disabled in production even when the flag is set")
}
if !synthetic.Enabled("dev", getenv) {
t.Fatal("explicit development opt-in should enable the synthetic route")
}
if synthetic.Enabled("dev", func(string) string { return "" }) {
t.Fatal("synthetic route must default to disabled")
}
}
func TestSyntheticAdministratorCheckUsesJWTClaims(t *testing.T) {
gin.SetMode(gin.TestMode)
admin, _ := gin.CreateTestContext(nil)
admin.Set(jwt.JwtPayloadKey, jwt.MapClaims{jwt.RoleKey: "admin"})
if !synthetic.IsAdministrator(admin) {
t.Fatal("admin role should be authorized")
}
operator, _ := gin.CreateTestContext(nil)
operator.Set(jwt.JwtPayloadKey, jwt.MapClaims{jwt.RoleKey: "operator"})
if synthetic.IsAdministrator(operator) {
t.Fatal("non-admin role must be rejected")
}
}