feat: deliver Sense audits to Bell (T-016)
This commit is contained in:
+8
-3
@@ -1,6 +1,6 @@
|
||||
# Sense M1/M2 接入骨架
|
||||
|
||||
本目录是 YoVision Sense 的 M1/M2 接入骨架。数据库保存期望态,ONVIF 和 MediaMTX 通过端口隔离;M1 默认使用 SQLite,T-009~T-012 增加 PostgreSQL 双 schema、Area 准入、本地审计 Outbox、Control API v1、多实例调和 fencing 和孤儿受控处置。默认关闭真实 ONVIF 与公共业务路由;T-006 的真实样机结论仅覆盖已批准的精确海康基线,不能据此宣称多品牌兼容。
|
||||
本目录是 YoVision Sense 的 M1/M2 接入骨架。数据库保存期望态,ONVIF 和 MediaMTX 通过端口隔离;M1 默认使用 SQLite,T-009~T-016 增加 PostgreSQL 双 schema、Area 准入、本地审计 Outbox、Control API v1、多实例调和 fencing、孤儿受控处置和到 Bell 的审计 relay。默认关闭真实 ONVIF、公共业务路由和 relay;T-006 的真实样机结论仅覆盖已批准的精确海康基线,不能据此宣称多品牌兼容。
|
||||
|
||||
## 常用命令
|
||||
|
||||
@@ -42,6 +42,11 @@ Unix 将构建产物改为 `bin/sense-api`。服务默认监听 `127.0.0.1:8080`
|
||||
| `SENSE_CONTROL_AUTH_FILE` | 空 | 仓库外绝对路径;version 1 JSON 只保存 token SHA-256、主体、tenant、Site scope 和权限 |
|
||||
| `SENSE_CONTROL_CURSOR_KEY_FILE` | 空 | 仓库外绝对路径;内容为至少 32 字节随机值的无填充 base64url |
|
||||
| `SENSE_CONTROL_ALLOW_INSECURE_HTTP` | `false` | Control API 非回环明文监听的独立风险接受;正常部署应保持回环并在受控代理终止 TLS |
|
||||
| `SENSE_AUDIT_RELAY_ENABLED` | `false` | 显式开启 PostgreSQL Outbox → Bell relay;SQLite 不支持 |
|
||||
| `SENSE_AUDIT_RELAY_URL` | 空 | 精确指向 Bell `/internal/v1/audit-events:batch`;非回环必须 HTTPS |
|
||||
| `SENSE_AUDIT_RELAY_KEY_FILE` | 空 | 仓库外绝对路径 version 1 JSON key 文件,secret 至少 32 字节 |
|
||||
| `SENSE_AUDIT_RELAY_KEY_ID` | 空 | 本实例用于签名的 key ID |
|
||||
| `SENSE_AUDIT_RELAY_INTERVAL` | `1s` | 队列轮询间隔,最短 1 秒 |
|
||||
|
||||
设备台账只保存 `env://<key>` 凭据引用。真实适配器从进程环境读取以下变量,不把秘密写入 SQLite、日志或 MediaMTX 错误:
|
||||
|
||||
@@ -60,7 +65,7 @@ MediaMTX `v1.19.3` 应作为独立二进制启动并只在可信网络开放 API
|
||||
|
||||
初始化与增量 SQL 位于 `deploy/postgres/`,由高权限部署步骤按文件名前缀执行;Sense 进程不会自动创建角色、schema 或 Bell 对象。`bell_app` 拥有 Site/Area、配额、`capture_policy` 及两个版本化视图,`sense_app` 只能读取两个视图,不能读取或写入 Bell 源表。T-011 的 v4 schema 增加资源版本、24 小时幂等收据和 batch operation;T-012 的 v5 schema 增加数据库时钟租约、Path 历史归属及脱敏孤儿报告/处置结果。表中不保存 MediaMTX source URI。应用登录角色和密码由部署环境创建,不进入仓库。
|
||||
|
||||
PostgreSQL 新建设备必须携带匹配 tenant/Site 的 `area_id`。具有 `video_capture` 能力的设备在创建、移动 Area 和从 disabled 切到 enabled 时执行 Area 准入;`non_imaging_only` 拒绝成像设备但允许非成像设备。投影缺失、非法或版本回退只拒绝新变更,不关闭已有流。创建、配置修改和期望态受理都与对应脱敏 Outbox 事实同事务;停用后调和器只删除该设备的精确 MediaMTX path 并收敛为 offline,不枚举未知 path。Bell relay 尚未实现。
|
||||
PostgreSQL 新建设备必须携带匹配 tenant/Site 的 `area_id`。具有 `video_capture` 能力的设备在创建、移动 Area 和从 disabled 切到 enabled 时执行 Area 准入;`non_imaging_only` 拒绝成像设备但允许非成像设备。投影缺失、非法或版本回退只拒绝新变更,不关闭已有流。创建、配置修改和期望态受理都与对应脱敏 Outbox 事实同事务;停用后调和器只删除该设备的精确 MediaMTX path 并收敛为 offline,不枚举未知 path。开启 relay 后,Sense 用数据库 lease/fencing 批量投递,成功、重试和 dead letter 均保留本地事实;Sense 不访问 Bell schema。
|
||||
|
||||
Windows 本机集成测试从仓库根目录执行:
|
||||
|
||||
@@ -78,7 +83,7 @@ $env:SENSE_DB_DSN = '由部署环境私下设置'
|
||||
go run ./cmd/sense-api
|
||||
```
|
||||
|
||||
PostgreSQL 启动会检查 Sense v5 migration、当前角色对两个 Bell 投影视图和本地控制/对账表的最小权限;权限过宽、视图不可读或 schema 未安装时初始化失败。默认 SQLite 路径和 `cmd/sense-lab` 保持不变,但 SQLite 不实现生产 Area/Outbox、多实例租约或孤儿处置语义,Control API 与孤儿扫描在 SQLite 下不会启动。
|
||||
PostgreSQL 基础启动会检查 Sense v5 migration、当前角色对两个 Bell 投影视图和本地控制/对账表的最小权限;启用 relay 时额外要求 v6 和本地 Outbox 权限,并确认当前 Sense 登录不能访问 Bell 全局审计表。默认 SQLite 路径和 `cmd/sense-lab` 保持不变,但 SQLite 不实现生产 Area/Outbox、多实例租约、孤儿处置或 relay 语义,Control API、孤儿扫描与 relay 在 SQLite 下不会启动。
|
||||
|
||||
### 孤儿报告与受控处置
|
||||
|
||||
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"yovision/sense/internal/auditrelay"
|
||||
"yovision/sense/internal/auth"
|
||||
"yovision/sense/internal/config"
|
||||
"yovision/sense/internal/controlapi"
|
||||
@@ -55,6 +56,28 @@ func run(logger *slog.Logger) error {
|
||||
return err
|
||||
}
|
||||
defer repository.Close()
|
||||
var auditWorker *auditrelay.Worker
|
||||
if cfg.AuditRelayEnabled {
|
||||
relayStore, ok := repository.(auditrelay.Repository)
|
||||
if !ok {
|
||||
return errors.New("selected repository does not support audit relay")
|
||||
}
|
||||
if err := relayStore.AuditRelayReady(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
secret, err := auditrelay.LoadKey(cfg.AuditRelayKeyFile, cfg.AuditRelayKeyID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
client, err := auditrelay.NewClient(cfg.AuditRelayURL, cfg.AuditRelayKeyID, secret, nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
auditWorker, err = auditrelay.NewWorker(relayStore, client, instanceID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
var controlHandler http.Handler
|
||||
if cfg.ControlAPIEnabled {
|
||||
controlStore, ok := repository.(store.ControlRepository)
|
||||
@@ -116,6 +139,9 @@ func run(logger *slog.Logger) error {
|
||||
if orphanScanner != nil {
|
||||
startBackground(func() { orphanScanner.Run(ctx, cfg.OrphanScanInterval, report) })
|
||||
}
|
||||
if auditWorker != nil {
|
||||
startBackground(func() { auditWorker.Run(ctx, cfg.AuditRelayInterval, report) })
|
||||
}
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("GET /healthz", func(writer http.ResponseWriter, _ *http.Request) {
|
||||
@@ -146,7 +172,8 @@ func run(logger *slog.Logger) error {
|
||||
go func() {
|
||||
logger.Info("Sense listening", "address", cfg.HTTPAddress, "version", version,
|
||||
"instance_id", instanceID,
|
||||
"control_api_enabled", cfg.ControlAPIEnabled)
|
||||
"control_api_enabled", cfg.ControlAPIEnabled,
|
||||
"audit_relay_enabled", cfg.AuditRelayEnabled)
|
||||
serverErrors <- server.ListenAndServe()
|
||||
}()
|
||||
|
||||
|
||||
@@ -0,0 +1,132 @@
|
||||
package auditrelay
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func testEnvelope(id string) Envelope {
|
||||
return Envelope{SchemaVersion: 1, Event: Event{
|
||||
EventID: id, EventType: "device.created", TenantID: "tenant", SiteID: "site", DeviceID: "camera-1",
|
||||
Actor: Actor{Type: "system", ID: "sense"}, AggregateGeneration: 1,
|
||||
ProjectionVersions: ProjectionVersions{}, Data: json.RawMessage(`{"kind":"device_created"}`), OccurredAt: time.Unix(1, 0).UTC(),
|
||||
}}
|
||||
}
|
||||
|
||||
func TestClientSignsCanonicalRequest(t *testing.T) {
|
||||
secret := bytesOf(32, 7)
|
||||
server := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, request *http.Request) {
|
||||
body := make([]byte, request.ContentLength)
|
||||
_, _ = request.Body.Read(body)
|
||||
canonical := CanonicalString(request.Method, request.URL.Path, request.Header.Get(HeaderTimestamp), request.Header.Get(HeaderNonce), body)
|
||||
if request.Header.Get(HeaderSignature) != Signature(secret, canonical) {
|
||||
t.Error("request signature did not match canonical vector")
|
||||
}
|
||||
_ = json.NewEncoder(writer).Encode(BatchResponse{Results: []Result{{EventID: "audit_00000000000000000000000000000001", Status: "accepted"}}})
|
||||
}))
|
||||
defer server.Close()
|
||||
client, err := NewClient(server.URL+RelayPath, "sense-a", secret, server.Client())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
client.now = func() time.Time { return time.Unix(1_800_000_000, 0) }
|
||||
client.nonce = func() (string, error) { return "AAAAAAAAAAAAAAAAAAAAAA", nil }
|
||||
results, err := client.Send(context.Background(), []Envelope{testEnvelope("audit_00000000000000000000000000000001")})
|
||||
if err != nil || len(results) != 1 || results[0].Status != "accepted" {
|
||||
t.Fatalf("unexpected relay result: %+v, %v", results, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEndpointAndKeySecurity(t *testing.T) {
|
||||
if _, err := ValidateEndpoint("http://example.com" + RelayPath); err == nil {
|
||||
t.Fatal("remote plaintext relay URL was accepted")
|
||||
}
|
||||
if _, err := ValidateEndpoint("https://example.com" + RelayPath + "?secret=x"); err == nil {
|
||||
t.Fatal("relay URL query was accepted")
|
||||
}
|
||||
secret := bytesOf(32, 9)
|
||||
path := filepath.Join(t.TempDir(), "keys.json")
|
||||
document := map[string]any{"version": 1, "keys": []any{map[string]any{"key_id": "sense-a", "secret_base64url": base64.RawURLEncoding.EncodeToString(secret)}}}
|
||||
raw, _ := json.Marshal(document)
|
||||
if err := os.WriteFile(path, raw, 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
loaded, err := LoadKey(path, "sense-a")
|
||||
if err != nil || sha256.Sum256(loaded) != sha256.Sum256(secret) {
|
||||
t.Fatalf("external key was not loaded: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
type fakeRepository struct {
|
||||
queued []QueuedEvent
|
||||
completed []Completion
|
||||
}
|
||||
|
||||
func (*fakeRepository) AuditRelayReady(context.Context) error { return nil }
|
||||
func (f *fakeRepository) ClaimAuditRelayBatch(context.Context, string, int, time.Duration) ([]QueuedEvent, error) {
|
||||
return f.queued, nil
|
||||
}
|
||||
func (f *fakeRepository) CompleteAuditRelayBatch(_ context.Context, _ string, values []Completion) error {
|
||||
f.completed = append([]Completion(nil), values...)
|
||||
return nil
|
||||
}
|
||||
|
||||
type fakeSender struct {
|
||||
results []Result
|
||||
err error
|
||||
}
|
||||
|
||||
func (f fakeSender) Send(context.Context, []Envelope) ([]Result, error) { return f.results, f.err }
|
||||
|
||||
func TestWorkerDispositionAndBackoff(t *testing.T) {
|
||||
first := "audit_00000000000000000000000000000001"
|
||||
second := "audit_00000000000000000000000000000002"
|
||||
repository := &fakeRepository{queued: []QueuedEvent{{Envelope: testEnvelope(first), LeaseToken: 1, AttemptCount: 1}, {Envelope: testEnvelope(second), LeaseToken: 2, AttemptCount: 10}}}
|
||||
code := "schema_invalid"
|
||||
worker, _ := NewWorker(repository, fakeSender{results: []Result{{EventID: first, Status: "accepted"}, {EventID: second, Status: "rejected", ErrorCode: &code}}}, "worker")
|
||||
if err := worker.RelayOnce(context.Background()); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if repository.completed[0].Disposition != Delivered || repository.completed[1].Disposition != DeadLetter {
|
||||
t.Fatalf("unexpected dispositions: %+v", repository.completed)
|
||||
}
|
||||
|
||||
repository.completed = nil
|
||||
worker.sender = fakeSender{err: errors.New("network")}
|
||||
if err := worker.RelayOnce(context.Background()); err == nil {
|
||||
t.Fatal("network failure was hidden")
|
||||
}
|
||||
if repository.completed[0].RetryAfter != time.Second || repository.completed[1].RetryAfter != 300*time.Second {
|
||||
t.Fatalf("retry bounds drifted: %+v", repository.completed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWorkerRetriesWholeBatchForIncompleteResponse(t *testing.T) {
|
||||
first := "audit_00000000000000000000000000000001"
|
||||
second := "audit_00000000000000000000000000000002"
|
||||
repository := &fakeRepository{queued: []QueuedEvent{{Envelope: testEnvelope(first), LeaseToken: 1, AttemptCount: 1}, {Envelope: testEnvelope(second), LeaseToken: 2, AttemptCount: 2}}}
|
||||
worker, _ := NewWorker(repository, fakeSender{results: []Result{{EventID: first, Status: "accepted"}}}, "worker")
|
||||
if err := worker.RelayOnce(context.Background()); err == nil {
|
||||
t.Fatal("incomplete response was accepted")
|
||||
}
|
||||
if len(repository.completed) != 2 || repository.completed[0].Disposition != Retry || repository.completed[1].Disposition != Retry {
|
||||
t.Fatalf("incomplete response partially completed the batch: %+v", repository.completed)
|
||||
}
|
||||
}
|
||||
|
||||
func bytesOf(size int, value byte) []byte {
|
||||
result := make([]byte, size)
|
||||
for index := range result {
|
||||
result[index] = value
|
||||
}
|
||||
return result
|
||||
}
|
||||
@@ -0,0 +1,131 @@
|
||||
package auditrelay
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/hmac"
|
||||
"crypto/rand"
|
||||
"crypto/sha256"
|
||||
"encoding/base64"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"net/http"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
HeaderKeyID = "X-YoVision-Key-Id"
|
||||
HeaderTimestamp = "X-YoVision-Timestamp"
|
||||
HeaderNonce = "X-YoVision-Nonce"
|
||||
HeaderSignature = "X-YoVision-Signature"
|
||||
RelayPath = "/internal/v1/audit-events:batch"
|
||||
)
|
||||
|
||||
type Client struct {
|
||||
endpoint *url.URL
|
||||
keyID string
|
||||
secret []byte
|
||||
httpClient *http.Client
|
||||
now func() time.Time
|
||||
nonce func() (string, error)
|
||||
}
|
||||
|
||||
func NewClient(rawURL, keyID string, secret []byte, client *http.Client) (*Client, error) {
|
||||
endpoint, err := ValidateEndpoint(rawURL)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if keyID == "" || len(secret) < 32 {
|
||||
return nil, errors.New("audit relay key ID and 32-byte secret are required")
|
||||
}
|
||||
if client == nil {
|
||||
client = &http.Client{Timeout: 10 * time.Second}
|
||||
}
|
||||
return &Client{endpoint: endpoint, keyID: keyID, secret: append([]byte(nil), secret...), httpClient: client, now: time.Now, nonce: randomNonce}, nil
|
||||
}
|
||||
|
||||
func ValidateEndpoint(rawURL string) (*url.URL, error) {
|
||||
parsed, err := url.Parse(rawURL)
|
||||
if err != nil || parsed.Host == "" || parsed.Path != RelayPath || parsed.RawQuery != "" || parsed.Fragment != "" || parsed.User != nil {
|
||||
return nil, errors.New("invalid Bell audit relay URL")
|
||||
}
|
||||
host := parsed.Hostname()
|
||||
ip := net.ParseIP(host)
|
||||
loopback := strings.EqualFold(host, "localhost") || (ip != nil && ip.IsLoopback())
|
||||
if parsed.Scheme != "https" && !(parsed.Scheme == "http" && loopback) {
|
||||
return nil, errors.New("Bell audit relay URL requires HTTPS outside loopback")
|
||||
}
|
||||
return parsed, nil
|
||||
}
|
||||
|
||||
func CanonicalString(method, path, timestamp, nonce string, body []byte) string {
|
||||
digest := sha256.Sum256(body)
|
||||
return strings.Join([]string{method, path, timestamp, nonce, hex.EncodeToString(digest[:])}, "\n")
|
||||
}
|
||||
|
||||
func Signature(secret []byte, canonical string) string {
|
||||
mac := hmac.New(sha256.New, secret)
|
||||
_, _ = mac.Write([]byte(canonical))
|
||||
return base64.RawURLEncoding.EncodeToString(mac.Sum(nil))
|
||||
}
|
||||
|
||||
func randomNonce() (string, error) {
|
||||
value := make([]byte, 16)
|
||||
if _, err := rand.Read(value); err != nil {
|
||||
return "", err
|
||||
}
|
||||
return base64.RawURLEncoding.EncodeToString(value), nil
|
||||
}
|
||||
|
||||
func (c *Client) Send(ctx context.Context, events []Envelope) ([]Result, error) {
|
||||
if len(events) < 1 || len(events) > MaxBatchSize {
|
||||
return nil, errors.New("audit relay batch must contain 1 to 100 events")
|
||||
}
|
||||
body, err := json.Marshal(BatchRequest{Events: events})
|
||||
if err != nil || len(body) > MaxBodyBytes {
|
||||
return nil, errors.New("encode audit relay batch")
|
||||
}
|
||||
requestContext, cancel := context.WithTimeout(ctx, 10*time.Second)
|
||||
defer cancel()
|
||||
timestamp := strconv.FormatInt(c.now().UTC().Unix(), 10)
|
||||
nonce, err := c.nonce()
|
||||
if err != nil {
|
||||
return nil, errors.New("generate audit relay nonce")
|
||||
}
|
||||
request, err := http.NewRequestWithContext(requestContext, http.MethodPost, c.endpoint.String(), bytes.NewReader(body))
|
||||
if err != nil {
|
||||
return nil, errors.New("create audit relay request")
|
||||
}
|
||||
request.Header.Set("Content-Type", "application/json")
|
||||
request.Header.Set(HeaderKeyID, c.keyID)
|
||||
request.Header.Set(HeaderTimestamp, timestamp)
|
||||
request.Header.Set(HeaderNonce, nonce)
|
||||
request.Header.Set(HeaderSignature, Signature(c.secret, CanonicalString(http.MethodPost, RelayPath, timestamp, nonce, body)))
|
||||
response, err := c.httpClient.Do(request)
|
||||
if err != nil {
|
||||
return nil, errors.New("send audit relay request")
|
||||
}
|
||||
defer response.Body.Close()
|
||||
if response.StatusCode != http.StatusOK {
|
||||
_, _ = io.Copy(io.Discard, io.LimitReader(response.Body, 4096))
|
||||
return nil, fmt.Errorf("Bell audit relay returned HTTP %d", response.StatusCode)
|
||||
}
|
||||
decoder := json.NewDecoder(io.LimitReader(response.Body, MaxBodyBytes+1))
|
||||
decoder.DisallowUnknownFields()
|
||||
var decoded BatchResponse
|
||||
if err := decoder.Decode(&decoded); err != nil {
|
||||
return nil, errors.New("decode audit relay response")
|
||||
}
|
||||
var trailing any
|
||||
if err := decoder.Decode(&trailing); !errors.Is(err, io.EOF) {
|
||||
return nil, errors.New("audit relay response contains trailing data")
|
||||
}
|
||||
return decoded.Results, nil
|
||||
}
|
||||
@@ -0,0 +1,53 @@
|
||||
package auditrelay
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/base64"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"io"
|
||||
"os"
|
||||
"regexp"
|
||||
)
|
||||
|
||||
var keyIDPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$`)
|
||||
|
||||
type keyFile struct {
|
||||
Version int `json:"version"`
|
||||
Keys []struct {
|
||||
KeyID string `json:"key_id"`
|
||||
Secret string `json:"secret_base64url"`
|
||||
} `json:"keys"`
|
||||
}
|
||||
|
||||
func LoadKey(path, keyID string) ([]byte, error) {
|
||||
raw, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return nil, errors.New("read audit relay key file")
|
||||
}
|
||||
var document keyFile
|
||||
decoder := json.NewDecoder(bytes.NewReader(raw))
|
||||
decoder.DisallowUnknownFields()
|
||||
if err := decoder.Decode(&document); err != nil || document.Version != 1 || len(document.Keys) == 0 {
|
||||
return nil, errors.New("invalid audit relay key file")
|
||||
}
|
||||
var trailing any
|
||||
if err := decoder.Decode(&trailing); !errors.Is(err, io.EOF) {
|
||||
return nil, errors.New("invalid audit relay key file")
|
||||
}
|
||||
values := make(map[string][]byte, len(document.Keys))
|
||||
for _, candidate := range document.Keys {
|
||||
secret, err := base64.RawURLEncoding.DecodeString(candidate.Secret)
|
||||
if err != nil || !keyIDPattern.MatchString(candidate.KeyID) || len(secret) < 32 {
|
||||
return nil, errors.New("invalid audit relay secret")
|
||||
}
|
||||
if _, duplicate := values[candidate.KeyID]; duplicate {
|
||||
return nil, errors.New("duplicate audit relay key ID")
|
||||
}
|
||||
values[candidate.KeyID] = secret
|
||||
}
|
||||
if secret, exists := values[keyID]; exists {
|
||||
return secret, nil
|
||||
}
|
||||
return nil, errors.New("audit relay key ID not found")
|
||||
}
|
||||
@@ -0,0 +1,89 @@
|
||||
// Package auditrelay delivers Sense-owned audit facts to Bell without sharing databases.
|
||||
package auditrelay
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
MaxBatchSize = 100
|
||||
MaxBodyBytes = 1 << 20
|
||||
LeaseDuration = 30 * time.Second
|
||||
)
|
||||
|
||||
type Actor struct {
|
||||
Type string `json:"type"`
|
||||
ID string `json:"id"`
|
||||
}
|
||||
|
||||
type ProjectionVersions struct {
|
||||
QuotaSourceVersion *int64 `json:"quota_source_version"`
|
||||
AreaPolicySourceVersion *int64 `json:"area_policy_source_version"`
|
||||
}
|
||||
|
||||
type Event struct {
|
||||
EventID string `json:"event_id"`
|
||||
EventType string `json:"event_type"`
|
||||
TenantID string `json:"tenant_id"`
|
||||
SiteID string `json:"site_id"`
|
||||
DeviceID string `json:"device_id"`
|
||||
Actor Actor `json:"actor"`
|
||||
Reason *string `json:"reason"`
|
||||
TraceID *string `json:"trace_id"`
|
||||
AggregateGeneration int64 `json:"aggregate_generation"`
|
||||
ProjectionVersions ProjectionVersions `json:"projection_versions"`
|
||||
Data json.RawMessage `json:"data"`
|
||||
OccurredAt time.Time `json:"occurred_at"`
|
||||
}
|
||||
|
||||
type Envelope struct {
|
||||
SchemaVersion int `json:"schema_version"`
|
||||
Event Event `json:"event"`
|
||||
}
|
||||
|
||||
type BatchRequest struct {
|
||||
Events []Envelope `json:"events"`
|
||||
}
|
||||
|
||||
type Result struct {
|
||||
EventID string `json:"event_id"`
|
||||
Status string `json:"status"`
|
||||
ErrorCode *string `json:"error_code,omitempty"`
|
||||
}
|
||||
|
||||
type BatchResponse struct {
|
||||
Results []Result `json:"results"`
|
||||
}
|
||||
|
||||
type QueuedEvent struct {
|
||||
Envelope
|
||||
LeaseToken int64
|
||||
AttemptCount int
|
||||
}
|
||||
|
||||
type Disposition string
|
||||
|
||||
const (
|
||||
Delivered Disposition = "delivered"
|
||||
DeadLetter Disposition = "dead_letter"
|
||||
Retry Disposition = "retry"
|
||||
)
|
||||
|
||||
type Completion struct {
|
||||
EventID string
|
||||
LeaseToken int64
|
||||
Disposition Disposition
|
||||
ErrorCode string
|
||||
RetryAfter time.Duration
|
||||
}
|
||||
|
||||
var ErrLeaseLost = errors.New("audit relay lease lost")
|
||||
|
||||
type Repository interface {
|
||||
AuditRelayReady(context.Context) error
|
||||
ClaimAuditRelayBatch(context.Context, string, int, time.Duration) ([]QueuedEvent, error)
|
||||
CompleteAuditRelayBatch(context.Context, string, []Completion) error
|
||||
}
|
||||
@@ -0,0 +1,126 @@
|
||||
package auditrelay
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"regexp"
|
||||
"time"
|
||||
)
|
||||
|
||||
var stableErrorCode = regexp.MustCompile(`^[a-z][a-z0-9_]{0,63}$`)
|
||||
|
||||
type Sender interface {
|
||||
Send(context.Context, []Envelope) ([]Result, error)
|
||||
}
|
||||
|
||||
type Worker struct {
|
||||
repository Repository
|
||||
sender Sender
|
||||
owner string
|
||||
}
|
||||
|
||||
func NewWorker(repository Repository, sender Sender, owner string) (*Worker, error) {
|
||||
if repository == nil || sender == nil || owner == "" {
|
||||
return nil, errors.New("audit relay worker dependencies are required")
|
||||
}
|
||||
return &Worker{repository: repository, sender: sender, owner: owner}, nil
|
||||
}
|
||||
|
||||
func RetryDelay(attempt int) time.Duration {
|
||||
if attempt < 1 {
|
||||
attempt = 1
|
||||
}
|
||||
if attempt > 9 {
|
||||
return 300 * time.Second
|
||||
}
|
||||
delay := time.Second << (attempt - 1)
|
||||
if delay > 300*time.Second {
|
||||
return 300 * time.Second
|
||||
}
|
||||
return delay
|
||||
}
|
||||
|
||||
func (w *Worker) RelayOnce(ctx context.Context) error {
|
||||
queued, err := w.repository.ClaimAuditRelayBatch(ctx, w.owner, MaxBatchSize, LeaseDuration)
|
||||
if err != nil || len(queued) == 0 {
|
||||
return err
|
||||
}
|
||||
events := make([]Envelope, len(queued))
|
||||
for index := range queued {
|
||||
events[index] = queued[index].Envelope
|
||||
}
|
||||
results, sendErr := w.sender.Send(ctx, events)
|
||||
if sendErr != nil {
|
||||
completions := make([]Completion, len(queued))
|
||||
for index, value := range queued {
|
||||
completions[index] = Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: Retry, ErrorCode: "delivery_failed", RetryAfter: RetryDelay(value.AttemptCount)}
|
||||
}
|
||||
if err := w.repository.CompleteAuditRelayBatch(ctx, w.owner, completions); err != nil {
|
||||
return err
|
||||
}
|
||||
return sendErr
|
||||
}
|
||||
if len(results) != len(queued) {
|
||||
return w.retryAll(ctx, queued, "invalid_response")
|
||||
}
|
||||
expected := make(map[string]bool, len(queued))
|
||||
for _, value := range queued {
|
||||
expected[value.Event.EventID] = true
|
||||
}
|
||||
byID := make(map[string]Result, len(results))
|
||||
for _, result := range results {
|
||||
validStatus := ((result.Status == "accepted" || result.Status == "duplicate") && result.ErrorCode == nil) ||
|
||||
(result.Status == "rejected" && result.ErrorCode != nil && stableErrorCode.MatchString(*result.ErrorCode))
|
||||
if !expected[result.EventID] || !validStatus {
|
||||
return w.retryAll(ctx, queued, "invalid_response")
|
||||
}
|
||||
if _, duplicate := byID[result.EventID]; duplicate {
|
||||
return w.retryAll(ctx, queued, "invalid_response")
|
||||
}
|
||||
byID[result.EventID] = result
|
||||
}
|
||||
completions := make([]Completion, 0, len(queued))
|
||||
for _, value := range queued {
|
||||
result, ok := byID[value.Event.EventID]
|
||||
if !ok {
|
||||
return w.retryAll(ctx, queued, "invalid_response")
|
||||
}
|
||||
switch result.Status {
|
||||
case "accepted", "duplicate":
|
||||
completions = append(completions, Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: Delivered})
|
||||
case "rejected":
|
||||
completions = append(completions, Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: DeadLetter, ErrorCode: *result.ErrorCode})
|
||||
}
|
||||
}
|
||||
return w.repository.CompleteAuditRelayBatch(ctx, w.owner, completions)
|
||||
}
|
||||
|
||||
func (w *Worker) retryAll(ctx context.Context, queued []QueuedEvent, code string) error {
|
||||
values := make([]Completion, len(queued))
|
||||
for index, value := range queued {
|
||||
values[index] = Completion{EventID: value.Event.EventID, LeaseToken: value.LeaseToken, Disposition: Retry, ErrorCode: code, RetryAfter: RetryDelay(value.AttemptCount)}
|
||||
}
|
||||
if err := w.repository.CompleteAuditRelayBatch(ctx, w.owner, values); err != nil {
|
||||
return err
|
||||
}
|
||||
return fmt.Errorf("audit relay %s", code)
|
||||
}
|
||||
|
||||
func (w *Worker) Run(ctx context.Context, interval time.Duration, report func(error)) {
|
||||
if interval <= 0 {
|
||||
interval = time.Second
|
||||
}
|
||||
for {
|
||||
if err := w.RelayOnce(ctx); err != nil && ctx.Err() == nil && report != nil {
|
||||
report(err)
|
||||
}
|
||||
timer := time.NewTimer(interval)
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
timer.Stop()
|
||||
return
|
||||
case <-timer.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -26,6 +26,7 @@ const (
|
||||
defaultOrphanScanPeriod = time.Minute
|
||||
defaultONVIFMode = "disabled"
|
||||
defaultControlAuthMode = "static-sha256"
|
||||
defaultAuditRelayPeriod = time.Second
|
||||
)
|
||||
|
||||
var instanceIDPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$`)
|
||||
@@ -53,6 +54,11 @@ type Config struct {
|
||||
ControlAuthFile string
|
||||
ControlCursorKeyFile string
|
||||
ControlAllowInsecureHTTP bool
|
||||
AuditRelayEnabled bool
|
||||
AuditRelayURL string
|
||||
AuditRelayKeyFile string
|
||||
AuditRelayKeyID string
|
||||
AuditRelayInterval time.Duration
|
||||
}
|
||||
|
||||
func Load() (Config, error) {
|
||||
@@ -92,6 +98,14 @@ func Load() (Config, error) {
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
}
|
||||
auditRelayEnabled, err := boolEnv("SENSE_AUDIT_RELAY_ENABLED", false)
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
}
|
||||
auditRelayInterval, err := durationEnv("SENSE_AUDIT_RELAY_INTERVAL", defaultAuditRelayPeriod)
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
}
|
||||
metricsEnabled, err := boolEnv("SENSE_METRICS_ENABLED", true)
|
||||
if err != nil {
|
||||
return Config{}, err
|
||||
@@ -130,6 +144,11 @@ func Load() (Config, error) {
|
||||
ControlAuthFile: stringEnv("SENSE_CONTROL_AUTH_FILE", ""),
|
||||
ControlCursorKeyFile: stringEnv("SENSE_CONTROL_CURSOR_KEY_FILE", ""),
|
||||
ControlAllowInsecureHTTP: controlAllowInsecure,
|
||||
AuditRelayEnabled: auditRelayEnabled,
|
||||
AuditRelayURL: stringEnv("SENSE_AUDIT_RELAY_URL", ""),
|
||||
AuditRelayKeyFile: stringEnv("SENSE_AUDIT_RELAY_KEY_FILE", ""),
|
||||
AuditRelayKeyID: stringEnv("SENSE_AUDIT_RELAY_KEY_ID", ""),
|
||||
AuditRelayInterval: auditRelayInterval,
|
||||
}
|
||||
if err := cfg.Validate(); err != nil {
|
||||
return Config{}, err
|
||||
@@ -228,6 +247,31 @@ func (c Config) Validate() error {
|
||||
return fmt.Errorf("non-loopback Control API requires SENSE_CONTROL_ALLOW_INSECURE_HTTP=true")
|
||||
}
|
||||
}
|
||||
if c.AuditRelayEnabled {
|
||||
if databaseDriver != postgresDatabaseDriver {
|
||||
return fmt.Errorf("Sense audit relay requires SENSE_DB_DRIVER=postgres")
|
||||
}
|
||||
if c.AuditRelayKeyFile == "" || !filepath.IsAbs(c.AuditRelayKeyFile) {
|
||||
return fmt.Errorf("SENSE_AUDIT_RELAY_KEY_FILE must be an absolute external path")
|
||||
}
|
||||
if !instanceIDPattern.MatchString(c.AuditRelayKeyID) {
|
||||
return fmt.Errorf("invalid SENSE_AUDIT_RELAY_KEY_ID")
|
||||
}
|
||||
if c.AuditRelayInterval < time.Second {
|
||||
return fmt.Errorf("SENSE_AUDIT_RELAY_INTERVAL must be at least 1s")
|
||||
}
|
||||
relayURL, err := url.Parse(c.AuditRelayURL)
|
||||
if err != nil || relayURL.Host == "" || relayURL.Path != "/internal/v1/audit-events:batch" ||
|
||||
relayURL.RawQuery != "" || relayURL.Fragment != "" || relayURL.User != nil {
|
||||
return fmt.Errorf("invalid SENSE_AUDIT_RELAY_URL")
|
||||
}
|
||||
relayHost := relayURL.Hostname()
|
||||
relayIP := net.ParseIP(relayHost)
|
||||
relayLoopback := relayHost == "localhost" || (relayIP != nil && relayIP.IsLoopback())
|
||||
if relayURL.Scheme != "https" && !(relayURL.Scheme == "http" && relayLoopback) {
|
||||
return fmt.Errorf("SENSE_AUDIT_RELAY_URL requires HTTPS outside loopback")
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -166,3 +166,31 @@ func TestValidateControlAPINonLoopbackNeedsSeparateRiskAcceptance(t *testing.T)
|
||||
t.Fatalf("explicit non-loopback Control API risk acceptance failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestValidateAuditRelaySecurityBoundary(t *testing.T) {
|
||||
base := Config{
|
||||
HTTPAddress: "127.0.0.1:8080", DatabaseDriver: "postgres",
|
||||
DatabaseDSN: "postgres://sense-runtime@127.0.0.1/yovision?sslmode=disable",
|
||||
MediaMTXURL: "http://127.0.0.1:9997", ReconcileInterval: time.Second, ProbeInterval: time.Second,
|
||||
AuditRelayEnabled: true, AuditRelayURL: "http://127.0.0.1:8081/internal/v1/audit-events:batch",
|
||||
AuditRelayKeyFile: filepath.Join(t.TempDir(), "relay-keys.json"), AuditRelayKeyID: "sense-a", AuditRelayInterval: time.Second,
|
||||
}
|
||||
if err := base.Validate(); err != nil {
|
||||
t.Fatalf("valid loopback relay rejected: %v", err)
|
||||
}
|
||||
remoteHTTP := base
|
||||
remoteHTTP.AuditRelayURL = "http://bell.example/internal/v1/audit-events:batch"
|
||||
if err := remoteHTTP.Validate(); err == nil {
|
||||
t.Fatal("remote plaintext relay was accepted")
|
||||
}
|
||||
sqlite := base
|
||||
sqlite.DatabaseDriver, sqlite.DatabaseDSN = "sqlite", "file:test.db"
|
||||
if err := sqlite.Validate(); err == nil {
|
||||
t.Fatal("SQLite audit relay was accepted")
|
||||
}
|
||||
relativeKey := base
|
||||
relativeKey.AuditRelayKeyFile = "relay-keys.json"
|
||||
if err := relativeKey.Validate(); err == nil {
|
||||
t.Fatal("repository-relative relay key was accepted")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,176 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"regexp"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"yovision/sense/internal/auditrelay"
|
||||
)
|
||||
|
||||
var relayErrorCode = regexp.MustCompile(`^[a-z][a-z0-9_]{0,63}$`)
|
||||
|
||||
func (s *Postgres) AuditRelayReady(ctx context.Context) error {
|
||||
var version int64
|
||||
if err := s.db.QueryRowContext(ctx, `SELECT COALESCE(MAX(version), 0) FROM sense.schema_migrations`).Scan(&version); err != nil || version < 6 {
|
||||
return errors.New("postgres Sense schema migration v6 is required for audit relay")
|
||||
}
|
||||
var canUseOutbox, canReadBell, canWriteBell bool
|
||||
if err := s.db.QueryRowContext(ctx, `SELECT
|
||||
has_table_privilege(current_user, 'sense.device_operation_outbox', 'SELECT,INSERT,UPDATE,DELETE'),
|
||||
has_table_privilege(current_user, 'bell.audit_events', 'SELECT'),
|
||||
has_table_privilege(current_user, 'bell.audit_events', 'INSERT,UPDATE,DELETE')`).Scan(
|
||||
&canUseOutbox, &canReadBell, &canWriteBell,
|
||||
); err != nil {
|
||||
return errors.New("verify audit relay privileges")
|
||||
}
|
||||
if !canUseOutbox || canReadBell || canWriteBell {
|
||||
return errors.New("Sense audit relay violates schema ownership boundary")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *Postgres) ClaimAuditRelayBatch(
|
||||
ctx context.Context,
|
||||
owner string,
|
||||
limit int,
|
||||
lease time.Duration,
|
||||
) ([]auditrelay.QueuedEvent, error) {
|
||||
if strings.TrimSpace(owner) == "" || limit < 1 || limit > auditrelay.MaxBatchSize || lease <= 0 {
|
||||
return nil, errors.New("invalid audit relay claim")
|
||||
}
|
||||
rows, err := s.db.QueryContext(ctx, `WITH due AS (
|
||||
SELECT event_id FROM sense.device_operation_outbox
|
||||
WHERE delivered_at IS NULL AND dead_lettered_at IS NULL
|
||||
AND COALESCE(next_attempt_at, available_at) <= clock_timestamp()
|
||||
AND (relay_lease_until IS NULL OR relay_lease_until <= clock_timestamp())
|
||||
ORDER BY COALESCE(next_attempt_at, available_at), event_id
|
||||
FOR UPDATE SKIP LOCKED LIMIT $1
|
||||
)
|
||||
UPDATE sense.device_operation_outbox AS outbox SET
|
||||
relay_lease_owner = $2,
|
||||
relay_lease_token = outbox.relay_lease_token + 1,
|
||||
relay_lease_until = clock_timestamp() + ($3 * interval '1 second'),
|
||||
attempt_count = outbox.attempt_count + 1,
|
||||
last_error_code = NULL
|
||||
FROM due WHERE outbox.event_id = due.event_id
|
||||
RETURNING outbox.event_id, outbox.event_type, outbox.tenant_id,
|
||||
outbox.site_id, outbox.device_id, outbox.actor_type, outbox.actor_id,
|
||||
outbox.reason, outbox.trace_id, outbox.aggregate_generation,
|
||||
outbox.quota_source_version, outbox.area_policy_source_version,
|
||||
outbox.payload, outbox.occurred_at, outbox.relay_lease_token,
|
||||
outbox.attempt_count`, limit, owner, lease.Seconds())
|
||||
if err != nil {
|
||||
return nil, errors.New("claim audit relay batch")
|
||||
}
|
||||
defer rows.Close()
|
||||
values := make([]auditrelay.QueuedEvent, 0)
|
||||
for rows.Next() {
|
||||
var value auditrelay.QueuedEvent
|
||||
var reason, trace sql.NullString
|
||||
var quota, area sql.NullInt64
|
||||
var payload []byte
|
||||
if err := rows.Scan(
|
||||
&value.Event.EventID, &value.Event.EventType, &value.Event.TenantID,
|
||||
&value.Event.SiteID, &value.Event.DeviceID, &value.Event.Actor.Type,
|
||||
&value.Event.Actor.ID, &reason, &trace, &value.Event.AggregateGeneration,
|
||||
"a, &area, &payload, &value.Event.OccurredAt, &value.LeaseToken,
|
||||
&value.AttemptCount,
|
||||
); err != nil {
|
||||
return nil, errors.New("scan audit relay claim")
|
||||
}
|
||||
if reason.Valid {
|
||||
value.Event.Reason = &reason.String
|
||||
}
|
||||
if trace.Valid {
|
||||
value.Event.TraceID = &trace.String
|
||||
}
|
||||
if quota.Valid {
|
||||
value.Event.ProjectionVersions.QuotaSourceVersion = "a.Int64
|
||||
}
|
||||
if area.Valid {
|
||||
value.Event.ProjectionVersions.AreaPolicySourceVersion = &area.Int64
|
||||
}
|
||||
if !json.Valid(payload) {
|
||||
return nil, errors.New("invalid audit payload in outbox")
|
||||
}
|
||||
value.Event.Data = append(json.RawMessage(nil), payload...)
|
||||
value.SchemaVersion = 1
|
||||
if value.Event.EventType == "device.configuration.accepted" {
|
||||
value.SchemaVersion = 2
|
||||
}
|
||||
values = append(values, value)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, errors.New("iterate audit relay claims")
|
||||
}
|
||||
return values, nil
|
||||
}
|
||||
|
||||
func (s *Postgres) CompleteAuditRelayBatch(
|
||||
ctx context.Context,
|
||||
owner string,
|
||||
values []auditrelay.Completion,
|
||||
) error {
|
||||
if strings.TrimSpace(owner) == "" || len(values) == 0 || len(values) > auditrelay.MaxBatchSize {
|
||||
return errors.New("invalid audit relay completion")
|
||||
}
|
||||
tx, err := s.db.BeginTx(ctx, nil)
|
||||
if err != nil {
|
||||
return errors.New("begin audit relay completion")
|
||||
}
|
||||
defer tx.Rollback()
|
||||
for _, value := range values {
|
||||
code := value.ErrorCode
|
||||
if code != "" && !relayErrorCode.MatchString(code) {
|
||||
code = "invalid_response"
|
||||
}
|
||||
var result sql.Result
|
||||
switch value.Disposition {
|
||||
case auditrelay.Delivered:
|
||||
result, err = tx.ExecContext(ctx, `UPDATE sense.device_operation_outbox SET
|
||||
delivered_at = clock_timestamp(), last_error_code = NULL,
|
||||
relay_lease_owner = NULL, relay_lease_until = NULL
|
||||
WHERE event_id = $1 AND relay_lease_owner = $2 AND relay_lease_token = $3
|
||||
AND relay_lease_until > clock_timestamp()`, value.EventID, owner, value.LeaseToken)
|
||||
case auditrelay.DeadLetter:
|
||||
if code == "" {
|
||||
return errors.New("dead-letter completion requires an error code")
|
||||
}
|
||||
result, err = tx.ExecContext(ctx, `UPDATE sense.device_operation_outbox SET
|
||||
dead_lettered_at = clock_timestamp(), last_error_code = $4,
|
||||
relay_lease_owner = NULL, relay_lease_until = NULL
|
||||
WHERE event_id = $1 AND relay_lease_owner = $2 AND relay_lease_token = $3
|
||||
AND relay_lease_until > clock_timestamp()`, value.EventID, owner, value.LeaseToken, code)
|
||||
case auditrelay.Retry:
|
||||
if code == "" || value.RetryAfter < time.Second || value.RetryAfter > 300*time.Second {
|
||||
return errors.New("invalid audit relay retry completion")
|
||||
}
|
||||
result, err = tx.ExecContext(ctx, `UPDATE sense.device_operation_outbox SET
|
||||
next_attempt_at = clock_timestamp() + ($4 * interval '1 second'),
|
||||
last_error_code = $5, relay_lease_owner = NULL, relay_lease_until = NULL
|
||||
WHERE event_id = $1 AND relay_lease_owner = $2 AND relay_lease_token = $3
|
||||
AND relay_lease_until > clock_timestamp()`, value.EventID, owner, value.LeaseToken, value.RetryAfter.Seconds(), code)
|
||||
default:
|
||||
return errors.New("invalid audit relay disposition")
|
||||
}
|
||||
if err != nil {
|
||||
return errors.New("persist audit relay completion")
|
||||
}
|
||||
affected, err := result.RowsAffected()
|
||||
if err != nil {
|
||||
return errors.New("read audit relay completion")
|
||||
}
|
||||
if affected != 1 {
|
||||
return auditrelay.ErrLeaseLost
|
||||
}
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return errors.New("commit audit relay completion")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,45 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"yovision/sense/internal/auditrelay"
|
||||
)
|
||||
|
||||
func TestPostgresAuditRelayClaimUsesFencing(t *testing.T) {
|
||||
repository, admin := openPostgresTestStore(t)
|
||||
ctx := context.Background()
|
||||
if err := repository.AuditRelayReady(ctx); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
insertBellSite(t, admin, "relay-tenant", "relay-site", 1)
|
||||
device := videoDevice(9001, "relay-tenant", "relay-site")
|
||||
if err := repository.CreateDevice(ctx, device); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
first, err := repository.ClaimAuditRelayBatch(ctx, "worker-a", 100, 30*time.Second)
|
||||
if err != nil || len(first) != 1 || first[0].SchemaVersion != 1 || first[0].AttemptCount != 1 {
|
||||
t.Fatalf("first claim: %+v %v", first, err)
|
||||
}
|
||||
if _, err := admin.ExecContext(ctx, `UPDATE sense.device_operation_outbox SET relay_lease_until=clock_timestamp()-interval '1 second' WHERE event_id=$1`, first[0].Event.EventID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
second, err := repository.ClaimAuditRelayBatch(ctx, "worker-b", 100, 30*time.Second)
|
||||
if err != nil || len(second) != 1 || second[0].LeaseToken <= first[0].LeaseToken {
|
||||
t.Fatalf("reclaim: %+v %v", second, err)
|
||||
}
|
||||
err = repository.CompleteAuditRelayBatch(ctx, "worker-a", []auditrelay.Completion{{EventID: first[0].Event.EventID, LeaseToken: first[0].LeaseToken, Disposition: auditrelay.Delivered}})
|
||||
if !errors.Is(err, auditrelay.ErrLeaseLost) {
|
||||
t.Fatalf("stale worker completion returned %v", err)
|
||||
}
|
||||
if err := repository.CompleteAuditRelayBatch(ctx, "worker-b", []auditrelay.Completion{{EventID: second[0].Event.EventID, LeaseToken: second[0].LeaseToken, Disposition: auditrelay.Delivered}}); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
var delivered bool
|
||||
if err := admin.QueryRowContext(ctx, `SELECT delivered_at IS NOT NULL FROM sense.device_operation_outbox WHERE event_id=$1`, first[0].Event.EventID).Scan(&delivered); err != nil || !delivered {
|
||||
t.Fatalf("delivered=%v err=%v", delivered, err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user