From 4be2386421401e82db546c435ba6ab4b46588a79 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Tue, 11 Aug 2026 00:24:32 +0800 Subject: [PATCH] feat: deliver Sense audits to Bell (T-016) --- Bell/README.md | 16 +- Bell/cmd/bell-api/main.go | 152 ++++++++ Bell/cmd/bell-api/main_test.go | 33 ++ Bell/internal/audit/audit.go | 363 ++++++++++++++++++ Bell/internal/audit/audit_test.go | 105 +++++ Bell/internal/audit/key.go | 47 +++ Bell/internal/store/audit_postgres.go | 138 +++++++ Bell/internal/store/audit_postgres_test.go | 82 ++++ Sense/README.md | 11 +- Sense/cmd/sense-api/main.go | 29 +- Sense/internal/auditrelay/auditrelay_test.go | 132 +++++++ Sense/internal/auditrelay/client.go | 131 +++++++ Sense/internal/auditrelay/key.go | 53 +++ Sense/internal/auditrelay/model.go | 89 +++++ Sense/internal/auditrelay/worker.go | 126 ++++++ Sense/internal/config/config.go | 44 +++ Sense/internal/config/config_test.go | 28 ++ Sense/internal/store/audit_relay_postgres.go | 176 +++++++++ .../store/audit_relay_postgres_test.go | 45 +++ deploy/postgres/014_audit_relay.sql | 130 +++++++ .../postgres/015_privileges_audit_relay.sql | 11 + deploy/postgres/README.md | 5 +- deploy/postgres/tests/assertions.sql | 31 +- docs/00-ai-start-here.md | 2 +- docs/03-tech-stack.md | 8 +- docs/04-architecture.md | 11 +- docs/06-tasks.md | 1 + docs/api.md | 8 +- docs/contracts/README.md | 23 +- .../sense-audit-relay-v1.openapi.json | 94 +++++ docs/current-state.md | 19 +- docs/tasks/T-016.md | 12 +- scripts/test_postgres.ps1 | 4 +- tests/test_postgres_contract.py | 17 + tests/test_sense_audit_relay_contract.py | 54 +++ 35 files changed, 2188 insertions(+), 42 deletions(-) create mode 100644 Bell/cmd/bell-api/main.go create mode 100644 Bell/cmd/bell-api/main_test.go create mode 100644 Bell/internal/audit/audit.go create mode 100644 Bell/internal/audit/audit_test.go create mode 100644 Bell/internal/audit/key.go create mode 100644 Bell/internal/store/audit_postgres.go create mode 100644 Bell/internal/store/audit_postgres_test.go create mode 100644 Sense/internal/auditrelay/auditrelay_test.go create mode 100644 Sense/internal/auditrelay/client.go create mode 100644 Sense/internal/auditrelay/key.go create mode 100644 Sense/internal/auditrelay/model.go create mode 100644 Sense/internal/auditrelay/worker.go create mode 100644 Sense/internal/store/audit_relay_postgres.go create mode 100644 Sense/internal/store/audit_relay_postgres_test.go create mode 100644 deploy/postgres/014_audit_relay.sql create mode 100644 deploy/postgres/015_privileges_audit_relay.sql create mode 100644 docs/contracts/sense-audit-relay-v1.openapi.json create mode 100644 tests/test_sense_audit_relay_contract.py diff --git a/Bell/README.md b/Bell/README.md index df7df08..00b0808 100644 --- a/Bell/README.md +++ b/Bell/README.md @@ -1,13 +1,21 @@ -# Bell 事件存储基础 +# Bell 事件存储与内部审计入口 -Bell 当前只实现 M3 的事件域基础,不包含可部署 HTTP 服务: +Bell 当前实现 M3 的事件域基础及 Sense 审计 relay 的最小内部 HTTP 服务: - Bell 在可信 ingress 内为不含 `id` 的候选事实生成 `evt_` ULID。 - 最终事件同时通过冻结 v0.1 JSON Schema 与六项代码级断言。 - PostgreSQL `bell.events` 保存不可变事实;后续 outcome 追加到 `bell.event_outcomes`。 -- `bell_runtime` 只拥有两张表的 `SELECT/INSERT`,没有 `UPDATE/DELETE/TRUNCATE` 或 migration owner 权限。 +- `bell_runtime` 对三张不可变事实表只有 `SELECT/INSERT`,没有 `UPDATE/DELETE/TRUNCATE` 或 migration owner 权限;仅可在短期 `audit_relay_receipts` 表查询、插入和清理过期收据。 +- `cmd/bell-api` 默认只监听 `127.0.0.1:8081`,接收 HMAC 签名的 `/internal/v1/audit-events:batch`,把脱敏设备操作事实追加到 `bell.audit_events`。 +- `(key_id, nonce)` 收据保存 10 分钟;相同摘要重放原结果,不同摘要返回冲突。非回环监听必须配置 TLS 证书和私钥。 -Brain→Bell transport、认证、公共事件 API、规则、Alert 和证据对象存储仍需后续任务冻结,不能把 `internal/event` 的 Go 类型当成公共网络协议。 +Brain→Bell transport、公共认证/事件 API、规则、Alert 和证据对象存储仍需后续任务冻结。审计 relay 只服务 Sense,不得把 `internal/event` 的 Go 类型或该 HMAC 适配器当成公共协议。 + +启动内部 receiver 前必须私下设置 `BELL_DB_DSN` 和仓库外绝对路径 `BELL_AUDIT_KEYS_FILE`。远端监听还必须设置 `BELL_TLS_CERT_FILE`、`BELL_TLS_KEY_FILE`;仓库不保存 DSN、key 或证书: + +```powershell +go -C Bell run ./cmd/bell-api +``` ## 验证 diff --git a/Bell/cmd/bell-api/main.go b/Bell/cmd/bell-api/main.go new file mode 100644 index 0000000..769a7e8 --- /dev/null +++ b/Bell/cmd/bell-api/main.go @@ -0,0 +1,152 @@ +package main + +import ( + "context" + "crypto/tls" + "errors" + "fmt" + "log/slog" + "net" + "net/http" + "os" + "os/signal" + "path/filepath" + "syscall" + "time" + + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/stdlib" + + "yovision/bell/internal/audit" + "yovision/bell/internal/store" +) + +var version = "dev" + +type configuration struct { + address string + dsn string + keyFile string + tlsCert string + tlsKey string +} + +func main() { + logger := slog.New(slog.NewJSONHandler(os.Stdout, nil)) + if err := run(logger); err != nil { + logger.Error("Bell stopped", "error", err) + os.Exit(1) + } +} + +func loadConfiguration() (configuration, error) { + value := configuration{ + address: envOr("BELL_HTTP_ADDR", "127.0.0.1:8081"), + dsn: os.Getenv("BELL_DB_DSN"), + keyFile: os.Getenv("BELL_AUDIT_KEYS_FILE"), + tlsCert: os.Getenv("BELL_TLS_CERT_FILE"), + tlsKey: os.Getenv("BELL_TLS_KEY_FILE"), + } + if value.dsn == "" { + return configuration{}, errors.New("BELL_DB_DSN is required") + } + if value.keyFile == "" || !filepath.IsAbs(value.keyFile) { + return configuration{}, errors.New("BELL_AUDIT_KEYS_FILE must be an absolute external path") + } + host, _, err := net.SplitHostPort(value.address) + if err != nil { + return configuration{}, errors.New("invalid BELL_HTTP_ADDR") + } + ip := net.ParseIP(host) + loopback := host == "localhost" || (ip != nil && ip.IsLoopback()) + if !loopback && (value.tlsCert == "" || value.tlsKey == "" || !filepath.IsAbs(value.tlsCert) || !filepath.IsAbs(value.tlsKey)) { + return configuration{}, errors.New("non-loopback Bell bind requires absolute TLS certificate and key paths") + } + if (value.tlsCert == "") != (value.tlsKey == "") { + return configuration{}, errors.New("Bell TLS certificate and key must be configured together") + } + return value, nil +} + +func run(logger *slog.Logger) error { + cfg, err := loadConfiguration() + if err != nil { + return err + } + pgConfig, err := pgx.ParseConfig(cfg.dsn) + if err != nil { + return errors.New("invalid Bell postgres DSN") + } + if pgConfig.RuntimeParams == nil { + pgConfig.RuntimeParams = make(map[string]string) + } + pgConfig.RuntimeParams["application_name"] = "yovision-bell" + db := stdlib.OpenDB(*pgConfig) + db.SetMaxOpenConns(16) + db.SetMaxIdleConns(4) + db.SetConnMaxLifetime(30 * time.Minute) + defer db.Close() + + ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) + defer stop() + repository, err := store.OpenPostgres(ctx, db) + if err != nil { + return err + } + if err := repository.AuditRelayReady(ctx); err != nil { + return err + } + keys, err := audit.LoadKeys(cfg.keyFile) + if err != nil { + return err + } + handler, err := audit.NewHandler(repository, keys) + if err != nil { + return err + } + mux := http.NewServeMux() + mux.Handle(audit.RelayPath, handler) + mux.HandleFunc("GET /healthz", func(writer http.ResponseWriter, _ *http.Request) { + writeStatus(writer, http.StatusOK, "ok") + }) + mux.HandleFunc("GET /readyz", func(writer http.ResponseWriter, request *http.Request) { + if err := repository.AuditRelayReady(request.Context()); err != nil { + writeStatus(writer, http.StatusServiceUnavailable, "not_ready") + return + } + writeStatus(writer, http.StatusOK, "ready") + }) + server := &http.Server{Addr: cfg.address, Handler: mux, ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 15 * time.Second, WriteTimeout: 15 * time.Second, IdleTimeout: 60 * time.Second, TLSConfig: &tls.Config{MinVersion: tls.VersionTLS12}} + serverErrors := make(chan error, 1) + go func() { + logger.Info("Bell listening", "address", cfg.address, "version", version, "tls_enabled", cfg.tlsCert != "") + if cfg.tlsCert != "" { + serverErrors <- server.ListenAndServeTLS(cfg.tlsCert, cfg.tlsKey) + return + } + serverErrors <- server.ListenAndServe() + }() + select { + case <-ctx.Done(): + case serverErr := <-serverErrors: + if !errors.Is(serverErr, http.ErrServerClosed) { + return fmt.Errorf("serve Bell HTTP: %w", serverErr) + } + } + shutdownContext, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + return server.Shutdown(shutdownContext) +} + +func writeStatus(writer http.ResponseWriter, status int, value string) { + writer.Header().Set("Content-Type", "application/json") + writer.WriteHeader(status) + _, _ = fmt.Fprintf(writer, `{"status":%q}`, value) +} + +func envOr(name, fallback string) string { + if value := os.Getenv(name); value != "" { + return value + } + return fallback +} diff --git a/Bell/cmd/bell-api/main_test.go b/Bell/cmd/bell-api/main_test.go new file mode 100644 index 0000000..ed8a7eb --- /dev/null +++ b/Bell/cmd/bell-api/main_test.go @@ -0,0 +1,33 @@ +package main + +import ( + "path/filepath" + "testing" +) + +func TestConfigurationRequiresDatabaseAndExternalKey(t *testing.T) { + t.Setenv("BELL_DB_DSN", "") + t.Setenv("BELL_AUDIT_KEYS_FILE", "") + if _, err := loadConfiguration(); err == nil { + t.Fatal("missing Bell database was accepted") + } + t.Setenv("BELL_DB_DSN", "postgres://bell@127.0.0.1/yovision") + t.Setenv("BELL_AUDIT_KEYS_FILE", "relative.json") + if _, err := loadConfiguration(); err == nil { + t.Fatal("relative Bell key file was accepted") + } +} + +func TestConfigurationRequiresTLSOutsideLoopback(t *testing.T) { + t.Setenv("BELL_DB_DSN", "postgres://bell@127.0.0.1/yovision") + t.Setenv("BELL_AUDIT_KEYS_FILE", filepath.Join(t.TempDir(), "keys.json")) + t.Setenv("BELL_HTTP_ADDR", "0.0.0.0:8081") + if _, err := loadConfiguration(); err == nil { + t.Fatal("remote plaintext Bell bind was accepted") + } + t.Setenv("BELL_TLS_CERT_FILE", filepath.Join(t.TempDir(), "server.crt")) + t.Setenv("BELL_TLS_KEY_FILE", filepath.Join(t.TempDir(), "server.key")) + if _, err := loadConfiguration(); err != nil { + t.Fatalf("remote TLS Bell bind rejected: %v", err) + } +} diff --git a/Bell/internal/audit/audit.go b/Bell/internal/audit/audit.go new file mode 100644 index 0000000..14b43d0 --- /dev/null +++ b/Bell/internal/audit/audit.go @@ -0,0 +1,363 @@ +// Package audit authenticates and validates Sense audit relay batches. +package audit + +import ( + "bytes" + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/base64" + "encoding/hex" + "encoding/json" + "errors" + "io" + "net/http" + "regexp" + "strconv" + "strings" + "time" + "unicode/utf8" +) + +const ( + RelayPath = "/internal/v1/audit-events:batch" + MaxBatchSize = 100 + MaxBodyBytes = 1 << 20 + HeaderKeyID = "X-YoVision-Key-Id" + HeaderTimestamp = "X-YoVision-Timestamp" + HeaderNonce = "X-YoVision-Nonce" + HeaderSignature = "X-YoVision-Signature" +) + +var ErrReplayConflict = errors.New("audit relay replay conflict") + +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 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 Candidate struct { + Envelope Envelope + RecordHash [sha256.Size]byte + ErrorCode string +} + +type Repository interface { + ProcessAuditBatch(context.Context, string, string, [sha256.Size]byte, []Candidate) ([]Result, error) +} + +type Handler struct { + repository Repository + keys map[string][]byte + now func() time.Time +} + +func NewHandler(repository Repository, keys map[string][]byte) (*Handler, error) { + if repository == nil || len(keys) == 0 { + return nil, errors.New("audit handler dependencies are required") + } + copyKeys := make(map[string][]byte, len(keys)) + for id, secret := range keys { + if !keyIDPattern.MatchString(id) || len(secret) < 32 { + return nil, errors.New("invalid audit handler key") + } + copyKeys[id] = append([]byte(nil), secret...) + } + return &Handler{repository: repository, keys: copyKeys, now: time.Now}, nil +} + +func (h *Handler) ServeHTTP(writer http.ResponseWriter, request *http.Request) { + if request.Method != http.MethodPost || request.URL.Path != RelayPath { + writeError(writer, http.StatusNotFound, "not_found") + return + } + body, err := io.ReadAll(io.LimitReader(request.Body, MaxBodyBytes+1)) + if err != nil || len(body) > MaxBodyBytes { + writeError(writer, http.StatusRequestEntityTooLarge, "payload_too_large") + return + } + keyID := request.Header.Get(HeaderKeyID) + timestamp := request.Header.Get(HeaderTimestamp) + nonce := request.Header.Get(HeaderNonce) + provided := request.Header.Get(HeaderSignature) + secret, ok := h.keys[keyID] + seconds, timestampErr := strconv.ParseInt(timestamp, 10, 64) + nonceBytes, nonceErr := base64.RawURLEncoding.DecodeString(nonce) + signatureBytes, signatureErr := base64.RawURLEncoding.DecodeString(provided) + if !ok || timestampErr != nil || len(timestamp) < 10 || nonceErr != nil || len(nonceBytes) < 16 || len(nonceBytes) > 48 || + signatureErr != nil || len(signatureBytes) != sha256.Size || absDuration(h.now().UTC().Sub(time.Unix(seconds, 0).UTC())) > 300*time.Second { + writeError(writer, http.StatusUnauthorized, "unauthorized") + return + } + expected := signature(secret, canonicalString(request.Method, request.URL.EscapedPath(), timestamp, nonce, body)) + if !hmac.Equal(signatureBytes, expected) { + writeError(writer, http.StatusUnauthorized, "unauthorized") + return + } + candidates, err := decodeCandidates(body) + if err != nil { + writeError(writer, http.StatusBadRequest, "invalid_batch") + return + } + requestHash := sha256.Sum256(body) + results, err := h.repository.ProcessAuditBatch(request.Context(), keyID, nonce, requestHash, candidates) + if errors.Is(err, ErrReplayConflict) { + writeError(writer, http.StatusConflict, "replay_conflict") + return + } + if err != nil { + writeError(writer, http.StatusServiceUnavailable, "temporarily_unavailable") + return + } + writeJSON(writer, http.StatusOK, BatchResponse{Results: results}) +} + +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) []byte { + mac := hmac.New(sha256.New, secret) + _, _ = mac.Write([]byte(canonical)) + return mac.Sum(nil) +} + +func absDuration(value time.Duration) time.Duration { + if value < 0 { + return -value + } + return value +} + +type rawBatch struct { + Events []json.RawMessage `json:"events"` +} + +func decodeCandidates(body []byte) ([]Candidate, error) { + decoder := json.NewDecoder(bytes.NewReader(body)) + decoder.DisallowUnknownFields() + var batch rawBatch + if err := decoder.Decode(&batch); err != nil || len(batch.Events) < 1 || len(batch.Events) > MaxBatchSize { + return nil, errors.New("invalid audit batch") + } + var trailing any + if err := decoder.Decode(&trailing); !errors.Is(err, io.EOF) { + return nil, errors.New("invalid audit batch trailing data") + } + values := make([]Candidate, len(batch.Events)) + for index, raw := range batch.Events { + values[index].RecordHash = sha256.Sum256(raw) + if !hasExactEnvelopeShape(raw) { + values[index].ErrorCode = "schema_invalid" + continue + } + itemDecoder := json.NewDecoder(bytes.NewReader(raw)) + itemDecoder.DisallowUnknownFields() + if err := itemDecoder.Decode(&values[index].Envelope); err != nil { + values[index].ErrorCode = "schema_invalid" + continue + } + values[index].ErrorCode = validateEnvelope(values[index].Envelope) + } + return values, nil +} + +var ( + eventIDPattern = regexp.MustCompile(`^audit_[0-9a-f]{32}$`) + logicalIDPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$`) + keyIDPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$`) +) + +func validateEnvelope(value Envelope) string { + event := value.Event + if (value.SchemaVersion != 1 && value.SchemaVersion != 2) || !eventIDPattern.MatchString(event.EventID) || + !logicalIDPattern.MatchString(event.TenantID) || !logicalIDPattern.MatchString(event.SiteID) || !logicalIDPattern.MatchString(event.DeviceID) || + event.AggregateGeneration < 1 || event.OccurredAt.IsZero() || strings.TrimSpace(event.Actor.ID) == "" || utf8.RuneCountInString(event.Actor.ID) > 200 || + (event.Actor.Type != "user" && event.Actor.Type != "service" && event.Actor.Type != "system") || + (event.Reason != nil && utf8.RuneCountInString(*event.Reason) > 500) || (event.TraceID != nil && utf8.RuneCountInString(*event.TraceID) > 128) || + (event.ProjectionVersions.QuotaSourceVersion != nil && *event.ProjectionVersions.QuotaSourceVersion < 1) || + (event.ProjectionVersions.AreaPolicySourceVersion != nil && *event.ProjectionVersions.AreaPolicySourceVersion < 1) { + return "schema_invalid" + } + if event.EventType != "device.created" && event.EventType != "device.desired_state.accepted" && + (event.EventType != "device.configuration.accepted" || value.SchemaVersion != 2) { + return "schema_invalid" + } + var data map[string]any + if err := json.Unmarshal(event.Data, &data); err != nil || data == nil { + return "schema_invalid" + } + expected := map[string]string{ + "device.created": "device_created", + "device.desired_state.accepted": "desired_state_accepted", + "device.configuration.accepted": "configuration_accepted", + }[event.EventType] + if data["kind"] != expected || !validateData(event.EventType, data) || containsSensitiveKey(data) { + return "payload_invalid" + } + return "" +} + +func hasExactEnvelopeShape(raw []byte) bool { + var envelope map[string]json.RawMessage + if err := json.Unmarshal(raw, &envelope); err != nil || !exactRawKeys(envelope, "schema_version", "event") { + return false + } + var event map[string]json.RawMessage + if err := json.Unmarshal(envelope["event"], &event); err != nil || !exactRawKeys(event, + "event_id", "event_type", "tenant_id", "site_id", "device_id", "actor", "reason", "trace_id", + "aggregate_generation", "projection_versions", "data", "occurred_at", + ) { + return false + } + var actor, projections map[string]json.RawMessage + return json.Unmarshal(event["actor"], &actor) == nil && exactRawKeys(actor, "type", "id") && + json.Unmarshal(event["projection_versions"], &projections) == nil && exactRawKeys(projections, "quota_source_version", "area_policy_source_version") +} + +func exactRawKeys(value map[string]json.RawMessage, expected ...string) bool { + if len(value) != len(expected) { + return false + } + for _, key := range expected { + if _, exists := value[key]; !exists { + return false + } + } + return true +} + +func validateData(eventType string, data map[string]any) bool { + switch eventType { + case "device.created": + if !exactAnyKeys(data, "kind", "area_id", "modality", "capabilities", "desired_state") || !logicalIDPattern.MatchString(stringValue(data["area_id"])) { + return false + } + if !member(stringValue(data["modality"]), "video", "radar", "contact", "button", "wearable", "other") || !member(stringValue(data["desired_state"]), "disabled", "enabled") { + return false + } + return validStringSet(data["capabilities"], 16, "video_capture", "audio_capture", "spatial_rule", "telemetry") + case "device.desired_state.accepted": + if !exactAnyKeys(data, "kind", "previous_desired_state", "desired_state", "changed") { + return false + } + _, changed := data["changed"].(bool) + return changed && member(stringValue(data["previous_desired_state"]), "disabled", "enabled") && member(stringValue(data["desired_state"]), "disabled", "enabled") + case "device.configuration.accepted": + if !exactAnyKeys(data, "kind", "changed", "changed_fields", "area_id") || !logicalIDPattern.MatchString(stringValue(data["area_id"])) { + return false + } + _, changed := data["changed"].(bool) + return changed && validStringSet(data["changed_fields"], 5, "name", "area_id", "endpoint_ref", "credential_ref", "profile_token") + default: + return false + } +} + +func exactAnyKeys(value map[string]any, expected ...string) bool { + if len(value) != len(expected) { + return false + } + for _, key := range expected { + if _, exists := value[key]; !exists { + return false + } + } + return true +} + +func stringValue(value any) string { + result, _ := value.(string) + return result +} + +func member(value string, allowed ...string) bool { + for _, candidate := range allowed { + if value == candidate { + return true + } + } + return false +} + +func validStringSet(value any, maximum int, allowed ...string) bool { + items, ok := value.([]any) + if !ok || len(items) > maximum { + return false + } + seen := make(map[string]bool, len(items)) + for _, item := range items { + text, ok := item.(string) + if !ok || !member(text, allowed...) || seen[text] { + return false + } + seen[text] = true + } + return true +} + +func containsSensitiveKey(value any) bool { + forbidden := map[string]bool{"password": true, "stream_uri": true, "mediamtx_config": true} + switch typed := value.(type) { + case map[string]any: + for key, child := range typed { + if forbidden[strings.ToLower(key)] || containsSensitiveKey(child) { + return true + } + } + case []any: + for _, child := range typed { + if containsSensitiveKey(child) { + return true + } + } + } + return false +} + +func writeError(writer http.ResponseWriter, status int, code string) { + writeJSON(writer, status, map[string]string{"error": code}) +} + +func writeJSON(writer http.ResponseWriter, status int, value any) { + writer.Header().Set("Content-Type", "application/json") + writer.Header().Set("Cache-Control", "no-store") + writer.WriteHeader(status) + _ = json.NewEncoder(writer).Encode(value) +} diff --git a/Bell/internal/audit/audit_test.go b/Bell/internal/audit/audit_test.go new file mode 100644 index 0000000..7cc0461 --- /dev/null +++ b/Bell/internal/audit/audit_test.go @@ -0,0 +1,105 @@ +package audit + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/base64" + "encoding/json" + "net/http" + "net/http/httptest" + "strconv" + "testing" + "time" +) + +type recordingRepository struct { + candidates []Candidate + results []Result + err error +} + +func (r *recordingRepository) ProcessAuditBatch(_ context.Context, _, _ string, _ [sha256.Size]byte, values []Candidate) ([]Result, error) { + r.candidates = values + return r.results, r.err +} + +func validBody(t *testing.T) []byte { + t.Helper() + value := map[string]any{"events": []any{map[string]any{ + "schema_version": 1, + "event": map[string]any{ + "event_id": "audit_00000000000000000000000000000001", "event_type": "device.created", + "tenant_id": "tenant", "site_id": "site", "device_id": "camera-1", + "actor": map[string]any{"type": "system", "id": "sense"}, "reason": nil, "trace_id": nil, + "aggregate_generation": 1, + "projection_versions": map[string]any{"quota_source_version": 1, "area_policy_source_version": 1}, + "data": map[string]any{"kind": "device_created", "area_id": "area", "modality": "video", "capabilities": []any{"video_capture"}, "desired_state": "enabled"}, + "occurred_at": "2026-08-11T00:00:00Z", + }, + }}} + raw, err := json.Marshal(value) + if err != nil { + t.Fatal(err) + } + return raw +} + +func signedRequest(t *testing.T, body, secret []byte, timestamp time.Time, nonce string) *http.Request { + t.Helper() + request := httptest.NewRequest(http.MethodPost, RelayPath, bytes.NewReader(body)) + stamp := strconv.FormatInt(timestamp.Unix(), 10) + request.Header.Set(HeaderKeyID, "sense-a") + request.Header.Set(HeaderTimestamp, stamp) + request.Header.Set(HeaderNonce, nonce) + request.Header.Set(HeaderSignature, base64.RawURLEncoding.EncodeToString(signature(secret, canonicalString(http.MethodPost, RelayPath, stamp, nonce, body)))) + return request +} + +func TestHandlerAuthenticatesAndReturnsPerItemResults(t *testing.T) { + secret := bytes.Repeat([]byte{3}, 32) + repository := &recordingRepository{results: []Result{{EventID: "audit_00000000000000000000000000000001", Status: "accepted"}}} + handler, err := NewHandler(repository, map[string][]byte{"sense-a": secret}) + if err != nil { + t.Fatal(err) + } + now := time.Date(2026, 8, 11, 0, 1, 0, 0, time.UTC) + handler.now = func() time.Time { return now } + response := httptest.NewRecorder() + handler.ServeHTTP(response, signedRequest(t, validBody(t), secret, now, "AAAAAAAAAAAAAAAAAAAAAA")) + if response.Code != http.StatusOK || len(repository.candidates) != 1 || repository.candidates[0].ErrorCode != "" { + t.Fatalf("valid batch rejected: status=%d candidates=%+v", response.Code, repository.candidates) + } +} + +func TestHandlerRejectsStaleOrTamperedRequests(t *testing.T) { + secret := bytes.Repeat([]byte{4}, 32) + repository := &recordingRepository{} + handler, _ := NewHandler(repository, map[string][]byte{"sense-a": secret}) + now := time.Date(2026, 8, 11, 0, 10, 0, 0, time.UTC) + handler.now = func() time.Time { return now } + for _, request := range []*http.Request{ + signedRequest(t, validBody(t), secret, now.Add(-301*time.Second), "BBBBBBBBBBBBBBBBBBBBBB"), + signedRequest(t, append(validBody(t), ' '), bytes.Repeat([]byte{5}, 32), now, "CCCCCCCCCCCCCCCCCCCCCC"), + } { + response := httptest.NewRecorder() + handler.ServeHTTP(response, request) + if response.Code != http.StatusUnauthorized { + t.Fatalf("unsafe request returned %d", response.Code) + } + } +} + +func TestDecodeCandidatesRejectsSensitiveItemWithoutRejectingBatch(t *testing.T) { + body := validBody(t) + var value map[string]any + _ = json.Unmarshal(body, &value) + events := value["events"].([]any) + event := events[0].(map[string]any)["event"].(map[string]any) + event["data"].(map[string]any)["password"] = "must-not-persist" + body, _ = json.Marshal(value) + candidates, err := decodeCandidates(body) + if err != nil || len(candidates) != 1 || candidates[0].ErrorCode != "payload_invalid" { + t.Fatalf("unexpected per-item validation: %+v %v", candidates, err) + } +} diff --git a/Bell/internal/audit/key.go b/Bell/internal/audit/key.go new file mode 100644 index 0000000..5cd750b --- /dev/null +++ b/Bell/internal/audit/key.go @@ -0,0 +1,47 @@ +package audit + +import ( + "bytes" + "encoding/base64" + "encoding/json" + "errors" + "io" + "os" +) + +type keyDocument struct { + Version int `json:"version"` + Keys []struct { + KeyID string `json:"key_id"` + Secret string `json:"secret_base64url"` + } `json:"keys"` +} + +func LoadKeys(path string) (map[string][]byte, error) { + raw, err := os.ReadFile(path) + if err != nil { + return nil, errors.New("read Bell audit key file") + } + var document keyDocument + 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 Bell audit key file") + } + var trailing any + if err := decoder.Decode(&trailing); !errors.Is(err, io.EOF) { + return nil, errors.New("invalid Bell audit key file") + } + values := make(map[string][]byte, len(document.Keys)) + for _, item := range document.Keys { + secret, err := base64.RawURLEncoding.DecodeString(item.Secret) + if err != nil || !keyIDPattern.MatchString(item.KeyID) || len(secret) < 32 { + return nil, errors.New("invalid Bell audit key") + } + if _, exists := values[item.KeyID]; exists { + return nil, errors.New("duplicate Bell audit key ID") + } + values[item.KeyID] = secret + } + return values, nil +} diff --git a/Bell/internal/store/audit_postgres.go b/Bell/internal/store/audit_postgres.go new file mode 100644 index 0000000..c837ed4 --- /dev/null +++ b/Bell/internal/store/audit_postgres.go @@ -0,0 +1,138 @@ +package store + +import ( + "bytes" + "context" + "crypto/sha256" + "database/sql" + "encoding/json" + "errors" + "fmt" + + "yovision/bell/internal/audit" +) + +func (p *Postgres) AuditRelayReady(ctx context.Context) error { + var version int64 + if err := p.db.QueryRowContext(ctx, `SELECT COALESCE(MAX(version), 0) FROM bell.schema_migrations`).Scan(&version); err != nil || version < 4 { + return errors.New("postgres Bell schema migration v4 is required for audit relay") + } + var auditSelect, auditInsert, auditUpdate, auditDelete, auditTruncate bool + var receiptUse bool + if err := p.db.QueryRowContext(ctx, `SELECT + has_table_privilege(current_user, 'bell.audit_events', 'SELECT'), + has_table_privilege(current_user, 'bell.audit_events', 'INSERT'), + has_table_privilege(current_user, 'bell.audit_events', 'UPDATE'), + has_table_privilege(current_user, 'bell.audit_events', 'DELETE'), + has_table_privilege(current_user, 'bell.audit_events', 'TRUNCATE'), + has_table_privilege(current_user, 'bell.audit_relay_receipts', 'SELECT,INSERT,DELETE')`).Scan( + &auditSelect, &auditInsert, &auditUpdate, &auditDelete, &auditTruncate, &receiptUse, + ); err != nil { + return errors.New("verify Bell audit relay privileges") + } + if !auditSelect || !auditInsert || auditUpdate || auditDelete || auditTruncate || !receiptUse { + return errors.New("Bell audit relay privileges violate append-only boundary") + } + return nil +} + +func (p *Postgres) ProcessAuditBatch( + ctx context.Context, + keyID, nonce string, + requestHash [sha256.Size]byte, + candidates []audit.Candidate, +) ([]audit.Result, error) { + if len(candidates) < 1 || len(candidates) > audit.MaxBatchSize { + return nil, errors.New("invalid audit candidate batch") + } + tx, err := p.db.BeginTx(ctx, nil) + if err != nil { + return nil, errors.New("begin Bell audit batch") + } + defer tx.Rollback() + if _, err := tx.ExecContext(ctx, `SELECT pg_advisory_xact_lock(hashtext($1), hashtext($2))`, keyID, nonce); err != nil { + return nil, errors.New("lock Bell audit receipt") + } + if _, err := tx.ExecContext(ctx, `DELETE FROM bell.audit_relay_receipts WHERE expires_at <= clock_timestamp()`); err != nil { + return nil, errors.New("expire Bell audit receipts") + } + var existingHash, existingBody []byte + err = tx.QueryRowContext(ctx, `SELECT request_hash, response_body::text + FROM bell.audit_relay_receipts WHERE key_id=$1 AND nonce=$2`, keyID, nonce).Scan(&existingHash, &existingBody) + if err == nil { + if !bytes.Equal(existingHash, requestHash[:]) { + return nil, audit.ErrReplayConflict + } + var response audit.BatchResponse + if err := json.Unmarshal(existingBody, &response); err != nil { + return nil, errors.New("decode stored Bell audit receipt") + } + if err := tx.Commit(); err != nil { + return nil, errors.New("commit Bell audit replay") + } + return response.Results, nil + } + if !errors.Is(err, sql.ErrNoRows) { + return nil, errors.New("read Bell audit receipt") + } + results := make([]audit.Result, 0, len(candidates)) + for _, candidate := range candidates { + if candidate.ErrorCode != "" { + code := candidate.ErrorCode + results = append(results, audit.Result{EventID: candidate.Envelope.Event.EventID, Status: "rejected", ErrorCode: &code}) + continue + } + event := candidate.Envelope.Event + payload, err := json.Marshal(event) + if err != nil { + return nil, errors.New("encode Bell audit fact") + } + result, err := tx.ExecContext(ctx, `INSERT INTO bell.audit_events( + source_system, event_id, schema_version, event_type, tenant_id, site_id, + device_id, actor_type, actor_id, reason, trace_id, aggregate_generation, + quota_source_version, area_policy_source_version, payload, occurred_at, record_hash + ) VALUES ('sense',$1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14::jsonb,$15,$16) + ON CONFLICT (source_system,event_id) DO NOTHING`, + event.EventID, candidate.Envelope.SchemaVersion, event.EventType, event.TenantID, + event.SiteID, event.DeviceID, event.Actor.Type, event.Actor.ID, event.Reason, + event.TraceID, event.AggregateGeneration, + event.ProjectionVersions.QuotaSourceVersion, + event.ProjectionVersions.AreaPolicySourceVersion, + payload, event.OccurredAt, candidate.RecordHash[:]) + if err != nil { + return nil, fmt.Errorf("insert Bell audit fact: %w", err) + } + affected, err := result.RowsAffected() + if err != nil { + return nil, errors.New("read Bell audit insert result") + } + if affected == 1 { + results = append(results, audit.Result{EventID: event.EventID, Status: "accepted"}) + continue + } + var storedHash []byte + if err := tx.QueryRowContext(ctx, `SELECT record_hash FROM bell.audit_events + WHERE source_system='sense' AND event_id=$1`, event.EventID).Scan(&storedHash); err != nil { + return nil, errors.New("read existing Bell audit fact") + } + if bytes.Equal(storedHash, candidate.RecordHash[:]) { + results = append(results, audit.Result{EventID: event.EventID, Status: "duplicate"}) + } else { + code := "id_conflict" + results = append(results, audit.Result{EventID: event.EventID, Status: "rejected", ErrorCode: &code}) + } + } + encoded, err := json.Marshal(audit.BatchResponse{Results: results}) + if err != nil { + return nil, errors.New("encode Bell audit response") + } + if _, err := tx.ExecContext(ctx, `INSERT INTO bell.audit_relay_receipts( + key_id, nonce, request_hash, response_status, response_body, expires_at + ) VALUES ($1,$2,$3,200,$4::jsonb,clock_timestamp() + interval '10 minutes')`, keyID, nonce, requestHash[:], encoded); err != nil { + return nil, errors.New("insert Bell audit receipt") + } + if err := tx.Commit(); err != nil { + return nil, errors.New("commit Bell audit batch") + } + return results, nil +} diff --git a/Bell/internal/store/audit_postgres_test.go b/Bell/internal/store/audit_postgres_test.go new file mode 100644 index 0000000..d32dee6 --- /dev/null +++ b/Bell/internal/store/audit_postgres_test.go @@ -0,0 +1,82 @@ +package store + +import ( + "context" + "crypto/sha256" + "database/sql" + "encoding/json" + "errors" + "os" + "testing" + "time" + + _ "github.com/jackc/pgx/v5/stdlib" + + "yovision/bell/internal/audit" +) + +func auditCandidate(t *testing.T, eventID, actorID string) audit.Candidate { + t.Helper() + data := json.RawMessage(`{"kind":"device_created","area_id":"area","modality":"video","capabilities":["video_capture"],"desired_state":"enabled"}`) + value := audit.Envelope{SchemaVersion: 1, Event: audit.Event{ + EventID: eventID, EventType: "device.created", TenantID: "tenant", SiteID: "site", DeviceID: "camera-1", + Actor: audit.Actor{Type: "system", ID: actorID}, AggregateGeneration: 1, + ProjectionVersions: audit.ProjectionVersions{}, Data: data, OccurredAt: time.Date(2026, 8, 11, 0, 0, 0, 0, time.UTC), + }} + raw, err := json.Marshal(value) + if err != nil { + t.Fatal(err) + } + return audit.Candidate{Envelope: value, RecordHash: sha256.Sum256(raw)} +} + +func TestPostgresAuditBatchReceiptAndImmutableFact(t *testing.T) { + dsn := os.Getenv("YOVISION_TEST_BELL_POSTGRES_DSN") + if dsn == "" { + t.Skip("YOVISION_TEST_BELL_POSTGRES_DSN is not set") + } + db, err := sql.Open("pgx", dsn) + if err != nil { + t.Fatal(err) + } + defer db.Close() + ctx := context.Background() + repository, err := OpenPostgres(ctx, db) + if err != nil { + t.Fatal(err) + } + if err := repository.AuditRelayReady(ctx); err != nil { + t.Fatal(err) + } + + requestHash := sha256.Sum256([]byte("request-one")) + eventID := "audit_10000000000000000000000000000001" + results, err := repository.ProcessAuditBatch(ctx, "sense-a", "AAAAAAAAAAAAAAAAAAAAAA", requestHash, []audit.Candidate{auditCandidate(t, eventID, "sense")}) + if err != nil || len(results) != 1 || results[0].Status != "accepted" { + t.Fatalf("first batch: %+v %v", results, err) + } + replayed, err := repository.ProcessAuditBatch(ctx, "sense-a", "AAAAAAAAAAAAAAAAAAAAAA", requestHash, []audit.Candidate{auditCandidate(t, eventID, "ignored-by-receipt")}) + if err != nil || replayed[0].Status != "accepted" { + t.Fatalf("receipt replay: %+v %v", replayed, err) + } + different := sha256.Sum256([]byte("request-two")) + if _, err := repository.ProcessAuditBatch(ctx, "sense-a", "AAAAAAAAAAAAAAAAAAAAAA", different, []audit.Candidate{auditCandidate(t, eventID, "sense")}); !errors.Is(err, audit.ErrReplayConflict) { + t.Fatalf("expected replay conflict, got %v", err) + } + + duplicate, err := repository.ProcessAuditBatch(ctx, "sense-a", "BBBBBBBBBBBBBBBBBBBBBB", different, []audit.Candidate{auditCandidate(t, eventID, "sense")}) + if err != nil || duplicate[0].Status != "duplicate" { + t.Fatalf("event duplicate: %+v %v", duplicate, err) + } + conflictHash := sha256.Sum256([]byte("request-three")) + conflict, err := repository.ProcessAuditBatch(ctx, "sense-a", "CCCCCCCCCCCCCCCCCCCCCC", conflictHash, []audit.Candidate{auditCandidate(t, eventID, "other")}) + if err != nil || conflict[0].Status != "rejected" || conflict[0].ErrorCode == nil || *conflict[0].ErrorCode != "id_conflict" { + t.Fatalf("event conflict: %+v %v", conflict, err) + } + if _, err := db.ExecContext(ctx, `UPDATE bell.audit_events SET actor_id='mutated' WHERE event_id=$1`, eventID); err == nil { + t.Fatal("runtime updated immutable audit fact") + } + if _, err := db.ExecContext(ctx, `DELETE FROM bell.audit_events WHERE event_id=$1`, eventID); err == nil { + t.Fatal("runtime deleted immutable audit fact") + } +} diff --git a/Sense/README.md b/Sense/README.md index c856bcc..b7f36c3 100644 --- a/Sense/README.md +++ b/Sense/README.md @@ -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://` 凭据引用。真实适配器从进程环境读取以下变量,不把秘密写入 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 下不会启动。 ### 孤儿报告与受控处置 diff --git a/Sense/cmd/sense-api/main.go b/Sense/cmd/sense-api/main.go index 7699160..7366e88 100644 --- a/Sense/cmd/sense-api/main.go +++ b/Sense/cmd/sense-api/main.go @@ -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() }() diff --git a/Sense/internal/auditrelay/auditrelay_test.go b/Sense/internal/auditrelay/auditrelay_test.go new file mode 100644 index 0000000..8b5645c --- /dev/null +++ b/Sense/internal/auditrelay/auditrelay_test.go @@ -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 +} diff --git a/Sense/internal/auditrelay/client.go b/Sense/internal/auditrelay/client.go new file mode 100644 index 0000000..33a1904 --- /dev/null +++ b/Sense/internal/auditrelay/client.go @@ -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 +} diff --git a/Sense/internal/auditrelay/key.go b/Sense/internal/auditrelay/key.go new file mode 100644 index 0000000..7644345 --- /dev/null +++ b/Sense/internal/auditrelay/key.go @@ -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") +} diff --git a/Sense/internal/auditrelay/model.go b/Sense/internal/auditrelay/model.go new file mode 100644 index 0000000..71e6e89 --- /dev/null +++ b/Sense/internal/auditrelay/model.go @@ -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 +} diff --git a/Sense/internal/auditrelay/worker.go b/Sense/internal/auditrelay/worker.go new file mode 100644 index 0000000..09e5f60 --- /dev/null +++ b/Sense/internal/auditrelay/worker.go @@ -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: + } + } +} diff --git a/Sense/internal/config/config.go b/Sense/internal/config/config.go index a5ec521..9c9ecc8 100644 --- a/Sense/internal/config/config.go +++ b/Sense/internal/config/config.go @@ -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 } diff --git a/Sense/internal/config/config_test.go b/Sense/internal/config/config_test.go index 719bd93..bcc5425 100644 --- a/Sense/internal/config/config_test.go +++ b/Sense/internal/config/config_test.go @@ -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") + } +} diff --git a/Sense/internal/store/audit_relay_postgres.go b/Sense/internal/store/audit_relay_postgres.go new file mode 100644 index 0000000..5f868fb --- /dev/null +++ b/Sense/internal/store/audit_relay_postgres.go @@ -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 +} diff --git a/Sense/internal/store/audit_relay_postgres_test.go b/Sense/internal/store/audit_relay_postgres_test.go new file mode 100644 index 0000000..c23829c --- /dev/null +++ b/Sense/internal/store/audit_relay_postgres_test.go @@ -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) + } +} diff --git a/deploy/postgres/014_audit_relay.sql b/deploy/postgres/014_audit_relay.sql new file mode 100644 index 0000000..e6bc3ff --- /dev/null +++ b/deploy/postgres/014_audit_relay.sql @@ -0,0 +1,130 @@ +-- Sense Outbox relay fencing and Bell global audit facts. + +ALTER TABLE sense.device_operation_outbox + ADD COLUMN IF NOT EXISTS relay_lease_owner text, + ADD COLUMN IF NOT EXISTS relay_lease_token bigint NOT NULL DEFAULT 0, + ADD COLUMN IF NOT EXISTS relay_lease_until timestamptz, + ADD COLUMN IF NOT EXISTS last_error_code text, + ADD COLUMN IF NOT EXISTS dead_lettered_at timestamptz; + +ALTER TABLE sense.device_operation_outbox + DROP CONSTRAINT IF EXISTS sense_outbox_relay_lease_pair; +ALTER TABLE sense.device_operation_outbox + ADD CONSTRAINT sense_outbox_relay_lease_pair CHECK ( + (relay_lease_owner IS NULL AND relay_lease_until IS NULL) + OR (relay_lease_owner IS NOT NULL AND btrim(relay_lease_owner) <> '' + AND relay_lease_until IS NOT NULL) + ); +ALTER TABLE sense.device_operation_outbox + DROP CONSTRAINT IF EXISTS sense_outbox_relay_token_nonnegative; +ALTER TABLE sense.device_operation_outbox + ADD CONSTRAINT sense_outbox_relay_token_nonnegative CHECK (relay_lease_token >= 0); +ALTER TABLE sense.device_operation_outbox + DROP CONSTRAINT IF EXISTS sense_outbox_relay_error_length; +ALTER TABLE sense.device_operation_outbox + ADD CONSTRAINT sense_outbox_relay_error_length CHECK ( + last_error_code IS NULL OR ( + char_length(last_error_code) BETWEEN 1 AND 64 + AND last_error_code ~ '^[a-z][a-z0-9_]*$' + ) + ); +ALTER TABLE sense.device_operation_outbox + DROP CONSTRAINT IF EXISTS sense_outbox_relay_terminal_state; +ALTER TABLE sense.device_operation_outbox + ADD CONSTRAINT sense_outbox_relay_terminal_state CHECK ( + delivered_at IS NULL OR dead_lettered_at IS NULL + ); + +CREATE INDEX IF NOT EXISTS sense_outbox_relay_due_idx + ON sense.device_operation_outbox( + COALESCE(next_attempt_at, available_at), event_id + ) + WHERE delivered_at IS NULL AND dead_lettered_at IS NULL; + +CREATE TABLE IF NOT EXISTS bell.audit_events ( + source_system text NOT NULL, + event_id text NOT NULL, + schema_version smallint NOT NULL, + event_type text NOT NULL, + tenant_id text NOT NULL, + site_id text NOT NULL, + device_id text NOT NULL, + actor_type text NOT NULL, + actor_id text NOT NULL, + reason text, + trace_id text, + aggregate_generation bigint NOT NULL, + quota_source_version bigint, + area_policy_source_version bigint, + payload jsonb NOT NULL, + occurred_at timestamptz NOT NULL, + received_at timestamptz NOT NULL DEFAULT clock_timestamp(), + record_hash bytea NOT NULL, + PRIMARY KEY (source_system, event_id), + CONSTRAINT bell_audit_source_system CHECK (source_system = 'sense'), + CONSTRAINT bell_audit_event_id CHECK (event_id ~ '^audit_[0-9a-f]{32}$'), + CONSTRAINT bell_audit_schema_version CHECK (schema_version IN (1, 2)), + CONSTRAINT bell_audit_event_type CHECK (event_type IN ( + 'device.created', + 'device.desired_state.accepted', + 'device.configuration.accepted' + )), + CONSTRAINT bell_audit_identity_not_blank CHECK ( + btrim(tenant_id) <> '' AND btrim(site_id) <> '' AND btrim(device_id) <> '' + AND btrim(actor_id) <> '' + ), + CONSTRAINT bell_audit_actor_type CHECK (actor_type IN ('user', 'service', 'system')), + CONSTRAINT bell_audit_generation_positive CHECK (aggregate_generation >= 1), + CONSTRAINT bell_audit_projection_versions CHECK ( + (quota_source_version IS NULL OR quota_source_version >= 1) + AND (area_policy_source_version IS NULL OR area_policy_source_version >= 1) + ), + CONSTRAINT bell_audit_reason_length CHECK (reason IS NULL OR char_length(reason) <= 500), + CONSTRAINT bell_audit_trace_length CHECK (trace_id IS NULL OR char_length(trace_id) <= 128), + CONSTRAINT bell_audit_payload_object CHECK (jsonb_typeof(payload) = 'object'), + CONSTRAINT bell_audit_record_hash_length CHECK (octet_length(record_hash) = 32) +); +ALTER TABLE bell.audit_events OWNER TO bell_app; + +DROP TRIGGER IF EXISTS bell_audit_events_immutable ON bell.audit_events; +CREATE TRIGGER bell_audit_events_immutable +BEFORE UPDATE OR DELETE ON bell.audit_events +FOR EACH ROW EXECUTE FUNCTION bell.reject_immutable_change(); + +CREATE INDEX IF NOT EXISTS bell_audit_scope_time_idx + ON bell.audit_events(tenant_id, site_id, occurred_at DESC, event_id DESC); +CREATE INDEX IF NOT EXISTS bell_audit_device_time_idx + ON bell.audit_events(tenant_id, site_id, device_id, occurred_at DESC, event_id DESC); + +CREATE TABLE IF NOT EXISTS bell.audit_relay_receipts ( + key_id text NOT NULL, + nonce text NOT NULL, + request_hash bytea NOT NULL, + response_status integer NOT NULL, + response_body jsonb NOT NULL, + received_at timestamptz NOT NULL DEFAULT clock_timestamp(), + expires_at timestamptz NOT NULL, + PRIMARY KEY (key_id, nonce), + CONSTRAINT bell_audit_receipt_key_id CHECK ( + char_length(key_id) BETWEEN 1 AND 64 AND key_id ~ '^[A-Za-z0-9][A-Za-z0-9._-]*$' + ), + CONSTRAINT bell_audit_receipt_nonce CHECK ( + char_length(nonce) BETWEEN 22 AND 64 AND nonce ~ '^[A-Za-z0-9_-]+$' + ), + CONSTRAINT bell_audit_receipt_hash_length CHECK (octet_length(request_hash) = 32), + CONSTRAINT bell_audit_receipt_status CHECK (response_status = 200), + CONSTRAINT bell_audit_receipt_body CHECK (jsonb_typeof(response_body) = 'object'), + CONSTRAINT bell_audit_receipt_ttl CHECK ( + expires_at >= received_at + interval '10 minutes' + AND expires_at <= received_at + interval '11 minutes' + ) +); +ALTER TABLE bell.audit_relay_receipts OWNER TO bell_app; + +CREATE INDEX IF NOT EXISTS bell_audit_receipt_expiry_idx + ON bell.audit_relay_receipts(expires_at, key_id, nonce); + +INSERT INTO bell.schema_migrations(version) VALUES (4) +ON CONFLICT (version) DO NOTHING; +INSERT INTO sense.schema_migrations(version) VALUES (6) +ON CONFLICT (version) DO NOTHING; diff --git a/deploy/postgres/015_privileges_audit_relay.sql b/deploy/postgres/015_privileges_audit_relay.sql new file mode 100644 index 0000000..2acfce6 --- /dev/null +++ b/deploy/postgres/015_privileges_audit_relay.sql @@ -0,0 +1,11 @@ +-- Bell runtime may append global audit facts and maintain only short-lived +-- idempotency receipts. Sense keeps ownership of its local Outbox. + +REVOKE ALL ON TABLE bell.audit_events, bell.audit_relay_receipts FROM PUBLIC; +REVOKE ALL ON TABLE bell.audit_events, bell.audit_relay_receipts FROM bell_runtime; + +GRANT SELECT, INSERT ON TABLE bell.audit_events TO bell_runtime; +GRANT SELECT, INSERT, DELETE ON TABLE bell.audit_relay_receipts TO bell_runtime; + +REVOKE ALL ON TABLE bell.audit_events, bell.audit_relay_receipts FROM sense_app; +GRANT SELECT, INSERT, UPDATE, DELETE ON TABLE sense.device_operation_outbox TO sense_app; diff --git a/deploy/postgres/README.md b/deploy/postgres/README.md index 4967d80..04b5ab3 100644 --- a/deploy/postgres/README.md +++ b/deploy/postgres/README.md @@ -1,6 +1,6 @@ # YoVision PostgreSQL 初始化 -本目录实现 T-009~T-012、T-015 的 PostgreSQL `17.10` schema。SQL 必须按文件名前缀顺序执行:`001`~`004` 创建 NOLOGIN 权限角色、Bell/Sense 初始对象和配额权限;`005`~`007` 增量增加 Area 与审计;`008`~`009` 增加 Control API 状态;`010`~`011` 增加调和 fencing、MediaMTX Path 历史归属、孤儿报告/受控处置结果;`012`~`013` 增加 Bell 不可变事件、append-only outcome 和独立 `bell_runtime` 最小权限。全部 SQL 可重放。对象 owner/迁移角色为 `bell_app`/`sense_app`;应用登录角色及密码由部署环境或密钥系统创建,Sense 登录加入 `sense_app`,Bell 运行登录只加入 `bell_runtime`,仓库不保存登录凭据。 +本目录实现 T-009~T-012、T-015~T-016 的 PostgreSQL `17.10` schema。SQL 必须按文件名前缀顺序执行:`001`~`004` 创建 NOLOGIN 权限角色、Bell/Sense 初始对象和配额权限;`005`~`007` 增量增加 Area 与本地审计;`008`~`009` 增加 Control API 状态;`010`~`011` 增加调和 fencing、MediaMTX Path 历史归属、孤儿报告/受控处置结果;`012`~`013` 增加 Bell 不可变事件、append-only outcome 和独立 `bell_runtime` 最小权限;`014`~`015` 增加审计 Outbox relay fencing、Bell 全局审计事实、短期防重收据和双方最小权限。全部 SQL 可重放。对象 owner/迁移角色为 `bell_app`/`sense_app`;应用登录角色及密码由部署环境或密钥系统创建,Sense 登录加入 `sense_app`,Bell 运行登录只加入 `bell_runtime`,仓库不保存登录凭据。 生产/共享实例必须由管理员先备份并在 YoVision 专用数据库中执行。Sense 进程不会用高权限自动建库或建角色。示例只使用私有环境变量,不把实际 DSN 写入脚本或日志: @@ -20,7 +20,8 @@ Get-ChildItem deploy/postgres/[0-9][0-9][0-9]_*.sql | - `bell_app` 拥有 `bell.sites`/`bell.areas`、版本 trigger、`bell.site_quota_v1` 和 `bell.area_policy_v1`。 - `sense_app` 拥有 `sense` schema,只获得 `bell` schema 的 `USAGE` 和两个投影视图的 `SELECT`。 - `sense_app` 对 Bell 源表、Bell migration 表和 trigger function 没有权限;启动检查发现权限过宽时拒绝运行。 -- `sense.device_operation_outbox` 是本地持久化审计事实,不是 Bell 全局审计真相;relay 的传输、签名、确认和留存尚未实现。 +- `sense.device_operation_outbox` 是本地持久化队列,Bell 全局审计真相只写入 `bell.audit_events`。Sense 只领取/确认本地行,不获得 Bell 表权限;Bell 只通过签名 HTTP ingress 收取,不读取 Outbox。 +- `bell.audit_events` 与事件事实一样不可更新/删除;`bell.audit_relay_receipts` 只为 10 分钟 nonce 幂等窗口保留,Bell runtime 仅可在这张限定表中查询、插入和清理过期记录。 - `bell.events` 与 `bell.event_outcomes` 由 `bell_app` 拥有;`bell_runtime` 只获得 `SELECT/INSERT`,没有 owner、`UPDATE`、`DELETE` 或 `TRUNCATE` 权限,数据库 trigger 再拒绝 owner 路径的意外事实改写。 - `sense.control_idempotency_receipts` 不保存原始 Idempotency-Key,只保存 scope/request SHA-256 和脱敏响应快照;`batch_operations`/items 只保存逻辑 ID、状态和稳定错误,不保存连接秘密。 - 调和与孤儿租约使用 PostgreSQL `clock_timestamp()` 和 fencing token;过期 worker 不能提交完成/失败或扫描报告。`media_path_ownership`、扫描和处置表不保存 endpoint、credential 或 source URI;数据库约束禁止为 `unowned` finding 写删除结果。 diff --git a/deploy/postgres/tests/assertions.sql b/deploy/postgres/tests/assertions.sql index 1f6f1b8..db87870 100644 --- a/deploy/postgres/tests/assertions.sql +++ b/deploy/postgres/tests/assertions.sql @@ -71,8 +71,8 @@ BEGIN OR NOT has_table_privilege('sense_app', 'sense.device_operation_outbox', 'DELETE') THEN RAISE EXCEPTION 'sense_app lacks access to its local audit Outbox'; END IF; - IF (SELECT max(version) FROM bell.schema_migrations) <> 3 - OR (SELECT max(version) FROM sense.schema_migrations) <> 5 THEN + IF (SELECT max(version) FROM bell.schema_migrations) <> 4 + OR (SELECT max(version) FROM sense.schema_migrations) <> 6 THEN RAISE EXCEPTION 'schema migration version drift'; END IF; @@ -201,3 +201,30 @@ BEGIN END IF; END $reconcile_safety$; + +DO $audit_relay$ +BEGIN + IF NOT has_table_privilege('yovision_t012_sense', 'sense.device_operation_outbox', 'SELECT,INSERT,UPDATE,DELETE') + OR has_table_privilege('yovision_t012_sense', 'bell.audit_events', 'SELECT,INSERT,UPDATE,DELETE') + OR has_table_privilege('yovision_t012_sense', 'bell.audit_relay_receipts', 'SELECT,INSERT,UPDATE,DELETE') THEN + RAISE EXCEPTION 'Sense runtime violates audit relay schema ownership'; + END IF; + IF NOT has_table_privilege('yovision_t015_bell', 'bell.audit_events', 'SELECT,INSERT') + OR has_table_privilege('yovision_t015_bell', 'bell.audit_events', 'UPDATE,DELETE,TRUNCATE') + OR NOT has_table_privilege('yovision_t015_bell', 'bell.audit_relay_receipts', 'SELECT,INSERT,DELETE') + OR has_table_privilege('yovision_t015_bell', 'bell.audit_relay_receipts', 'UPDATE,TRUNCATE') THEN + RAISE EXCEPTION 'Bell runtime violates audit relay privileges'; + END IF; + IF has_table_privilege('public', 'bell.audit_events', 'SELECT,INSERT,UPDATE,DELETE,TRUNCATE') + OR has_table_privilege('public', 'bell.audit_relay_receipts', 'SELECT,INSERT,UPDATE,DELETE,TRUNCATE') THEN + RAISE EXCEPTION 'Bell audit relay tables leaked to PUBLIC'; + END IF; + IF NOT EXISTS ( + SELECT 1 FROM information_schema.columns + WHERE table_schema='sense' AND table_name='device_operation_outbox' + AND column_name='relay_lease_token' + ) THEN + RAISE EXCEPTION 'Sense audit relay fencing columns are missing'; + END IF; +END +$audit_relay$; diff --git a/docs/00-ai-start-here.md b/docs/00-ai-start-here.md index e6d6be7..f97d2a7 100644 --- a/docs/00-ai-start-here.md +++ b/docs/00-ai-start-here.md @@ -40,7 +40,7 @@ MVP 以默认 16 路跑通一个场景的端到端闭环;架构、数据和 UI ## 当前阶段 -当前为 **M0 指定型号实机准入、M1 Sense 五路混合源集成和 M2 本地 16 路软件基线均已完成,M3 已建立 Bell 不可变事件存储基础**。后续本地开发统一使用已准入的一台海康样机,多路软件闭环使用独立合成 RTSP 源补足;真实多设备证据延后到客户/借用/租赁条件具备时执行。客户网络尚未提供,T-013 WireGuard 继续后置,不阻塞 Sense Outbox → Bell 审计 relay。 +当前为 **M0 指定型号实机准入、M1 Sense 五路混合源集成和 M2 本地 16 路软件基线均已完成,M3 已建立 Bell 不可变事件存储及 Sense→Bell 全局审计 relay 基础**。后续本地开发统一使用已准入的一台海康样机,多路软件闭环使用独立合成 RTSP 源补足;真实多设备证据延后到客户/借用/租赁条件具备时执行。客户网络尚未提供,T-013 WireGuard 继续后置,不阻塞 Brain/Bell 本地事件链开发。 优先路径: diff --git a/docs/03-tech-stack.md b/docs/03-tech-stack.md index 7f5803f..850c990 100644 --- a/docs/03-tech-stack.md +++ b/docs/03-tech-stack.md @@ -67,6 +67,10 @@ T-014 没有增加生产依赖。Windows 容量脚本冻结并核对 Sense 模 T-015 不冻结 Brain→Bell transport,也不产生可部署 Bell API 二进制。内部 factory 接收不含 `id` 的候选事实,由 Bell 生成 ULID 后才形成最终 v0.1 事件;不得把该 Go 类型当成公共网络协议。 +### 1.4 Sense 审计 relay(T-016) + +T-016 不增加第三方依赖:两端使用 Go 标准库 HTTP、HMAC-SHA256、SHA-256、base64url 和 constant-time compare,数据库继续使用已冻结的 PostgreSQL 17.10/pgx。`cmd/bell-api` 只提供回环 health/ready 和 Sense 审计内部端点;非回环监听必须同时提供绝对路径 TLS 证书/私钥。HMAC key 使用仓库外 version 1 JSON 文件,secret 至少 32 字节;该适配器不替代未来 Bell 公共 JWT/OIDC。 + ## 2. 外部项目边界 - MiBeeNvr:只用于 M0 隔离实验室、ONVIF兼容性和交互参考,不作为生产依赖。 @@ -78,7 +82,7 @@ T-015 不冻结 Brain→Bell transport,也不产生可部署 Bell API 二进 - Python、Savant/DeepStream 的精确版本;Go、MediaMTX 与 PostgreSQL 已分别为 Sense M1/M2 冻结,后续阶段可按升级流程调整。 - Bell 前端框架和组件库。 -- 事件投递 transport 从 HTTP 起步还是直接采用消息总线。 +- Brain→Bell 业务事件投递 transport;Sense→Bell 审计 relay 已独立冻结为内部 HTTP,不能据此默认 Brain transport。 - 目标 GPU/边缘硬件、解码能力和每 worker 的 `max_sources`。 - MinIO/S3 的精确版本、加密实现,以及客户/法务确认后的最终生命周期策略。 - 短信/语音供应商及生产双路径组合;是否开发原生 App 最早在 M4 根据试点反馈决定。 @@ -110,7 +114,7 @@ go -C Sense build ./... go -C Sense run ./cmd/sense-api ``` -Bell 事件域基础单独执行(当前没有可启动 API): +Bell 事件域与内部审计 receiver 单独执行: ```powershell go -C Bell mod download diff --git a/docs/04-architecture.md b/docs/04-architecture.md index b879e16..f589be5 100644 --- a/docs/04-architecture.md +++ b/docs/04-architecture.md @@ -68,7 +68,7 @@ Sense ── 视频流/触发信号 ──> Brain 9. 投递状态机只依赖 Bell provider 接口,不直接依赖某家短信或语音 SDK;生产前至少两条独立路径并能故障切换。 10. 设备领域模型使用 `modality + capabilities`,页面不以摄像头作为唯一根实体;未实现协议适配器明确为 `adapter_not_ready`,不得用模拟遥测伪装交付。 11. Tenant/Site/Area/RBAC、配额、`capture_policy` 与全局审计属于 Bell;Sense Control API v1 只管理 Device 期望态与收敛查询,Sense 只读消费版本化投影并在设备写路径执行,投影不可用时只阻断相关新变更,不静默切断已有链路。 -12. Sense 的设备操作审计先写本地持久化 Outbox,再由幂等 relay 异步送入 Bell 全局审计;不得使用“先执行高风险操作、再尽力入队”的顺序。T-010 已冻结脱敏本地事件并实现原子写入;relay 的 transport、签名、确认、重放窗口与留存仍须独立冻结。 +12. Sense 的设备操作审计先写本地持久化 Outbox,再由幂等 relay 异步送入 Bell 全局审计;不得使用“先执行高风险操作、再尽力入队”的顺序。T-016 已实现 HMAC/nonce 内部 HTTP relay、数据库时钟 lease/fencing、逐项确认与 Bell 不可变全局事实;Sense 不获得 Bell schema 权限,Bell 不读取 Sense Outbox。 ## 6. 容量架构 @@ -88,7 +88,8 @@ T-014 已在单台 Windows 主机上用隔离 PostgreSQL、真实 Control API、 - MediaMTX Path 扫描把“Sense 历史拥有但当前失配”和“从未归属 Sense”分开;未知归属永不自动删除。历史拥有项也只允许在 15 分钟二次快照、1~128 项和 `候选 × 100 <= 当前 Path 总数 × 10` 全部通过时由本地运维命令逐项处置,不提供绕过。 - `bell.site_quota_v1` 行缺失、数值越界、版本回退或读取失败只阻止视频设备新增/启用,不中断已有流;降低配额导致超限时不自动停用,后续准入返回稳定错误并产生运维信号。多 Sense 实例使用 PostgreSQL transaction-scoped advisory lock 串行化同 tenant/site 的计数与写入,不能用进程内锁替代。 - `bell.area_policy_v1` 缺失、非法、版本回退或读取失败时,PostgreSQL repository 拒绝相关新增/启用;`non_imaging_only` 允许非成像设备但拒绝具有 `video_capture` 的设备。已有设备保持原状态,策略冲突由 Bell 管理端显式迁移或取消。同库实时视图不以源记录年龄误判 freshness。 -- 设备创建和期望态受理在本地事务内同时写脱敏 `sense.device_operation_outbox`;Outbox 失败回滚业务写入,相同期望态不增加 generation 但仍审计。异步 relay 尚未实现。 +- 设备创建和期望态受理在本地事务内同时写脱敏 `sense.device_operation_outbox`;Outbox 失败回滚业务写入,相同期望态不增加 generation 但仍审计。relay 最多领取 100 行,以 30 秒数据库 lease 和单调 fencing token 防止过期 worker 确认;成功和 dead letter 都保留本地事实。 +- Bell 先验证时间窗、nonce 和 constant-time HMAC,再逐项校验 v1/v2 事件;同 nonce/同摘要重放原结果,同 nonce/不同摘要拒绝。`bell.audit_events` 只追加且不自动清理,只有 10 分钟幂等收据允许 Bell runtime 删除过期行。 - Bell 最终事件写入 `bell.events`;同平台 ID/同摘要仅视为幂等重放,同 ID/不同摘要拒绝。`bell_runtime` 只有 `SELECT/INSERT`,事件与 outcome 的 UPDATE/DELETE 另由数据库 trigger 拒绝;后续人工/自动 outcome 追加到独立表,不改写事件 payload。 - Brain 投递失败落本地队列重试,不阻塞实时推理主链路。 - Alert 先落库再投递,进程重启恢复未完成升级链。 @@ -98,7 +99,7 @@ T-014 已在单台 Windows 主机上用隔离 PostgreSQL、真实 Control API、 ## 8. 数据与契约 - Bell 核心实体:Tenant → Site → Area(含 `capture_policy`)以及 Role/Binding/Quota/Audit;Sense 核心实体:Device(含 `modality + capabilities`)→ StreamBinding/Zone,以及只记录已观察版本的 SiteQuota/AreaPolicyProjection。两个 schema 以稳定逻辑 ID 关联,不跨 schema 写入;配额 v1 为五列,Area v1 固定为 `tenant_id/site_id/area_id/capture_policy/source_version/source_updated_at` 六列。 -- Sense Control API v1 使用站点作用域路径、认证上下文 tenant、HMAC cursor 分页、PostgreSQL 幂等收据与资源 ETag;敏感连接引用只写不读。T-011 已实现 7 个 handler,并以 feature flag 限定到 PostgreSQL 路径;首版外部静态 SHA-256 注册表只实现认证 port 的私有部署适配器。正式签名和兼容规则以 [`contracts/`](contracts/) 为准,Bell 管理服务、JWT/OIDC 与 Outbox relay 仍未实现。 +- Sense Control API v1 使用站点作用域路径、认证上下文 tenant、HMAC cursor 分页、PostgreSQL 幂等收据与资源 ETag;敏感连接引用只写不读。T-011 已实现 7 个 handler,并以 feature flag 限定到 PostgreSQL 路径;首版外部静态 SHA-256 注册表只实现认证 port 的私有部署适配器。T-016 审计 relay 采用独立外部 HMAC key 文件和内部端点,不等同于 Bell 公共管理认证。正式签名和兼容规则以 [`contracts/`](contracts/) 为准。 - 业务实体:Rule → Event → Alert → DeliveryAttempt/Ack;Event 与 Alert 不合并。 - Bell 通知域分为三个聚合:Contact/Team 保存身份、成员关系和已验证通道;OnCallSchedule/ScheduleVersion/ShiftException 保存时区、轮换与例外;EscalationPolicy/Step 通过 `person / team / on_call_schedule` 类型化 `target_ref` 引用目标。三者共享逻辑 ID,不复制手机号、班次或轮换字段。 - 每个 DeliveryAttempt 创建时解析当时生效的排班版本,并保存实际收件人、通道、`schedule_version` 和解析时间快照;之后联系人或排班修改不得回写既有投递事实。 @@ -114,10 +115,10 @@ Sense/cmd + Sense/internal/{device,onvif,mtx,reconcile,orphan,metrics,probe,trig Brain/{pipeline,models,judge,emit,trigger,contracts} Bell/cmd + Bell/internal/{ingest,event,rule,alert,deliver,feedback,tenant,audit,store} Bell/{web,packs,contracts} -deploy/postgres/{001_roles.sql,...,013_privileges_bell_events.sql,tests} +deploy/postgres/{001_roles.sql,...,015_privileges_audit_relay.sql,tests} ``` -Sense 脚手架和 PostgreSQL `001`~`013` 已实现;Bell 已有事件校验/不可变存储 Go 基础,但没有可部署 API 服务,Brain 仍为目录占位。 +Sense 脚手架和 PostgreSQL `001`~`015` 已实现;Bell 已有事件校验/不可变存储 Go 基础和只面向 Sense 审计 relay 的最小 `bell-api`,但没有公共管理 API,Brain 仍为目录占位。 ## 10. 开发顺序 diff --git a/docs/06-tasks.md b/docs/06-tasks.md index a72bf72..b78884d 100644 --- a/docs/06-tasks.md +++ b/docs/06-tasks.md @@ -37,6 +37,7 @@ - Brain 模型接口、判定内核和 v0.1 mapper。 - Bell 事件校验、不可变存储和 ULID。 - T-015:建立 Bell Go 事件域基础,复制并校验冻结 v0.1 schema,由 Bell 生成平台 ULID,执行六项代码断言,并以 `bell_runtime` 最小权限保存不可变事件和 append-only outcome;不冻结 Brain transport 或公共 API。 +- T-016:冻结并实现 Sense Outbox → Bell 内部审计 relay;使用 HMAC、nonce 收据、数据库时钟 lease/fencing、逐项确认和 dead letter,在不共享 schema 权限的前提下写入 Bell 不可变全局审计事实。 - 规则引擎、场景包加载、预警状态机与双路径投递。 - 最小 Web/App 处置流程、RBAC 与审计。 - 现场误报基线和反馈队列。 diff --git a/docs/api.md b/docs/api.md index 1ec48ac..d570920 100644 --- a/docs/api.md +++ b/docs/api.md @@ -1,6 +1,6 @@ # API 与契约 -> Brain → Bell 事件契约 v0.1、Sense Control API v1、Bell 配额/Area 只读投影 v1 与 Sense 本地设备审计事件 v1/v2 已冻结;其他 API 仍在设计阶段。不得把本文的“待定”自行具体化为公共契约。 +> Brain → Bell 事件契约 v0.1、Sense Control API v1、Bell 配额/Area 只读投影 v1、Sense 本地设备审计事件 v1/v2 与 Sense→Bell 审计 relay v1 已冻结;其他 API 仍在设计阶段。不得把本文的“待定”自行具体化为公共契约。 ## 1. 已冻结:Brain → Bell 事件契约 @@ -25,13 +25,13 @@ T-015 已实现 Bell 消费端的内部组装与存储边界:可信 ingress | --- | --- | --- | --- | | Sense → Bell | 读取站点视频配额 | 同一 PostgreSQL 实例内只读 `bell.site_quota_v1`;默认 16、最大 128;失败时拒绝新增/启用但不影响已有流 | T-008 冻结,T-009 已实现数据路径 | | Sense → Bell | 读取 Area 成像准入 | 只读 `bell.area_policy_v1`;`video_allowed | non_imaging_only`;缺失/非法/回退失败关闭但不影响已有设备 | T-010 已冻结并实现数据路径 | -| Sense → Bell | 汇入设备操作审计 | 本地 Outbox 事件已冻结并原子落库;transport、签名、确认与留存未冻结 | T-010 本地基础已实现,relay 待设计 | +| Sense → Bell | 汇入设备操作审计 | `POST /internal/v1/audit-events:batch`;1~100 项、1 MiB、10 秒 deadline、HMAC/nonce、逐项确认;非回环必须 HTTPS | T-016 已冻结并实现 | | Bell → Sense | 请求事件证据/pre-roll 切片 | 幂等、按租户授权、异步结果、不得暴露原始凭据 | 待 M3 设计 | | Bell → Brain | outcome/误报反馈 | 原事件不可变;反馈可重试、去重、审计 | 待 M3 设计 | | Sense → Brain | 流绑定与设备型触发 | 分片可路由,触发入口与流控制解耦 | 待 M2/M3 设计 | | Worker → 控制面 | 注册、心跳、容量 | `max_sources` 来自 profile/压测,不固定为 16 | 待 M3 设计 | -冻结签名和失败语义见 [`contracts/README.md`](contracts/README.md)、[`contracts/site-quota-v1.sql`](contracts/site-quota-v1.sql)、[`contracts/area-policy-v1.sql`](contracts/area-policy-v1.sql) 与 [`contracts/sense-device-audit-v1.schema.json`](contracts/sense-device-audit-v1.schema.json)。Bell 拥有投影源数据和视图,Sense 数据库角色只有 `SELECT`;未来分库必须发布新版本,不能在 v1 下把本地视图静默替换为网络调用。 +冻结签名和失败语义见 [`contracts/README.md`](contracts/README.md)、[`contracts/sense-audit-relay-v1.openapi.json`](contracts/sense-audit-relay-v1.openapi.json)、[`contracts/site-quota-v1.sql`](contracts/site-quota-v1.sql)、[`contracts/area-policy-v1.sql`](contracts/area-policy-v1.sql) 与 [`contracts/sense-device-audit-v1.schema.json`](contracts/sense-device-audit-v1.schema.json)。Bell 拥有投影源数据和视图,Sense 数据库角色只有 `SELECT`;未来分库必须发布新版本,不能在 v1 下把本地视图静默替换为网络调用。 ## 3. 已冻结:Sense Control API v1 @@ -47,7 +47,7 @@ T-015 已实现 Bell 消费端的内部组装与存储边界:可信 ingress - `endpoint_ref`、`credential_ref`、`profile_token` 只写不读;设备 ID 由服务端生成。普通响应和错误不得包含凭据、完整流 URI、token 或 MediaMTX 内部配置。 - v1 不提供删除设备;停用设备保留历史。写入受理只表示期望态已持久化,不能表示实际态已收敛。 -T-008 冻结公共控制契约;T-009/T-010 建立 PostgreSQL 投影、准入和本地审计基础;T-011 已实现 7 个 HTTP handler、外部静态摘要认证适配器、tenant/Site scope、幂等收据、ETag/HMAC cursor 和持久化 batch operation。业务路由默认关闭且仅可在 PostgreSQL 上开启;Bell 管理服务、JWT/OIDC 和 Outbox relay 尚未实现。 +T-008 冻结公共控制契约;T-009/T-010 建立 PostgreSQL 投影、准入和本地审计基础;T-011 已实现 7 个 HTTP handler、外部静态摘要认证适配器、tenant/Site scope、幂等收据、ETag/HMAC cursor 和持久化 batch operation。业务路由默认关闭且仅可在 PostgreSQL 上开启;T-016 的审计 relay 是独立内部端点,不替代 Bell 管理服务或 JWT/OIDC。 ## 4. 待冻结的 Bell 公共 API diff --git a/docs/contracts/README.md b/docs/contracts/README.md index 669498a..7523b76 100644 --- a/docs/contracts/README.md +++ b/docs/contracts/README.md @@ -1,6 +1,6 @@ # Sense 控制面、准入投影与本地审计契约 v1 -> 冻结日期:2026-08-07。Control API 契约版本:`1.0.0`。Sense 是设备期望态的提供方;Bell 是 Tenant、Site、Area、RBAC、配额与全局审计的所有者。T-011 已实现 Control API handler 与 PostgreSQL 一致性边界;Bell 管理服务、JWT/OIDC 和 Outbox relay 仍未实现。 +> 冻结日期:2026-08-11。Control API 与审计 relay 契约版本:`1.0.0`。Sense 是设备期望态的提供方;Bell 是 Tenant、Site、Area、RBAC、配额与全局审计的所有者。T-016 已实现 Sense Outbox 到 Bell 的内部 relay;Bell 公共管理服务与 JWT/OIDC 仍未实现。 ## 契约文件 @@ -9,8 +9,9 @@ | [`sense-control-v1.openapi.json`](sense-control-v1.openapi.json) | Sense | Bell 管理面、受控集成方 | 设备查询、创建、修改、启停与批量操作 | | [`site-quota-v1.sql`](site-quota-v1.sql) | Bell | Sense | 单 PostgreSQL 实例内的站点视频配额只读投影 | | [`area-policy-v1.sql`](area-policy-v1.sql) | Bell | Sense | Area 归属与 `capture_policy` 只读投影 | -| [`sense-device-audit-v1.schema.json`](sense-device-audit-v1.schema.json) | Sense | 本地 Outbox;未来 Bell relay | 脱敏设备操作审计事实,不包含传输协议 | -| [`sense-device-audit-v2.schema.json`](sense-device-audit-v2.schema.json) | Sense | 本地 Outbox;未来 Bell relay | v1 后继,增加脱敏配置修改受理事实;v1 文件保持不变 | +| [`sense-device-audit-v1.schema.json`](sense-device-audit-v1.schema.json) | Sense | 本地 Outbox;Bell relay | 脱敏设备操作审计事实,不包含传输协议 | +| [`sense-device-audit-v2.schema.json`](sense-device-audit-v2.schema.json) | Sense | 本地 Outbox;Bell relay | v1 后继,增加脱敏配置修改受理事实;v1 文件保持不变 | +| [`sense-audit-relay-v1.openapi.json`](sense-audit-relay-v1.openapi.json) | Bell | Sense | 内部批量端点、HMAC、逐项确认、nonce 防重与重试边界 | OpenAPI 的 `/api/v1` 路径是公共控制面边界;`/healthz`、`/readyz` 仍是非业务运维探针。v1 不提供设备删除:停用设备使用期望态接口,保留设备、操作和审计历史。Site、Area、配额、RBAC 和审计聚合不由 Sense 提供 CRUD。 @@ -57,7 +58,19 @@ T-010 冻结 `bell.area_policy_v1` 的列顺序为 `tenant_id/site_id/area_id/ca `sense-device-audit-v1.schema.json` 继续冻结创建与期望态两类事实且不原地扩展严格枚举。T-011 新增 v2 后继,兼容 v1 两类事件并增加 `device.configuration.accepted`;该 payload 只保存是否变化、字段名和 Area 逻辑 ID,不保存字段值。主体类型为 `user | service | system`,投影版本与 generation 随事实保存;endpoint、credential、profile token、path、密码、完整流 URI 或 MediaMTX 配置始终禁止进入审计。 -PostgreSQL repository 必须在设备创建/期望态事务内写 `sense.device_operation_outbox`;Outbox 失败回滚业务写入。相同期望态不增加 generation,但仍产生独立审计事实。Schema 不是 Bell relay 协议:传输端点、签名、批量确认、重放窗口和留存由后续任务冻结。 +PostgreSQL repository 必须在设备创建/期望态事务内写 `sense.device_operation_outbox`;Outbox 失败回滚业务写入。相同期望态不增加 generation,但仍产生独立审计事实。事件 schema 继续只定义事实;传输由 `sense-audit-relay-v1.openapi.json` 独立冻结。 + +## Sense → Bell 审计 relay + +Sense 向 `/internal/v1/audit-events:batch` 每批发送 1~100 个事件,请求体不超过 1 MiB、deadline 10 秒。请求用外部文件中的至少 32 字节 secret 做 HMAC-SHA256,canonical string 为 method、path、Unix 秒、随机 nonce 与 body SHA-256 的换行拼接;非回环地址必须使用 HTTPS。Bell 允许 300 秒时钟偏差并将 `(key_id, nonce)` 收据保留 600 秒:相同请求摘要返回原结果,不同摘要返回 `409 replay_conflict`。 + +Outbox 用 30 秒数据库时钟 lease、单调 fencing token 和 `FOR UPDATE SKIP LOCKED` 协调实例。`accepted/duplicate` 才标记 delivered;逐项 `rejected` 进入 dead letter;网络、5xx、认证失败或缺失结果以 1 秒起步、最多 300 秒指数退避。Sense 不直接访问 Bell schema,Bell 不读取 Sense Outbox;全局 `bell.audit_events` 不自动删除,只有短期 relay receipt 自动过期。 + +两端读取同格式的仓库外 key 文件;Bell 可同时接受多个 key,Sense 用 `SENSE_AUDIT_RELAY_KEY_ID` 选择一个,便于先加新 key、切换发送端、再移除旧 key。占位结构如下,`secret_base64url` 必须替换为至少 32 个随机字节的无填充 base64url,不能提交真实值: + +```json +{"version":1,"keys":[{"key_id":"sense-a","secret_base64url":""}]} +``` ## 兼容与废弃 @@ -73,8 +86,10 @@ PostgreSQL repository 必须在设备创建/期望态事务内写 `sense.device_ ```powershell python -m json.tool docs/contracts/sense-control-v1.openapi.json | Out-Null python -m json.tool docs/contracts/sense-device-audit-v2.schema.json | Out-Null +python -m json.tool docs/contracts/sense-audit-relay-v1.openapi.json | Out-Null python -m unittest discover -s tests -p "test_sense_control_contract.py" python -m unittest discover -s tests -p "test_sense_control_implementation.py" +python -m unittest discover -s tests -p "test_sense_audit_relay_contract.py" ``` 测试同时校验 OpenAPI 结构、生成 server glue、HTTP handler 与 PostgreSQL migration/事务;它不替代 Bell 消费方联合验收或客户现场容量验证。 diff --git a/docs/contracts/sense-audit-relay-v1.openapi.json b/docs/contracts/sense-audit-relay-v1.openapi.json new file mode 100644 index 0000000..3dfc0f4 --- /dev/null +++ b/docs/contracts/sense-audit-relay-v1.openapi.json @@ -0,0 +1,94 @@ +{ + "openapi": "3.1.0", + "info": { + "title": "YoVision Sense Audit Relay", + "version": "1.0.0", + "description": "Internal, signed and idempotent delivery of redacted Sense device audit facts to Bell." + }, + "paths": { + "/internal/v1/audit-events:batch": { + "post": { + "operationId": "receiveSenseAuditBatch", + "description": "Accepts 1-100 events in a body no larger than 1048576 bytes. The request deadline is 10 seconds. HMAC clock skew is at most 300 seconds and nonce receipts live for 600 seconds.", + "parameters": [ + {"name": "X-YoVision-Key-Id", "in": "header", "required": true, "schema": {"type": "string", "pattern": "^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$"}}, + {"name": "X-YoVision-Timestamp", "in": "header", "required": true, "description": "Unix seconds", "schema": {"type": "string", "pattern": "^[0-9]{10,}$"}}, + {"name": "X-YoVision-Nonce", "in": "header", "required": true, "description": "16-48 random bytes encoded as unpadded base64url", "schema": {"type": "string", "minLength": 22, "maxLength": 64, "pattern": "^[A-Za-z0-9_-]+$"}}, + {"name": "X-YoVision-Signature", "in": "header", "required": true, "description": "Unpadded base64url HMAC-SHA256 over POST, path, timestamp, nonce and lowercase SHA-256 body digest joined by newlines", "schema": {"type": "string", "minLength": 43, "maxLength": 43}} + ], + "requestBody": { + "required": true, + "content": {"application/json": {"schema": {"$ref": "#/components/schemas/BatchRequest"}}} + }, + "responses": { + "200": {"description": "Stored, duplicate or permanently rejected per item. Same key ID, nonce and request digest returns the original response.", "content": {"application/json": {"schema": {"$ref": "#/components/schemas/BatchResponse"}}}}, + "400": {"description": "Malformed batch envelope"}, + "401": {"description": "Invalid key, timestamp, nonce or signature"}, + "409": {"description": "Same key ID and nonce used with a different request digest; error is replay_conflict"}, + "413": {"description": "Body exceeds 1048576 bytes"}, + "503": {"description": "Bell cannot atomically persist the batch and receipt"} + } + } + } + }, + "components": { + "schemas": { + "BatchRequest": { + "type": "object", "additionalProperties": false, "required": ["events"], + "properties": {"events": {"type": "array", "minItems": 1, "maxItems": 100, "items": {"$ref": "#/components/schemas/Envelope"}}} + }, + "Envelope": { + "type": "object", "additionalProperties": false, "required": ["schema_version", "event"], + "properties": { + "schema_version": {"type": "integer", "enum": [1, 2]}, + "event": {"$ref": "#/components/schemas/AuditEvent"} + } + }, + "AuditEvent": { + "type": "object", "additionalProperties": false, + "required": ["event_id", "event_type", "tenant_id", "site_id", "device_id", "actor", "reason", "trace_id", "aggregate_generation", "projection_versions", "data", "occurred_at"], + "properties": { + "event_id": {"type": "string", "pattern": "^audit_[0-9a-f]{32}$"}, + "event_type": {"type": "string", "enum": ["device.created", "device.desired_state.accepted", "device.configuration.accepted"]}, + "tenant_id": {"$ref": "#/components/schemas/LogicalId"}, + "site_id": {"$ref": "#/components/schemas/LogicalId"}, + "device_id": {"$ref": "#/components/schemas/LogicalId"}, + "actor": {"type": "object", "additionalProperties": false, "required": ["type", "id"], "properties": {"type": {"type": "string", "enum": ["user", "service", "system"]}, "id": {"type": "string", "minLength": 1, "maxLength": 200}}}, + "reason": {"type": ["string", "null"], "maxLength": 500}, + "trace_id": {"type": ["string", "null"], "maxLength": 128}, + "aggregate_generation": {"type": "integer", "minimum": 1}, + "projection_versions": {"type": "object", "additionalProperties": false, "required": ["quota_source_version", "area_policy_source_version"], "properties": {"quota_source_version": {"type": ["integer", "null"], "minimum": 1}, "area_policy_source_version": {"type": ["integer", "null"], "minimum": 1}}}, + "data": {"type": "object"}, + "occurred_at": {"type": "string", "format": "date-time"} + } + }, + "LogicalId": {"type": "string", "minLength": 1, "maxLength": 128, "pattern": "^[A-Za-z0-9][A-Za-z0-9._:-]*$"}, + "BatchResponse": { + "type": "object", "additionalProperties": false, "required": ["results"], + "properties": {"results": {"type": "array", "minItems": 1, "maxItems": 100, "items": {"$ref": "#/components/schemas/Result"}}} + }, + "Result": { + "type": "object", "additionalProperties": false, "required": ["event_id", "status"], + "properties": { + "event_id": {"type": "string"}, + "status": {"type": "string", "enum": ["accepted", "duplicate", "rejected"]}, + "error_code": {"type": "string", "pattern": "^[a-z][a-z0-9_]{0,63}$"} + } + } + } + }, + "x-yovision-signature": { + "canonical": "METHOD\\nPATH\\nTIMESTAMP\\nNONCE\\nLOWERCASE_SHA256_BODY", + "algorithm": "HMAC-SHA256", + "encoding": "base64url-no-padding", + "clock_skew_seconds": 300, + "receipt_ttl_seconds": 600 + }, + "x-yovision-delivery": { + "lease_seconds": 30, + "initial_retry_seconds": 1, + "maximum_retry_seconds": 300, + "successful_statuses": ["accepted", "duplicate"], + "permanent_status": "rejected" + } +} diff --git a/docs/current-state.md b/docs/current-state.md index 0fb3fd5..71de7de 100644 --- a/docs/current-state.md +++ b/docs/current-state.md @@ -1,28 +1,29 @@ # 当前实现状态 -> 快照日期:2026-08-10。只记录仓库现实与 blocker;任务实时状态到 Gitea Issue 查看。 +> 快照日期:2026-08-11。只记录仓库现实与 blocker;任务实时状态到 Gitea Issue 查看。 ## 当前阶段 -- 阶段:M0 指定摄像头型号准入、M1“一实机 + 四合成源”软件闭环和 M2 本地 16 路批量收敛/稳定基线已通过;M3 已建立 Bell 不可变事件存储基础。客户网络尚未提供,WireGuard T-013 后置,五条独立真实上游和生产 SLA 仍未验收。下一开发重点是 Sense Outbox → Bell 全局审计 relay。 -- 生产代码:Sense 已包含可构建进程、SQLite/PostgreSQL repository、Site/Area 准入、设备操作 Outbox、标准 ONVIF SOAP/WS-Security adapter、凭据引用、MediaMTX 生成客户端、Control API v1、对账/探活、数据库租约、孤儿只读扫描/受控命令、低基数指标和可重复 16 路容量脚本;Bell 已包含可构建 Go module、事件 v0.1 schema/语义校验、平台 ULID、不可变 PostgreSQL repository 和 append-only outcome,但仍没有可部署 API、管理服务/JWT、Outbox relay、规则/Alert 或 Web/H5。 +- 阶段:M0 指定摄像头型号准入、M1“一实机 + 四合成源”软件闭环和 M2 本地 16 路批量收敛/稳定基线已通过;M3 已建立 Bell 不可变事件存储及 Sense→Bell 全局审计 relay 基础。客户网络尚未提供,WireGuard T-013 后置,五条独立真实上游和生产 SLA 仍未验收。 +- 生产代码:Sense 已包含可构建进程、SQLite/PostgreSQL repository、Site/Area 准入、设备操作 Outbox、可选签名 relay、标准 ONVIF SOAP/WS-Security adapter、凭据引用、MediaMTX 生成客户端、Control API v1、对账/探活、数据库租约、孤儿只读扫描/受控命令、低基数指标和可重复 16 路容量脚本;Bell 已包含事件 v0.1 校验/不可变存储、append-only outcome,以及只服务 Sense 审计的最小 `bell-api` 和全局审计 repository,但仍没有 Brain 事件 ingress、公共管理服务/JWT、规则/Alert 或 Web/H5。 - 默认容量:16 路;单站点本阶段上限 128 路,必须横向分片。 ## 仓库现实 -- `Sense/` 已有 Go module 与 `cmd/sense-api`;`Bell/` 已有事件域 Go module 但没有 `cmd/bell-api`;`Brain/` 仍只有目录占位。Bell/Sense migration 统一位于根目录 `deploy/postgres/`。 +- `Sense/` 已有 Go module 与 `cmd/sense-api`;`Bell/` 已有事件域 Go module 和最小 `cmd/bell-api` 内部审计 receiver;`Brain/` 仍只有目录占位。Bell/Sense migration 统一位于根目录 `deploy/postgres/`。 - Sense 设备模型使用 `modality + capabilities`,SQLite 执行 v1 migration;视频配额默认 16、允许 1~128,17/128/129、新增/启用和“降低配额不关闭已有流”均有测试。 - T-009 冻结 PostgreSQL `17.10` 和 `pgx/v5 v5.10.0`,实现 `bell`/`sense` schema、NOLOGIN 权限角色、Bell Site 版本 trigger、`bell.site_quota_v1` 和 Sense PostgreSQL repository;同站点并发准入用事务级 advisory lock,配额缺失/越界/版本回退时失败关闭且不改变已有流。 - T-010 增量实现 `bell.areas`、`bell.area_policy_v1`、Area 版本观察和 `sense.device_operation_outbox`;`non_imaging_only` 拒绝成像设备创建/启用,失败不改变已有设备。设备创建/期望态受理与脱敏 Outbox 同事务,相同期望态不增加 generation 但仍审计。 -- Windows 隔离测试使用 `D:\pgsql17\bin` 启动随机回环端口临时集群,`001`~`013` migration 可重放;Sense repository、Bell event repository、权限、幂等冲突与不可变性测试通过后自动清理,现有 `D:\pgsql17\data` 和 5432 服务未被读取、停止或修改。 +- Windows 隔离测试使用 `D:\pgsql17\bin` 启动随机回环端口临时集群,`001`~`015` migration 可重放;Sense Outbox fencing、Bell event/global-audit repository、nonce 收据、权限、幂等冲突与不可变性测试通过后自动清理,现有 `D:\pgsql17\data` 和 5432 服务未被读取、停止或修改。 - T-015 冻结 Bell Go 1.26.5、JSON Schema v6.0.2 和 ULID v2.1.2;Bell 拒绝上游自报平台 ID,在内部 candidate 组装后执行冻结 v0.1 schema 与六项语义断言。`bell_runtime` 只允许追加/读取 `bell.events`、`bell.event_outcomes`;Brain transport、整数事件 ID 与现有文本逻辑 ID 的跨系统映射、公共 API 和生产隐私 resolver 仍未冻结,不能把内部 factory 当成已上线入口。 +- T-016 冻结 `sense-audit-relay-v1`:每批 1~100 项、1 MiB、10 秒 deadline、300 秒时钟窗、600 秒 nonce 收据、30 秒数据库 lease、1~300 秒退避。Sense 使用 `FOR UPDATE SKIP LOCKED` 和 fencing token;Bell constant-time 校验 HMAC,逐项返回 accepted/duplicate/rejected,并把全局事实追加到不可变 `bell.audit_events`。relay 默认关闭,非回环两端必须 HTTPS/TLS,key 只从仓库外文件读取。 - MediaMTX 固定为独立二进制 `v1.19.3`,官方 OpenAPI 已按 SHA-256 vendoring,并由固定 `oapi-codegen v2.8.0` 生成客户端;手写薄封装有 create/read/delete、幂等 ensure、探活和只返回名称的受限分页枚举测试。 - T-003 对账进度与指数退避持久化,覆盖取消和 SQLite 重启恢复;T-006 增加真实 ONVIF adapter、RTSP router、实验室播种/状态工具、故障代理和五路自动验收。T-012 的普通调和不枚举孤儿;独立 PostgreSQL 扫描默认只报告,未知归属永不删除。 - T-006 正式使用 1 台准入实机和 4 个独立合成 publisher 连续观察 `1806.6 s` / 180 次采样,四类恢复均通过,最大与最终 `unconverged` 均为 0;详细证据见 `docs/research/sense-5-stream-integration.md`。 - T-014 正式使用隔离 PostgreSQL、真实 Control API、两套 MediaMTX 和 16 个独立低码率合成 publisher,完成 17 路配额拒绝、三轮 `16 → 0 → 16` 批量收敛和固定四路故障恢复;稳定观察 `1800.1 s` / 180 次采样,最大与最终 `unconverged` 均为 0、最终在线 Path 16、帧错误 0。证据见 `docs/research/sense-16-stream-capacity.md`;不外推到真实 16 机、网络、录像、AI/GPU、64/128 路或生产 SLA。 - `docs/raw/01`~`08` 已记录需求、分析、方案、客户场景、事件比对和三系统职责。 - `docs/raw/contracts/event-v0.1.schema.json` 已冻结,并有多份示例与语义说明。 -- `docs/contracts/sense-control-v1.openapi.json` 的 7 个站点作用域/operation endpoint 已由 T-011 实现:外部静态 SHA-256 主体注册表、tenant/Site scope、HMAC cursor、ETag、PostgreSQL 24 小时幂等收据和最多 128 项 batch operation 均有代码与隔离集成测试。默认 SQLite 只暴露运维探针与低基数 `/metrics`,不注册业务路由;Bell 管理服务、JWT/OIDC 和 Outbox relay 尚未实现。 +- `docs/contracts/sense-control-v1.openapi.json` 的 7 个站点作用域/operation endpoint 已由 T-011 实现;`sense-audit-relay-v1.openapi.json` 已由 T-016 实现。默认 SQLite 只暴露运维探针与低基数 `/metrics`,不注册业务路由或 relay;Bell 管理服务与 JWT/OIDC 尚未实现。 - T-012 把 PostgreSQL schema 提升到 v5:due row 用数据库时钟、`FOR UPDATE SKIP LOCKED`、逐项续租和 fencing token 协调;MediaMTX Path 历史归属、15 分钟孤儿快照、最多 10%/128 项安全闸、无 bypass 的本地处置命令及 `/metrics` 已实现。SQLite 明确保留单实例开发语义。 - harness coding 文档、上下文清单、Gitea Issue/PR 模板和治理脚本已接入。 - Gitea 已初始化 12 个协作标签;`status/waiting` 用于依赖或外部条件未满足的未领取任务,实时可领取状态必须从 Gitea 查询,不在本文复制。 @@ -44,7 +45,7 @@ Windows: go -C Sense run ./cmd/sense-api ``` -Bell 事件域基础验证(当前没有可启动 API): +Bell 事件域和内部审计 receiver 验证: ```powershell go -C Bell test ./... @@ -86,12 +87,12 @@ Sense 默认监听 `127.0.0.1:8080`,提供 `/healthz`、`/readyz` 运维探针 - 人脸方向已延后至 M5 的 S4 成人园区候选试点;必要性/PIP 影响评估、单独同意与替代方式、合法底库来源和删除流程未完成,阻塞人脸能力上线。 - 短信/语音具体供应商未选;生产前必须选定两条独立投递路径并验证故障切换。 - Python/Savant 的精确版本、目标硬件和 Bell 前端栈尚未冻结;Sense M1 的 Go、SQLite driver、MediaMTX、生成器及生成运行时版本已在 T-003 冻结,PostgreSQL/pgx 版本已在 T-009 冻结。 -- 本机现有 PostgreSQL 5432 实例使用 SCRAM 且当前开发进程没有管理员密码;T-009~T-015 不绕过认证,自动验收使用隔离临时集群。向共享/生产实例安装 migration 前仍需管理员私下提供专用数据库、最小权限登录角色、外部 Control API 安全文件与备份方案。 +- 本机现有 PostgreSQL 5432 实例使用 SCRAM 且当前开发进程没有管理员密码;T-009~T-016 不绕过认证,自动验收使用隔离临时集群。向共享/生产实例安装 migration 前仍需管理员私下提供专用数据库、最小权限登录角色、外部 Control API/relay key 文件、TLS 证书与备份方案。 - 代码知识图谱在无业务代码阶段可能为空;工具不可用时使用 `rg` 处理文档与配置。 ## 下一步 -客户网络仍未提供,T-013 WireGuard 继续后置。Bell 不可变事件存储已由 T-015 落地;下一项是独立冻结并实现 Sense 设备操作 Outbox 到 Bell 全局审计的幂等 relay,不把 Brain 事件 transport、规则或 Alert 混入。客户授权、借用或租赁条件具备后再执行 T-007 五条独立真实上游现场门禁。T-014/T-015 不解除 T-007/T-013,也不形成真实多路或生产 SLA 承诺。 +客户网络仍未提供,T-013 WireGuard 继续后置。T-015/T-016 已分别落地 Bell 不可变事件存储和 Sense 全局审计 relay;下一项应独立冻结 Brain→Bell 业务事件 ingress,不把规则/Alert 或公共管理 API 混入。客户授权、借用或租赁条件具备后再执行 T-007 五条独立真实上游现场门禁。T-014~T-016 不解除 T-007/T-013,也不形成真实多路或生产 SLA 承诺。 ## 已知风险 diff --git a/docs/tasks/T-016.md b/docs/tasks/T-016.md index 39d796a..03c173d 100644 --- a/docs/tasks/T-016.md +++ b/docs/tasks/T-016.md @@ -3,12 +3,12 @@ id: T-016 title: 冻结并实现 Sense Outbox 到 Bell 全局审计幂等 relay phase: 3 deps: [T-015] -status: TODO +status: DONE created: 2026-08-10 issue: 55 -context_ref: null -claim_branch: null -work_branch: null +context_ref: aca22f4667f0da77a102eb7ba15a93c0128df656 +claim_branch: claims/T-016 +work_branch: agent/codex/T-016 write_paths: - docs/tasks/T-016.md - docs/contracts/README.md @@ -94,3 +94,7 @@ T-015 已提供 Bell 的独立运行角色和不可变事实模式。本任务 ## 执行记录 - 2026-08-10:在 T-015 合并并关闭后拆出本任务;实现尚未开始。 +- 2026-08-11:冻结 `sense-audit-relay-v1.openapi.json`,实现 Sense 外部 key/HMAC client、30 秒数据库 lease/fencing worker、逐项结果/dead letter/1~300 秒退避,以及 PostgreSQL v6 Outbox relay repository;默认关闭且远端 URL 强制 HTTPS。 +- 2026-08-11:实现 Bell 最小 `cmd/bell-api`、300 秒时间窗与 constant-time HMAC 校验、600 秒 nonce 幂等收据、逐项 v1/v2 校验,以及 PostgreSQL v4 不可变 `bell.audit_events`;Sense/Bell 登录权限保持单向隔离,只有 Bell 可清理过期收据。 +- 2026-08-11:`./init.ps1`、三条 Python 治理命令、Sense/Bell `test/vet/build` 与 `git diff --check` 全部通过;Python 共 68 项测试通过。 +- 2026-08-11:`./scripts/test_postgres.ps1 -PgRoot D:\pgsql17` 通过;PostgreSQL 17.10 临时集群将 `001`~`015` 重放两次,Sense fencing、Bell receipt/重复/冲突/不可变性和双方权限断言通过,随机端口与临时目录已清理,现有 5432 listener 未改变。 diff --git a/scripts/test_postgres.ps1 b/scripts/test_postgres.ps1 index 610f8a6..cd12450 100644 --- a/scripts/test_postgres.ps1 +++ b/scripts/test_postgres.ps1 @@ -91,7 +91,9 @@ try { '010_reconcile_safety.sql', '011_privileges_reconcile_safety.sql', '012_bell_events.sql', - '013_privileges_bell_events.sql' + '013_privileges_bell_events.sql', + '014_audit_relay.sql', + '015_privileges_audit_relay.sql' )) { Invoke-Checked $psql '-X' '-v' 'ON_ERROR_STOP=1' '-d' $adminDatabaseDSN '-f' (Join-Path $repoRoot "deploy\postgres\$name") } diff --git a/tests/test_postgres_contract.py b/tests/test_postgres_contract.py index b4743a1..1b7de4a 100644 --- a/tests/test_postgres_contract.py +++ b/tests/test_postgres_contract.py @@ -60,6 +60,8 @@ class PostgresContractTests(unittest.TestCase): "011_privileges_reconcile_safety.sql", "012_bell_events.sql", "013_privileges_bell_events.sql", + "014_audit_relay.sql", + "015_privileges_audit_relay.sql", ], names, ) @@ -129,6 +131,21 @@ class PostgresContractTests(unittest.TestCase): self.assertIn(marker, text) self.assertNotIn("D:\\pgsql17\\data", text) + def test_audit_relay_uses_fencing_and_separate_schema_ownership(self) -> None: + migration = normalized(migration_text("014_audit_relay.sql")) + privileges = normalized(migration_text("015_privileges_audit_relay.sql")) + for marker in ( + "relay_lease_owner", + "relay_lease_token", + "relay_lease_until", + "create table if not exists bell.audit_events", + "create table if not exists bell.audit_relay_receipts", + "expires_at >= received_at + interval '10 minutes'", + ): + self.assertIn(marker, migration) + self.assertIn("grant select, insert on table bell.audit_events to bell_runtime", privileges) + self.assertIn("revoke all on table bell.audit_events, bell.audit_relay_receipts from sense_app", privileges) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_sense_audit_relay_contract.py b/tests/test_sense_audit_relay_contract.py new file mode 100644 index 0000000..e849fae --- /dev/null +++ b/tests/test_sense_audit_relay_contract.py @@ -0,0 +1,54 @@ +"""Cross-language invariants for the Sense-to-Bell audit relay contract.""" + +from __future__ import annotations + +import json +import unittest +from pathlib import Path + + +ROOT = Path(__file__).resolve().parents[1] +CONTRACT = ROOT / "docs" / "contracts" / "sense-audit-relay-v1.openapi.json" + + +class SenseAuditRelayContractTests(unittest.TestCase): + @classmethod + def setUpClass(cls) -> None: + cls.document = json.loads(CONTRACT.read_text(encoding="utf-8")) + + def test_transport_and_response_are_frozen(self) -> None: + operation = self.document["paths"]["/internal/v1/audit-events:batch"]["post"] + self.assertEqual({"200", "400", "401", "409", "413", "503"}, set(operation["responses"])) + headers = {item["name"] for item in operation["parameters"]} + self.assertEqual( + {"X-YoVision-Key-Id", "X-YoVision-Timestamp", "X-YoVision-Nonce", "X-YoVision-Signature"}, + headers, + ) + + def test_numeric_safety_bounds_match_task(self) -> None: + events = self.document["components"]["schemas"]["BatchRequest"]["properties"]["events"] + self.assertEqual((1, 100), (events["minItems"], events["maxItems"])) + signature = self.document["x-yovision-signature"] + delivery = self.document["x-yovision-delivery"] + self.assertEqual(300, signature["clock_skew_seconds"]) + self.assertEqual(600, signature["receipt_ttl_seconds"]) + self.assertEqual((30, 1, 300), (delivery["lease_seconds"], delivery["initial_retry_seconds"], delivery["maximum_retry_seconds"])) + + def test_contract_and_implementations_share_literals(self) -> None: + sense = (ROOT / "Sense" / "internal" / "auditrelay" / "client.go").read_text(encoding="utf-8") + bell = (ROOT / "Bell" / "internal" / "audit" / "audit.go").read_text(encoding="utf-8") + for marker in ( + "/internal/v1/audit-events:batch", + "X-YoVision-Key-Id", + "X-YoVision-Timestamp", + "X-YoVision-Nonce", + "X-YoVision-Signature", + ): + self.assertIn(marker, sense) + self.assertIn(marker, bell) + self.assertIn("10 * time.Second", sense) + self.assertIn("300*time.Second", (ROOT / "Sense" / "internal" / "auditrelay" / "worker.go").read_text(encoding="utf-8")) + + +if __name__ == "__main__": + unittest.main() -- 2.55.0