Persist one confirmed call result per execution
This commit is contained in:
@@ -60,7 +60,8 @@
|
||||
|
||||
- 已新增 Agent 内存录音单次 PUT 组件 `UploadClient.UploadBytes`:复用受限授权、大小和 SHA-256 核对,成功路径不写临时录音或通话结果文件。OSS 明确拒绝保留状态码且不自动重试;传输结果不明时返回专用错误、停止自动重试并隐藏带签名的 URL。空授权头拒绝而非使服务崩溃。这里只验证组件,尚未连接实际录音或最终结果。
|
||||
- Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;真实 Dispatcher 授权及回报尚未接线。
|
||||
- 已验证:`go test ./... -count=1`、`go test -race ./internal/agent -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 唯一最终结果 outbox、真实 Agent↔Dispatcher 结果交付及 MQ/端到端验收,不能宣称 P06 通过。
|
||||
- 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 通过。
|
||||
|
||||
## 验收台账
|
||||
|
||||
|
||||
@@ -0,0 +1,123 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"crypto/sha256"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/configread"
|
||||
"git.ipao.vip/rogee/go-sip/internal/contract"
|
||||
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
||||
)
|
||||
|
||||
var ErrCurrentUploadUnverified = errors.New("uploaded recording has no Dispatcher-approved target")
|
||||
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 (s *CurrentStore) RecordCallResult(dispatcherID, sourceEventID string, payload []byte) (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")
|
||||
}
|
||||
var result struct {
|
||||
TaskID string `json:"task_id"`
|
||||
CallerProfileID string `json:"caller_profile_id"`
|
||||
Callee string `json:"callee"`
|
||||
TrunkID string `json:"trunk_id"`
|
||||
StartedAt string `json:"started_at"`
|
||||
EndedAt string `json:"ended_at"`
|
||||
Recording json.RawMessage `json:"recording"`
|
||||
}
|
||||
if err := json.Unmarshal(payload, &result); err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("decode final result: %w", err)
|
||||
}
|
||||
start, err := time.Parse(time.RFC3339Nano, result.StartedAt)
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result has invalid start time")
|
||||
}
|
||||
end, err := time.Parse(time.RFC3339Nano, result.EndedAt)
|
||||
if err != nil || end.Before(start) {
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result end precedes start or is invalid")
|
||||
}
|
||||
tx, err := s.db.Begin()
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, err
|
||||
}
|
||||
defer tx.Rollback()
|
||||
var tenantID int64
|
||||
var taskID, callee, trunkID, status string
|
||||
var snapshotJSON []byte
|
||||
err = tx.QueryRow(`SELECT tenant_id,task_id,callee,COALESCE(selected_trunk_id,''),status,snapshot_json FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, sourceEventID).Scan(&tenantID, &taskID, &callee, &trunkID, &status, &snapshotJSON)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return CurrentOutboxEvent{}, false, errors.New("no durable approved call matches final result")
|
||||
}
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("load completed call identity: %w", err)
|
||||
}
|
||||
if status != "finished" || trunkID == "" || len(snapshotJSON) == 0 {
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result requires confirmed call end and its frozen snapshot")
|
||||
}
|
||||
var snapshot struct {
|
||||
Task configread.CurrentTask `json:"task"`
|
||||
}
|
||||
if err := json.Unmarshal(snapshotJSON, &snapshot); err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("decode frozen call task: %w", err)
|
||||
}
|
||||
if snapshot.Task.DispatcherID != dispatcherID || snapshot.Task.TenantID != tenantID || snapshot.Task.TaskID != taskID || result.TaskID != taskID || result.Callee != callee || result.TrunkID != trunkID || result.CallerProfileID != snapshot.Task.CallerProfileID {
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result identity differs from the approved call")
|
||||
}
|
||||
var recording struct {
|
||||
Status string `json:"status"`
|
||||
}
|
||||
if err := json.Unmarshal(result.Recording, &recording); err != nil {
|
||||
return CurrentOutboxEvent{}, false, errors.New("final result recording fact is invalid")
|
||||
}
|
||||
if recording.Status == "uploaded" {
|
||||
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[:])
|
||||
body, err := json.Marshal(struct {
|
||||
EventID string `json:"event_id"`
|
||||
EventType string `json:"event_type"`
|
||||
DispatcherID string `json:"dispatcher_id"`
|
||||
TenantID int64 `json:"tenant_id"`
|
||||
IssuedAt string `json:"issued_at"`
|
||||
Payload json.RawMessage `json:"payload"`
|
||||
}{eventID, "call.execute.result", dispatcherID, tenantID, result.EndedAt, payload})
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("encode final result: %w", err)
|
||||
}
|
||||
if err := contract.ValidateCurrent("mq", body); err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("final result violates current MQ contract: %w", err)
|
||||
}
|
||||
inserted, err := tx.Exec(`INSERT INTO dispatcher_outbox(dispatcher_id,event_id,event_type,routing_key,body) VALUES(?,?,?,?,?) ON CONFLICT(dispatcher_id,event_id) DO NOTHING`, dispatcherID, eventID, "call.execute.result", route.BindingKey, body)
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("persist final result outbox: %w", err)
|
||||
}
|
||||
count, err := inserted.RowsAffected()
|
||||
if err != nil {
|
||||
return CurrentOutboxEvent{}, false, err
|
||||
}
|
||||
var stored CurrentOutboxEvent
|
||||
stored.EventID = eventID
|
||||
if err := tx.QueryRow(`SELECT event_type,routing_key,body FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, eventID).Scan(&stored.EventType, &stored.RoutingKey, &stored.Body); err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("inspect existing final result: %w", err)
|
||||
}
|
||||
if stored.EventType != "call.execute.result" || stored.RoutingKey != route.BindingKey || !bytes.Equal(stored.Body, body) {
|
||||
return CurrentOutboxEvent{}, false, ErrCurrentResultConflict
|
||||
}
|
||||
if err := tx.Commit(); err != nil {
|
||||
return CurrentOutboxEvent{}, false, fmt.Errorf("commit final result outbox: %w", err)
|
||||
}
|
||||
return stored, count == 1, nil
|
||||
}
|
||||
@@ -0,0 +1,204 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.ipao.vip/rogee/go-sip/internal/contract"
|
||||
"git.ipao.vip/rogee/go-sip/internal/tenant"
|
||||
)
|
||||
|
||||
func currentResultPayload(t *testing.T) []byte {
|
||||
t.Helper()
|
||||
payload := map[string]any{
|
||||
"task_id": "task-asr", "caller_profile_id": currentStoreSnapshot(t).Task.CallerProfileID,
|
||||
"callee": "15003164745", "trunk_id": "trunk-mock",
|
||||
"started_at": "2026-09-21T01:30:01Z", "ended_at": "2026-09-21T01:30:06Z", "duration_ms": 5000,
|
||||
"outcome": "no_answer", "reason_code": 486, "reason_message": "busy",
|
||||
"transcript": []any{}, "opt_out": false, "recording": map[string]any{},
|
||||
}
|
||||
raw, err := json.Marshal(payload)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
return raw
|
||||
}
|
||||
|
||||
func currentResultCall(t *testing.T, s *CurrentStore, id string, finish bool) CurrentExecuteCommand {
|
||||
t.Helper()
|
||||
cmd := currentCall(id)
|
||||
if _, _, err := s.RecordExecute(cmd); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.ReserveExecute(cmd.DispatcherID, cmd.EventID, currentReservation(), time.Date(2026, 9, 21, 1, 30, 0, 0, time.UTC)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := s.MarkExecuteDispatched(cmd.DispatcherID, cmd.EventID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if finish {
|
||||
if err := s.FinishExecute(cmd.DispatcherID, cmd.EventID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
return cmd
|
||||
}
|
||||
|
||||
func TestCurrentFinalResultNeedsConfirmedEndAndRemainsExactlyOne(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "result-call-1", false)
|
||||
payload := currentResultPayload(t)
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, payload); err == nil {
|
||||
t.Fatal("Agent report released a call before confirmed hangup")
|
||||
}
|
||||
if err := s.FinishExecute(cmd.DispatcherID, cmd.EventID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if occupancy, err := s.TrunkOccupancy(cmd.DispatcherID); err != nil || occupancy["trunk-mock"] != 0 {
|
||||
t.Fatalf("confirmed end did not release the call before OSS/result delivery: %+v err=%v", occupancy, err)
|
||||
}
|
||||
event, created, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, payload)
|
||||
if err != nil || !created || event.EventType != "call.execute.result" || event.EventID == cmd.EventID {
|
||||
t.Fatalf("exact final result not persisted: event=%+v created=%t err=%v", event, created, err)
|
||||
}
|
||||
route, err := tenant.CurrentResultRoute(cmd.DispatcherID)
|
||||
if err != nil || event.RoutingKey != route.BindingKey || contract.ValidateCurrent("mq", event.Body) != nil {
|
||||
t.Fatalf("final result was not Schema-valid or routed to shared SaaS queue: route=%+v err=%v", route, err)
|
||||
}
|
||||
outbox, err := s.ListPendingOutbox(cmd.DispatcherID)
|
||||
if err != nil || len(outbox) != 2 || !bytes.Equal(outbox[1].Body, event.Body) {
|
||||
t.Fatalf("call ack or final result was lost: outbox=%+v err=%v", outbox, err)
|
||||
}
|
||||
repeated, created, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, payload)
|
||||
if err != nil || created || !bytes.Equal(repeated.Body, event.Body) {
|
||||
t.Fatalf("redelivered identical result inserted twice: created=%t err=%v", created, err)
|
||||
}
|
||||
for _, id := range []string{cmd.EventID, event.EventID} {
|
||||
if err := s.MarkOutboxConfirmed(cmd.DispatcherID, id); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
if _, created, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, payload); err != nil || created {
|
||||
t.Fatalf("confirmed outbox row was replaced instead of retained: created=%t err=%v", created, err)
|
||||
}
|
||||
if pending, err := s.ListPendingOutbox(cmd.DispatcherID); err != nil || len(pending) != 0 {
|
||||
t.Fatalf("confirmed result was republished: pending=%+v err=%v", pending, err)
|
||||
}
|
||||
var total int
|
||||
if err := s.db.QueryRow(`SELECT COUNT(*) FROM dispatcher_outbox WHERE dispatcher_id=?`, cmd.DispatcherID).Scan(&total); err != nil || total != 2 {
|
||||
t.Fatalf("result history disappeared or multiplied: total=%d err=%v", total, err)
|
||||
}
|
||||
var claim map[string]any
|
||||
if err := json.Unmarshal(payload, &claim); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
claim["reason_message"] = "different result after receipt loss"
|
||||
changed, _ := json.Marshal(claim)
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, changed); err == nil {
|
||||
t.Fatal("same call identity accepted a second, conflicting final result")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentFinalResultRejectsDifferentBoundTaskCallOrTrunk(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "result-bound-1", true)
|
||||
for _, change := range []struct{ field, value string }{
|
||||
{"task_id", "other-task"}, {"callee", "15830461047"}, {"trunk_id", "other-trunk"}, {"caller_profile_id", "other-caller"},
|
||||
} {
|
||||
var claim map[string]any
|
||||
if err := json.Unmarshal(currentResultPayload(t), &claim); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
claim[change.field] = change.value
|
||||
raw, _ := json.Marshal(claim)
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, raw); err == nil {
|
||||
t.Fatalf("unbound %s was accepted in final result", change.field)
|
||||
}
|
||||
}
|
||||
if outbox, err := s.ListPendingOutbox(cmd.DispatcherID); err != nil || len(outbox) != 1 {
|
||||
t.Fatalf("invalid result escaped into SaaS outbox: outbox=%+v err=%v", outbox, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentFinalResultRejectsUnverifiedOSSAssetAndInventedState(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "result-unverified-1", true)
|
||||
var claim map[string]any
|
||||
if err := json.Unmarshal(currentResultPayload(t), &claim); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
claim["recording"] = map[string]any{"status": "uploaded", "bucket": "mock-bucket", "object_key": "tenant/rec.wav", "format": "wav", "channels": 1, "sample_rate_hz": 16000, "duration_ms": 3000, "size_bytes": 100, "checksum_sha256": strings.Repeat("a", 64)}
|
||||
uploaded, _ := json.Marshal(claim)
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, uploaded); !errors.Is(err, ErrCurrentUploadUnverified) {
|
||||
t.Fatalf("Agent invented an ungranted uploaded asset: %v", err)
|
||||
}
|
||||
claim["recording"] = map[string]any{"status": "unavailable"}
|
||||
invented, _ := json.Marshal(claim)
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, invented); err == nil {
|
||||
t.Fatal("unapproved recording status bypassed current MQ Schema")
|
||||
}
|
||||
if outbox, err := s.ListPendingOutbox(cmd.DispatcherID); err != nil || len(outbox) != 1 {
|
||||
t.Fatalf("unverified asset produced a result: outbox=%+v err=%v", outbox, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentFinalResultOutboxFailureRollsBackAndCanResume(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "result-fault-1", true)
|
||||
if _, err := s.db.Exec(`CREATE TRIGGER fail_final BEFORE INSERT ON dispatcher_outbox WHEN NEW.event_type='call.execute.result' BEGIN SELECT RAISE(ABORT,'injected final outbox failure'); END`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, _, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)); err == nil {
|
||||
t.Fatal("failed SQLite outbox insert was treated as delivered")
|
||||
}
|
||||
if outbox, err := s.ListPendingOutbox(cmd.DispatcherID); err != nil || len(outbox) != 1 {
|
||||
t.Fatalf("partial result leaked outside failed transaction: outbox=%+v err=%v", outbox, err)
|
||||
}
|
||||
if _, err := s.db.Exec(`DROP TRIGGER fail_final`); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if _, created, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t)); err != nil || !created {
|
||||
t.Fatalf("original result was not recoverable after SQLite failure: created=%t err=%v", created, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCurrentFinalResultRestartResumesSameOutboxBody(t *testing.T) {
|
||||
s := preparedCurrentCallStore(t)
|
||||
cmd := currentResultCall(t, s, "result-restart-1", true)
|
||||
original, created, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t))
|
||||
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("inspect private SQLite fixture: name=%q err=%v", databaseName, err)
|
||||
}
|
||||
if err := s.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
restarted, err := OpenCurrent(databasePath)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer restarted.Close()
|
||||
pending, err := restarted.ListPendingOutbox(cmd.DispatcherID)
|
||||
if err != nil || len(pending) != 2 || !bytes.Equal(pending[1].Body, original.Body) {
|
||||
t.Fatalf("restart lost the undelivered original MQ result: events=%+v err=%v", pending, err)
|
||||
}
|
||||
repeated, created, err := restarted.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t))
|
||||
if err != nil || created || repeated.EventID != original.EventID || !bytes.Equal(repeated.Body, original.Body) {
|
||||
t.Fatalf("retry created a second result after restart: created=%t err=%v", created, err)
|
||||
}
|
||||
if err := restarted.MarkOutboxConfirmed(cmd.DispatcherID, original.EventID); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
pending, err = restarted.ListPendingOutbox(cmd.DispatcherID)
|
||||
if err != nil || len(pending) != 1 || pending[0].EventType != "call.execute" {
|
||||
t.Fatalf("confirmed result was redelivered as a new event: events=%+v err=%v", pending, err)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user