Bind original recording uploads and confirm one final result
This commit is contained in:
@@ -60,8 +60,9 @@
|
||||
|
||||
- 已新增 Agent 内存录音单次 PUT 组件 `UploadClient.UploadBytes`:复用受限授权、大小和 SHA-256 核对,成功路径不写临时录音或通话结果文件。OSS 明确拒绝保留状态码且不自动重试;传输结果不明时返回专用错误、停止自动重试并隐藏带签名的 URL。空授权头拒绝而非使服务崩溃。这里只验证组件,尚未连接实际录音或最终结果。
|
||||
- Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;真实 Dispatcher 授权及回报尚未接线。
|
||||
- Dispatcher 的无录音最终结果隔离组件:`CurrentStore.RecordCallResult` 仅在确认通话结束后,按持久任务快照校验任务、被叫、主叫和已选线路,并以源执行事件固定生成唯一最终结果身份;消息通过严格 MQ Schema 校验后与 outbox 在同一事务写入。同内容重投/重启只恢复原消息,冲突结果、未授权的上传资产及 SQLite 写入失败均不会产生第二份结果。这里只验证无录音状态;Agent 实际回报和已上传对象的授权核验仍未连通。
|
||||
- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent -count=1`、`go test -race ./internal/store -count=1`、`go vet ./...`、`go build ./...`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`git diff --check`。尚未完成实际录音到直传/失败恢复的接线、无录音与生成失败的真实调用、已上传对象的 Dispatcher 授权绑定核验、真实 Agent↔Dispatcher 结果交付及 MQ/端到端验收,不能宣称 P06 通过。
|
||||
- Dispatcher 的无录音最终结果隔离组件:`CurrentStore.RecordCallResult` 仅在确认通话结束后,按持久任务快照校验任务、被叫、主叫和已选线路,并以源执行事件固定生成唯一最终结果身份;消息通过严格 MQ Schema 校验后与 outbox 在同一事务写入。同内容重投/重启只恢复原消息,冲突结果和 SQLite 写入失败均不会产生第二份结果。这里只验证隔离组件,Agent 实际回报尚未连通。
|
||||
- Dispatcher 的原始 OSS 目标及已上传结果隔离组件:新 SQLite 布局把一次通话的 upload_id、bucket、object_key、录音格式/时长/大小和 SHA-256 唯一绑定到已保留的执行;不保存临时 URL 或 TOKEN。旧布局拒绝启动并原样保留待交付 outbox,不自动迁移或清理。已签发录音目标不能通过空录音结果绕过上传;已有空录音最终结果不能再签发录音目标。`RecordUploadedCallResult` 仅接受与持久绑定完全一致的录音事实及 Agent 所报告的成功 PUT 状态,录音确认与唯一最终结果 outbox 同事务提交;丢失回报或 MQ 投递时重用原消息,已确认后拒绝再次签发 PUT 授权。Mock 证明的是本地状态约束,不是独立 OSS 校验或真实 Agent 身份验证。
|
||||
- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent -count=1`、`go test -race ./internal/store -count=1`、`go vet ./...`、`go build ./...`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-current-contracts.sh`、`PATH=/tmp/sip-go-agent-tools/bin:$PATH bash scripts/check-proto.sh`、`git diff --check`。尚未完成实际录音到直传/失败恢复的接线、无录音与生成失败的真实调用、Dispatcher 实际签发授权与 Agent↔Dispatcher 上传事实/最终结果交付及 MQ/端到端验收,不能宣称 P06 通过。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -71,6 +71,25 @@ CREATE TABLE IF NOT EXISTS dispatcher_outbox (
|
||||
confirmed INTEGER NOT NULL DEFAULT 0 CHECK(confirmed IN (0,1)),
|
||||
confirmed_at TEXT,
|
||||
PRIMARY KEY(dispatcher_id,event_id)
|
||||
);
|
||||
CREATE TABLE IF NOT EXISTS dispatcher_recordings (
|
||||
dispatcher_id TEXT NOT NULL,
|
||||
source_event_id TEXT NOT NULL,
|
||||
upload_id TEXT NOT NULL,
|
||||
recording_id TEXT NOT NULL,
|
||||
bucket TEXT NOT NULL,
|
||||
object_key TEXT NOT NULL,
|
||||
checksum_sha256 TEXT NOT NULL,
|
||||
size_bytes INTEGER NOT NULL CHECK(size_bytes > 0),
|
||||
format TEXT NOT NULL,
|
||||
channels INTEGER NOT NULL CHECK(channels > 0),
|
||||
sample_rate_hz INTEGER NOT NULL CHECK(sample_rate_hz > 0),
|
||||
duration_ms INTEGER NOT NULL CHECK(duration_ms >= 0),
|
||||
confirmed_at TEXT,
|
||||
PRIMARY KEY(dispatcher_id,source_event_id),
|
||||
UNIQUE(dispatcher_id,upload_id),
|
||||
UNIQUE(bucket,object_key),
|
||||
FOREIGN KEY(dispatcher_id,source_event_id) REFERENCES dispatcher_inbox(dispatcher_id,event_id)
|
||||
);`
|
||||
|
||||
func OpenCurrent(path string) (_ *CurrentStore, err error) {
|
||||
@@ -91,7 +110,7 @@ func OpenCurrent(path string) (_ *CurrentStore, err error) {
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("inspect SQLite schema before writing: %w", err)
|
||||
}
|
||||
allowed := map[string]bool{"dispatcher_state": true, "dispatcher_tasks": true, "dispatcher_configs": true, "dispatcher_inbox": true, "dispatcher_outbox": true}
|
||||
allowed := map[string]bool{"dispatcher_state": true, "dispatcher_tasks": true, "dispatcher_configs": true, "dispatcher_inbox": true, "dispatcher_outbox": true, "dispatcher_recordings": true}
|
||||
seen := make(map[string]bool, len(allowed))
|
||||
var unexpected []string
|
||||
for rows.Next() {
|
||||
@@ -123,7 +142,7 @@ func OpenCurrent(path string) (_ *CurrentStore, err error) {
|
||||
return nil, fmt.Errorf("unknown SQLite layout version %d without tables: preserve database before admission", version)
|
||||
}
|
||||
} else {
|
||||
if len(seen) != len(allowed) || version != 1 {
|
||||
if len(seen) != len(allowed) || version != 2 {
|
||||
return nil, fmt.Errorf("existing SQLite layout is incomplete or obsolete (tables=%d, version=%d): preserve data before admission", len(seen), version)
|
||||
}
|
||||
var requiredColumns int
|
||||
@@ -133,19 +152,36 @@ func OpenCurrent(path string) (_ *CurrentStore, err error) {
|
||||
if requiredColumns != 3 {
|
||||
return nil, errors.New("existing SQLite lacks durable SIP admission columns: preserve data before admission")
|
||||
}
|
||||
var recordingColumns int
|
||||
if err := db.QueryRow(`SELECT COUNT(*) FROM pragma_table_info('dispatcher_recordings') WHERE name IN ('dispatcher_id','source_event_id','upload_id','recording_id','bucket','object_key','checksum_sha256','size_bytes','format','channels','sample_rate_hz','duration_ms','confirmed_at')`).Scan(&recordingColumns); err != nil {
|
||||
return nil, fmt.Errorf("inspect original recording target layout: %w", err)
|
||||
}
|
||||
if recordingColumns != 13 {
|
||||
return nil, errors.New("existing SQLite lacks immutable recording target columns: preserve data before admission")
|
||||
}
|
||||
}
|
||||
for _, pragma := range []string{"PRAGMA busy_timeout=5000", "PRAGMA journal_mode=WAL", "PRAGMA synchronous=FULL", "PRAGMA foreign_keys=ON"} {
|
||||
if _, err := db.Exec(pragma); err != nil {
|
||||
return nil, fmt.Errorf("initialize SQLite %s: %w", pragma, err)
|
||||
}
|
||||
}
|
||||
if _, err := db.Exec(currentSchema); err != nil {
|
||||
return nil, fmt.Errorf("create current SQLite schema: %w", err)
|
||||
}
|
||||
if version == 0 {
|
||||
if _, err := db.Exec(`PRAGMA user_version=1`); err != nil {
|
||||
return nil, fmt.Errorf("commit current SQLite layout version: %w", err)
|
||||
if len(seen) == 0 {
|
||||
tx, err := db.Begin()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
if _, err := tx.Exec(currentSchema); err != nil {
|
||||
return nil, fmt.Errorf("create current SQLite schema: %w", err)
|
||||
}
|
||||
if _, err := tx.Exec(`PRAGMA user_version=2`); err != nil {
|
||||
return nil, fmt.Errorf("set current SQLite layout version: %w", err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return nil, fmt.Errorf("commit current SQLite schema: %w", err)
|
||||
}
|
||||
} else if _, err := db.Exec(currentSchema); err != nil {
|
||||
return nil, fmt.Errorf("verify current SQLite schema: %w", err)
|
||||
}
|
||||
return &CurrentStore{db: db}, nil
|
||||
}
|
||||
|
||||
@@ -14,13 +14,28 @@ import (
|
||||
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
||||
)
|
||||
|
||||
var ErrCurrentUploadUnverified = errors.New("uploaded recording has no Dispatcher-approved target")
|
||||
var ErrCurrentUploadUnverified = errors.New("recording outcome is not verified against its Dispatcher-approved upload")
|
||||
var ErrCurrentResultConflict = errors.New("call already has a different final result")
|
||||
|
||||
// RecordCallResult stores one final result only after confirmed call end. The
|
||||
// same source command always maps to one durable outbox identity; an identical
|
||||
// retry can resume delivery, but a different result cannot replace it.
|
||||
func currentResultEventID(dispatcherID, sourceEventID string) string {
|
||||
identity := sha256.Sum256([]byte("call.execute.result\x00" + dispatcherID + "\x00" + sourceEventID))
|
||||
return fmt.Sprintf("result-%x", identity[:])
|
||||
}
|
||||
|
||||
// RecordCallResult stores one no-recording final result only after confirmed
|
||||
// call end. A durable upload grant cannot be bypassed with an empty recording.
|
||||
func (s *CurrentStore) RecordCallResult(dispatcherID, sourceEventID string, payload []byte) (CurrentOutboxEvent, bool, error) {
|
||||
return s.recordCallResult(dispatcherID, sourceEventID, payload, nil)
|
||||
}
|
||||
|
||||
// RecordUploadedCallResult accepts the authenticated Agent's successful PUT
|
||||
// observation only for the precise, previously bound OSS asset. Both the
|
||||
// upload fact and sole SaaS result outbox entry commit in the same transaction.
|
||||
func (s *CurrentStore) RecordUploadedCallResult(dispatcherID, sourceEventID string, payload []byte, proof CurrentUploadProof) (CurrentOutboxEvent, bool, error) {
|
||||
return s.recordCallResult(dispatcherID, sourceEventID, payload, &proof)
|
||||
}
|
||||
|
||||
func (s *CurrentStore) recordCallResult(dispatcherID, sourceEventID string, payload []byte, proof *CurrentUploadProof) (CurrentOutboxEvent, bool, error) {
|
||||
if dispatcherID == "" || sourceEventID == "" || len(sourceEventID) > 255 || len(payload) == 0 {
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result requires a durable call identity and payload")
|
||||
}
|
||||
@@ -72,20 +87,37 @@ func (s *CurrentStore) RecordCallResult(dispatcherID, sourceEventID string, payl
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result identity differs from the approved call")
|
||||
}
|
||||
var recording struct {
|
||||
Status string `json:"status"`
|
||||
Status string `json:"status"`
|
||||
Bucket string `json:"bucket"`
|
||||
ObjectKey string `json:"object_key"`
|
||||
Format string `json:"format"`
|
||||
Channels int `json:"channels"`
|
||||
SampleRateHz int `json:"sample_rate_hz"`
|
||||
DurationMS int64 `json:"duration_ms"`
|
||||
SizeBytes int64 `json:"size_bytes"`
|
||||
ChecksumSHA256 string `json:"checksum_sha256"`
|
||||
}
|
||||
if err := json.Unmarshal(result.Recording, &recording); err != nil {
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result recording fact is invalid")
|
||||
}
|
||||
grant, grantErr := loadRecordingUpload(tx.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, dispatcherID, sourceEventID))
|
||||
if grantErr != nil && !errors.Is(grantErr, ErrCurrentUploadNotFound) {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("inspect original recording target: %w", grantErr)
|
||||
}
|
||||
if recording.Status == "uploaded" {
|
||||
if proof == nil || grantErr != nil || proof.StatusCode < 200 || proof.StatusCode >= 300 ||
|
||||
proof.UploadID != grant.UploadID || proof.RecordingID != grant.RecordingID || proof.SizeBytes != grant.SizeBytes || proof.SHA256 != grant.ChecksumSHA256 ||
|
||||
recording.Bucket != grant.Bucket || recording.ObjectKey != grant.ObjectKey || recording.Format != grant.Format || recording.Channels != grant.Channels || recording.SampleRateHz != grant.SampleRateHz || recording.DurationMS != grant.DurationMS || recording.SizeBytes != grant.SizeBytes || recording.ChecksumSHA256 != grant.ChecksumSHA256 {
|
||||
return CurrentOutboxEvent{}, false, ErrCurrentUploadUnverified
|
||||
}
|
||||
} else if proof != nil || grantErr == nil {
|
||||
return CurrentOutboxEvent{}, false, ErrCurrentUploadUnverified
|
||||
}
|
||||
route, err := tenant.CurrentResultRoute(dispatcherID)
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, err
|
||||
}
|
||||
identity := sha256.Sum256([]byte("call.execute.result\x00" + dispatcherID + "\x00" + sourceEventID))
|
||||
eventID := fmt.Sprintf("result-%x", identity[:])
|
||||
eventID := currentResultEventID(dispatcherID, sourceEventID)
|
||||
body, err := json.Marshal(struct {
|
||||
EventID string `json:"event_id"`
|
||||
EventType string `json:"event_type"`
|
||||
@@ -116,6 +148,22 @@ func (s *CurrentStore) RecordCallResult(dispatcherID, sourceEventID string, payl
|
||||
if stored.EventType != "call.execute.result" || stored.RoutingKey != route.BindingKey || !bytes.Equal(stored.Body, body) {
|
||||
return CurrentOutboxEvent{}, false, ErrCurrentResultConflict
|
||||
}
|
||||
if proof != nil {
|
||||
if grant.ConfirmedAt != "" && count == 1 {
|
||||
return CurrentOutboxEvent{}, false, ErrCurrentResultConflict
|
||||
}
|
||||
confirmed, err := tx.Exec(`UPDATE dispatcher_recordings SET confirmed_at=COALESCE(confirmed_at,?) WHERE dispatcher_id=? AND source_event_id=?`, time.Now().UTC().Format(time.RFC3339Nano), dispatcherID, sourceEventID)
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("persist confirmed recording with result: %w", err)
|
||||
}
|
||||
affected, err := confirmed.RowsAffected()
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("inspect recording confirmation: %w", err)
|
||||
}
|
||||
if affected != 1 {
|
||||
return CurrentOutboxEvent{}, false, errors.New("recording confirmation lost its original target")
|
||||
}
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("commit final result outbox: %w", err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,118 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"encoding/hex"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
)
|
||||
|
||||
var ErrCurrentUploadConflict = errors.New("recording upload differs from its original approved target")
|
||||
var ErrCurrentUploadNotFound = errors.New("no durable original recording target")
|
||||
var ErrCurrentUploadAlreadyConfirmed = errors.New("recording upload already confirmed; another PUT is forbidden")
|
||||
|
||||
// CurrentUploadProof is the authenticated Agent's successful PUT observation;
|
||||
// the Dispatcher also verifies every asset field against its durable grant.
|
||||
// It does not claim an independent OSS HEAD or SaaS application receipt.
|
||||
type CurrentUploadProof struct {
|
||||
UploadID string
|
||||
RecordingID string
|
||||
StatusCode int
|
||||
SizeBytes int64
|
||||
SHA256 string
|
||||
}
|
||||
|
||||
// CurrentRecordingGrant holds the original D-approved OSS destination. Signed
|
||||
// URLs, headers and short-lived tokens are deliberately never stored here.
|
||||
type CurrentRecordingGrant struct {
|
||||
DispatcherID string
|
||||
SourceEventID string
|
||||
UploadID string
|
||||
RecordingID string
|
||||
Bucket string
|
||||
ObjectKey string
|
||||
ChecksumSHA256 string
|
||||
SizeBytes int64
|
||||
Format string
|
||||
Channels int
|
||||
SampleRateHz int
|
||||
DurationMS int64
|
||||
ConfirmedAt string
|
||||
}
|
||||
|
||||
func (s *CurrentStore) BindRecordingUpload(proposed CurrentRecordingGrant) (CurrentRecordingGrant, bool, error) {
|
||||
if proposed.DispatcherID == "" || proposed.SourceEventID == "" || proposed.UploadID == "" || proposed.RecordingID == "" || proposed.Bucket == "" || proposed.ObjectKey == "" || proposed.SizeBytes <= 0 || proposed.Format != "wav" || (proposed.Channels != 1 && proposed.Channels != 2) || proposed.SampleRateHz != 16000 || proposed.DurationMS < 0 || proposed.ConfirmedAt != "" {
|
||||
return CurrentRecordingGrant{}, false, errors.New("recording grant lacks a bounded original asset and execution")
|
||||
}
|
||||
if len(proposed.ChecksumSHA256) != 64 || strings.ToLower(proposed.ChecksumSHA256) != proposed.ChecksumSHA256 {
|
||||
return CurrentRecordingGrant{}, false, errors.New("recording grant requires a canonical SHA-256 digest")
|
||||
}
|
||||
if _, err := hex.DecodeString(proposed.ChecksumSHA256); err != nil {
|
||||
return CurrentRecordingGrant{}, false, errors.New("recording grant has an invalid SHA-256 digest")
|
||||
}
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return CurrentRecordingGrant{}, false, err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var status, trunkID string
|
||||
var snapshot []byte
|
||||
err = tx.QueryRow(`SELECT status,COALESCE(selected_trunk_id,''),snapshot_json FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, proposed.DispatcherID, proposed.SourceEventID).Scan(&status, &trunkID, &snapshot)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return CurrentRecordingGrant{}, false, ErrCurrentUploadNotFound
|
||||
}
|
||||
if err != nil {
|
||||
return CurrentRecordingGrant{}, false, fmt.Errorf("load approved call for recording: %w", err)
|
||||
}
|
||||
if status != "dispatched" && status != "unknown" && status != "finished" || trunkID == "" || len(snapshot) == 0 {
|
||||
return CurrentRecordingGrant{}, false, errors.New("recording grant requires an already reserved approved call")
|
||||
}
|
||||
inserted, err := tx.Exec(`INSERT INTO dispatcher_recordings(dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms)
|
||||
VALUES(?,?,?,?,?,?,?,?,?,?,?,?) ON CONFLICT(dispatcher_id,source_event_id) DO NOTHING`, proposed.DispatcherID, proposed.SourceEventID, proposed.UploadID, proposed.RecordingID, proposed.Bucket, proposed.ObjectKey, proposed.ChecksumSHA256, proposed.SizeBytes, proposed.Format, proposed.Channels, proposed.SampleRateHz, proposed.DurationMS)
|
||||
if err != nil {
|
||||
return CurrentRecordingGrant{}, false, fmt.Errorf("persist immutable recording target: %w", err)
|
||||
}
|
||||
count, err := inserted.RowsAffected()
|
||||
if err != nil {
|
||||
return CurrentRecordingGrant{}, false, err
|
||||
}
|
||||
stored, err := loadRecordingUpload(tx.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, proposed.DispatcherID, proposed.SourceEventID))
|
||||
if err != nil {
|
||||
return CurrentRecordingGrant{}, false, fmt.Errorf("inspect persisted recording target: %w", err)
|
||||
}
|
||||
if stored.ConfirmedAt != "" {
|
||||
return CurrentRecordingGrant{}, false, ErrCurrentUploadAlreadyConfirmed
|
||||
}
|
||||
if stored != proposed {
|
||||
return CurrentRecordingGrant{}, false, ErrCurrentUploadConflict
|
||||
}
|
||||
var finalized int
|
||||
err = tx.QueryRow(`SELECT 1 FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, proposed.DispatcherID, currentResultEventID(proposed.DispatcherID, proposed.SourceEventID)).Scan(&finalized)
|
||||
if err == nil {
|
||||
return CurrentRecordingGrant{}, false, ErrCurrentResultConflict
|
||||
}
|
||||
if !errors.Is(err, sql.ErrNoRows) {
|
||||
return CurrentRecordingGrant{}, false, fmt.Errorf("check existing call result before granting an upload: %w", err)
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return CurrentRecordingGrant{}, false, fmt.Errorf("commit original recording target: %w", err)
|
||||
}
|
||||
return stored, count == 1, nil
|
||||
}
|
||||
|
||||
func (s *CurrentStore) LoadRecordingUpload(dispatcherID, sourceEventID string) (CurrentRecordingGrant, error) {
|
||||
if dispatcherID == "" || sourceEventID == "" {
|
||||
return CurrentRecordingGrant{}, ErrCurrentUploadNotFound
|
||||
}
|
||||
return loadRecordingUpload(s.db.QueryRow(`SELECT dispatcher_id,source_event_id,upload_id,recording_id,bucket,object_key,checksum_sha256,size_bytes,format,channels,sample_rate_hz,duration_ms,COALESCE(confirmed_at,'') FROM dispatcher_recordings WHERE dispatcher_id=? AND source_event_id=?`, dispatcherID, sourceEventID))
|
||||
}
|
||||
|
||||
func loadRecordingUpload(row *sql.Row) (CurrentRecordingGrant, error) {
|
||||
var grant CurrentRecordingGrant
|
||||
err := row.Scan(&grant.DispatcherID, &grant.SourceEventID, &grant.UploadID, &grant.RecordingID, &grant.Bucket, &grant.ObjectKey, &grant.ChecksumSHA256, &grant.SizeBytes, &grant.Format, &grant.Channels, &grant.SampleRateHz, &grant.DurationMS, &grant.ConfirmedAt)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return CurrentRecordingGrant{}, ErrCurrentUploadNotFound
|
||||
}
|
||||
return grant, err
|
||||
}
|
||||
@@ -0,0 +1,196 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func currentUploadedResultPayload(t *testing.T, binding CurrentRecordingGrant) []byte {
|
||||
t.Helper()
|
||||
var payload map[string]any
|
||||
if err := json.Unmarshal(currentResultPayload(t), &payload); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
payload["outcome"] = "answered"
|
||||
payload["reason_code"] = 200
|
||||
payload["reason_message"] = "completed"
|
||||
payload["recording"] = map[string]any{
|
||||
"status": "uploaded", "bucket": binding.Bucket, "object_key": binding.ObjectKey,
|
||||
"format": binding.Format, "channels": binding.Channels, "sample_rate_hz": binding.SampleRateHz,
|
||||
"duration_ms": binding.DurationMS, "size_bytes": binding.SizeBytes,
|
||||
"checksum_sha256": binding.ChecksumSHA256,
|
||||
}
|
||||
raw, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return raw
|
||||
}
|
||||
|
||||
func currentUploadProof(binding CurrentRecordingGrant) CurrentUploadProof {
|
||||
return CurrentUploadProof{UploadID: binding.UploadID, RecordingID: binding.RecordingID, StatusCode: 200, SizeBytes: binding.SizeBytes, SHA256: binding.ChecksumSHA256}
|
||||
}
|
||||
|
||||
func TestCurrentUploadedResultRequiresOriginalGrantAndConfirmedPUT(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "uploaded-result-1", true)
|
||||
binding := currentUploadBinding(cmd)
|
||||
payload := currentUploadedResultPayload(t, binding)
|
||||
proof := currentUploadProof(binding)
|
||||
if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof); !errors.Is(err, ErrCurrentUploadUnverified) {
|
||||
t.Fatalf("ungranted asset was reported as uploaded: %v", err)
|
||||
}
|
||||
if _, _, err := s.BindRecordingUpload(binding); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)); !errors.Is(err, ErrCurrentUploadUnverified) {
|
||||
t.Fatalf("failed OSS path bypassed itself with an empty recording: %v", err)
|
||||
}
|
||||
if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, CurrentUploadProof{}); !errors.Is(err, ErrCurrentUploadUnverified) {
|
||||
t.Fatalf("missing PUT confirmation was accepted: %v", err)
|
||||
}
|
||||
original, created, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof)
|
||||
if err != nil || !created || original.EventType != "call.execute.result" {
|
||||
t.Fatalf("confirmed grant was not durably reported: created=%t event=%+v err=%v", created, original, err)
|
||||
}
|
||||
stored, err := s.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID)
|
||||
if err != nil || stored.ConfirmedAt == "" {
|
||||
t.Fatalf("confirmed OSS fact was not committed with the outbox: stored=%+v err=%v", stored, err)
|
||||
}
|
||||
repeated, created, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof)
|
||||
if err != nil || created || repeated.EventID != original.EventID {
|
||||
t.Fatalf("same PUT fact created a second MQ result: created=%t err=%v", created, err)
|
||||
}
|
||||
if _, _, err := s.BindRecordingUpload(binding); !errors.Is(err, ErrCurrentUploadAlreadyConfirmed) {
|
||||
t.Fatalf("confirmed upload received another grant: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentUploadedResultRejectsChangedOSSMetadataOrProof(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "uploaded-changed-1", true)
|
||||
binding := currentUploadBinding(cmd)
|
||||
if _, _, err := s.BindRecordingUpload(binding); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
originalPayload := currentUploadedResultPayload(t, binding)
|
||||
proof := currentUploadProof(binding)
|
||||
for _, tc := range []struct {
|
||||
field string
|
||||
value any
|
||||
}{
|
||||
{"bucket", "other-bucket"}, {"object_key", "other/key.wav"}, {"checksum_sha256", strings.Repeat("b", 64)},
|
||||
{"size_bytes", binding.SizeBytes + 1}, {"format", "mp3"}, {"channels", 2}, {"sample_rate_hz", 8000}, {"duration_ms", binding.DurationMS + 1},
|
||||
} {
|
||||
var result map[string]any
|
||||
if err := json.Unmarshal(originalPayload, &result); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
result["recording"].(map[string]any)[tc.field] = tc.value
|
||||
altered, _ := json.Marshal(result)
|
||||
if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, altered, proof); err == nil {
|
||||
t.Fatalf("changed %s was reported as the original upload", tc.field)
|
||||
}
|
||||
}
|
||||
for _, change := range []struct {
|
||||
name string
|
||||
apply func(*CurrentUploadProof)
|
||||
}{
|
||||
{"status", func(p *CurrentUploadProof) { p.StatusCode = 503 }},
|
||||
{"upload ID", func(p *CurrentUploadProof) { p.UploadID = "new-upload" }},
|
||||
{"recording ID", func(p *CurrentUploadProof) { p.RecordingID = "new-recording" }},
|
||||
{"checksum", func(p *CurrentUploadProof) { p.SHA256 = strings.Repeat("b", 64) }},
|
||||
{"size", func(p *CurrentUploadProof) { p.SizeBytes++ }},
|
||||
} {
|
||||
altered := proof
|
||||
change.apply(&altered)
|
||||
if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, originalPayload, altered); !errors.Is(err, ErrCurrentUploadUnverified) {
|
||||
t.Fatalf("changed %s proof was accepted: %v", change.name, err)
|
||||
}
|
||||
}
|
||||
if outbox, err := s.ListPendingOutbox(cmd.DispatcherID); err != nil || len(outbox) != 1 {
|
||||
t.Fatalf("invalid upload escaped to SaaS: outbox=%+v err=%v", outbox, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentUploadedResultOutboxFailureDoesNotConfirmGrant(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "uploaded-outbox-fault-1", true)
|
||||
binding := currentUploadBinding(cmd)
|
||||
if _, _, err := s.BindRecordingUpload(binding); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := s.db.Exec(`CREATE TRIGGER fail_upload_final BEFORE INSERT ON dispatcher_outbox WHEN NEW.event_type='call.execute.result' BEGIN SELECT RAISE(ABORT,'injected outbox failure'); END`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
payload, proof := currentUploadedResultPayload(t, binding), currentUploadProof(binding)
|
||||
if _, _, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof); err == nil {
|
||||
t.Fatal("failed SQLite outbox was mistaken for delivered upload")
|
||||
}
|
||||
if stored, err := s.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID); err != nil || stored.ConfirmedAt != "" {
|
||||
t.Fatalf("partial confirmation escaped failed transaction: state=%+v err=%v", stored, err)
|
||||
}
|
||||
if _, err := s.db.Exec(`DROP TRIGGER fail_upload_final`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, created, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof); err != nil || !created {
|
||||
t.Fatalf("same uploaded fact could not resume without another PUT: created=%t err=%v", created, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentUploadCannotBeginAfterEmptyFinalResult(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "empty-result-before-grant", true)
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.BindRecordingUpload(currentUploadBinding(cmd)); !errors.Is(err, ErrCurrentResultConflict) {
|
||||
t.Fatalf("already final no-recording call gained a new upload target: %v", err)
|
||||
}
|
||||
if _, err := s.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID); !errors.Is(err, ErrCurrentUploadNotFound) {
|
||||
t.Fatalf("unexpected recording for finalized empty result: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentUploadedResultRestartReusesGrantAndPendingOutbox(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "uploaded-restart-1", true)
|
||||
binding := currentUploadBinding(cmd)
|
||||
if _, _, err := s.BindRecordingUpload(binding); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
payload, proof := currentUploadedResultPayload(t, binding), currentUploadProof(binding)
|
||||
original, created, err := s.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof)
|
||||
if err != nil || !created {
|
||||
t.Fatalf("persist original result: created=%t err=%v", created, err)
|
||||
}
|
||||
var seq int
|
||||
var databaseName, databasePath string
|
||||
if err := s.db.QueryRow(`PRAGMA database_list`).Scan(&seq, &databaseName, &databasePath); err != nil || databasePath == "" {
|
||||
t.Fatalf("locate isolated database: %v", err)
|
||||
}
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
restarted, err := OpenCurrent(databasePath)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer restarted.Close()
|
||||
grant, err := restarted.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID)
|
||||
if err != nil || grant.ConfirmedAt == "" || grant.ObjectKey != binding.ObjectKey {
|
||||
t.Fatalf("restarted D forgot uploaded fact or original OSS target: grant=%+v err=%v", grant, err)
|
||||
}
|
||||
pending, err := restarted.ListPendingOutbox(cmd.DispatcherID)
|
||||
if err != nil || len(pending) != 2 || string(pending[1].Body) != string(original.Body) {
|
||||
t.Fatalf("restarted D lost the sole undelivered MQ result: events=%+v err=%v", pending, err)
|
||||
}
|
||||
if repeated, created, err := restarted.RecordUploadedCallResult(cmd.DispatcherID, cmd.EventID, payload, proof); err != nil || created || repeated.EventID != original.EventID {
|
||||
t.Fatalf("restarted D created a second result: created=%t err=%v", created, err)
|
||||
}
|
||||
if _, _, err := restarted.BindRecordingUpload(binding); !errors.Is(err, ErrCurrentUploadAlreadyConfirmed) {
|
||||
t.Fatalf("restarted D issued another PUT grant: %v", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,156 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"errors"
|
||||
"path/filepath"
|
||||
"reflect"
|
||||
"strings"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func currentUploadBinding(cmd CurrentExecuteCommand) CurrentRecordingGrant {
|
||||
return CurrentRecordingGrant{
|
||||
DispatcherID: cmd.DispatcherID, SourceEventID: cmd.EventID,
|
||||
UploadID: "upload-" + cmd.EventID, RecordingID: "recording-" + cmd.EventID,
|
||||
Bucket: "mock-bucket", ObjectKey: "tenant/1001/" + cmd.EventID + "/rec.wav",
|
||||
ChecksumSHA256: strings.Repeat("a", 64), SizeBytes: 128,
|
||||
Format: "wav", Channels: 1, SampleRateHz: 16000, DurationMS: 3000,
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentUploadBindingPersistsOnlyOneOriginalAssetAcrossRestart(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "grant-stable-1", true)
|
||||
original := currentUploadBinding(cmd)
|
||||
stored, created, err := s.BindRecordingUpload(original)
|
||||
if err != nil || !created || !reflect.DeepEqual(stored, original) {
|
||||
t.Fatalf("original D-owned bucket/key not bound: stored=%+v created=%t err=%v", stored, created, err)
|
||||
}
|
||||
stored, created, err = s.BindRecordingUpload(original)
|
||||
if err != nil || created || !reflect.DeepEqual(stored, original) {
|
||||
t.Fatalf("explicit grant request changed the original target: stored=%+v created=%t err=%v", stored, created, err)
|
||||
}
|
||||
var version, seq int
|
||||
var databaseName, databasePath string
|
||||
if err := s.db.QueryRow(`PRAGMA user_version`).Scan(&version); err != nil || version != 2 {
|
||||
t.Fatalf("current schema is not the new, reviewed recording layout: version=%d err=%v", version, err)
|
||||
}
|
||||
if err := s.db.QueryRow(`PRAGMA database_list`).Scan(&seq, &databaseName, &databasePath); err != nil || databasePath == "" {
|
||||
t.Fatalf("locate isolated SQLite fixture: %v", err)
|
||||
}
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
restarted, err := OpenCurrent(databasePath)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer restarted.Close()
|
||||
loaded, err := restarted.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID)
|
||||
if err != nil || !reflect.DeepEqual(loaded, original) {
|
||||
t.Fatalf("restart changed the only OSS target: loaded=%+v err=%v", loaded, err)
|
||||
}
|
||||
var schema string
|
||||
if err := restarted.db.QueryRow(`SELECT sql FROM sqlite_master WHERE name='dispatcher_recordings'`).Scan(&schema); err != nil || strings.Contains(strings.ToLower(schema), "token") || strings.Contains(strings.ToLower(schema), "signed_url") || strings.Contains(strings.ToLower(schema), "target_url") {
|
||||
t.Fatalf("temporary OSS credential was persisted in grant table: err=%v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentUploadBindingRejectsChangedTargetAndUnreservedCalls(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
pending := currentCall("grant-not-reserved")
|
||||
if _, _, err := s.RecordExecute(pending); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.BindRecordingUpload(currentUploadBinding(pending)); err == nil {
|
||||
t.Fatal("unreserved instruction granted an OSS upload")
|
||||
}
|
||||
cmd := currentResultCall(t, s, "grant-immutable-1", true)
|
||||
original := currentUploadBinding(cmd)
|
||||
if _, _, err := s.BindRecordingUpload(original); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, change := range []struct {
|
||||
name string
|
||||
apply func(*CurrentRecordingGrant)
|
||||
}{
|
||||
{"upload ID", func(x *CurrentRecordingGrant) { x.UploadID = "new-upload" }},
|
||||
{"recording ID", func(x *CurrentRecordingGrant) { x.RecordingID = "other-recording" }},
|
||||
{"bucket", func(x *CurrentRecordingGrant) { x.Bucket = "changed-bucket" }},
|
||||
{"object key", func(x *CurrentRecordingGrant) { x.ObjectKey = "another/key.wav" }},
|
||||
{"checksum", func(x *CurrentRecordingGrant) { x.ChecksumSHA256 = strings.Repeat("b", 64) }},
|
||||
{"size", func(x *CurrentRecordingGrant) { x.SizeBytes++ }},
|
||||
{"channels", func(x *CurrentRecordingGrant) { x.Channels = 2 }},
|
||||
} {
|
||||
changed := original
|
||||
change.apply(&changed)
|
||||
if _, _, err := s.BindRecordingUpload(changed); !errors.Is(err, ErrCurrentUploadConflict) {
|
||||
t.Fatalf("%s replaced the immutable original: %v", change.name, err)
|
||||
}
|
||||
}
|
||||
if loaded, err := s.LoadRecordingUpload(cmd.DispatcherID, cmd.EventID); err != nil || !reflect.DeepEqual(loaded, original) {
|
||||
t.Fatalf("conflicting grant changed durable target: loaded=%+v err=%v", loaded, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentUploadBindingRejectsReusedUploadIDOrOSSObject(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
first := currentUploadBinding(currentResultCall(t, s, "grant-unique-1", true))
|
||||
second := currentUploadBinding(currentResultCall(t, s, "grant-unique-2", true))
|
||||
if _, _, err := s.BindRecordingUpload(first); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
second.UploadID = first.UploadID
|
||||
if _, _, err := s.BindRecordingUpload(second); err == nil {
|
||||
t.Fatal("two calls shared one upload identity")
|
||||
}
|
||||
second.UploadID = "upload-grant-unique-2"
|
||||
second.ObjectKey = first.ObjectKey
|
||||
if _, _, err := s.BindRecordingUpload(second); err == nil {
|
||||
t.Fatal("two calls shared one original OSS object")
|
||||
}
|
||||
if _, err := s.LoadRecordingUpload(second.DispatcherID, second.SourceEventID); !errors.Is(err, ErrCurrentUploadNotFound) {
|
||||
t.Fatalf("failed uniqueness check left a partial grant: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentSchemaRejectsPreviousCompleteLayoutWithoutErasingOutbox(t *testing.T) {
|
||||
path := filepath.Join(t.TempDir(), "previous-current.db")
|
||||
s, err := OpenCurrent(path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := s.db.Exec(`INSERT INTO dispatcher_outbox(dispatcher_id,event_id,event_type,routing_key,body) VALUES(?,?,?,?,?)`, currentDispatcherID, "not-delivered", "call.execute.result", "mock.route", []byte(`{"persisted":true}`)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := s.db.Exec(`DROP TABLE dispatcher_recordings`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, err := s.db.Exec(`PRAGMA user_version=1`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if reopened, err := OpenCurrent(path); err == nil {
|
||||
reopened.Close()
|
||||
t.Fatal("incomplete previous SQLite layout was silently modified or admitted")
|
||||
}
|
||||
old, err := sql.Open("sqlite", path)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer old.Close()
|
||||
var body []byte
|
||||
if err := old.QueryRow(`SELECT body FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, currentDispatcherID, "not-delivered").Scan(&body); err != nil || string(body) != `{"persisted":true}` {
|
||||
t.Fatalf("undelivered result was altered by refused admission: body=%s err=%v", body, err)
|
||||
}
|
||||
var count, version int
|
||||
if err := old.QueryRow(`SELECT COUNT(*) FROM sqlite_master WHERE name='dispatcher_recordings'`).Scan(&count); err != nil || count != 0 {
|
||||
t.Fatalf("refused database was automatically migrated: tables=%d err=%v", count, err)
|
||||
}
|
||||
if err := old.QueryRow(`PRAGMA user_version`).Scan(&version); err != nil || version != 1 {
|
||||
t.Fatalf("refused database layout was modified: version=%d err=%v", version, err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user