Persist ended calls racing originate acknowledgment

This commit is contained in:
2026-09-30 00:14:09 +08:00
parent 40eeed2ba6
commit 0f7415e366
4 changed files with 226 additions and 9 deletions
@@ -63,6 +63,7 @@
- Agent 失败恢复隔离组件:仅在 OSS PUT 明确失败后,把原 bucket/object_key 对应的录音和通话信息两文件写入私有目录并同步落盘;双文件缺失或损坏明确报错、不伪造结果。完成保存后固定 48 小时窗口,按 1 分钟递增至最长 1 小时重试;到期保留原文件。PUT 前持久写入 in-flight,结果不明或进程重启不会二次 PUT;确认上传后先持久记录成功,再经注入的 Mock 回报通话结果,回报失败/重启仅重发原结果。启动扫描识别遗漏文件并提供不暴露原路径的稳定摘要;真实 Dispatcher 授权及回报尚未接线。
- 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 身份验证。
- Dispatcher 的终结与外呼回执顺序竞争隔离修复:原流程在 Agent 接受执行的 RPC 返回后才写入外呼回执,快速结束或 RPC 超时可能先到;现以一次 SQLite 事务在确认通话已结束后补齐原回执并释放占用,未知执行仅在确认结束后释放。迟到的执行响应、超时和重复结束不会产生第二份回执;注入 outbox 写入失败保留原占用。并发竞争及结束后立即生成唯一最终结果有单元测试;真实 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 通过。
## 验收台账
+84 -8
View File
@@ -235,15 +235,66 @@ func (s *CurrentStore) MarkExecuteUnknown(dispatcherID, eventID string) error {
if err != nil {
return fmt.Errorf("persist unknown execution: %w", err)
}
return requireOneRow(result, "unknown execution")
if err := requireOneRow(result, "unknown execution"); err != nil {
// The Agent may have confirmed the real end while the originate RPC was
// still returning a timeout. Never turn a durable finished call back into
// an unknown occupancy or treat that late timeout as another originate.
var confirmed int
scanErr := s.db.QueryRow(`SELECT COUNT(*) FROM dispatcher_inbox i JOIN dispatcher_outbox o
ON o.dispatcher_id=i.dispatcher_id AND o.event_id=i.event_id AND o.event_type='call.execute'
WHERE i.dispatcher_id=? AND i.event_id=? AND i.status='finished'`, dispatcherID, eventID).Scan(&confirmed)
if scanErr != nil {
return errors.Join(err, fmt.Errorf("inspect call after unknown RPC: %w", scanErr))
}
if confirmed == 1 {
return nil
}
return err
}
return nil
}
// FinishExecute accepts only an authenticated, confirmed Agent end fact. An
// end can race the originate RPC acknowledgment, or resolve a previously
// unknown RPC outcome. In both cases the acknowledgment and occupancy release
// commit together; a missing outbox never frees capacity.
func (s *CurrentStore) FinishExecute(dispatcherID, eventID string) error {
result, err := s.db.Exec(`UPDATE dispatcher_inbox SET status='finished' WHERE dispatcher_id=? AND event_id=? AND status='dispatched'`, dispatcherID, eventID)
tx, err := s.db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
var status, issuedAt string
var tenantID int64
if err := tx.QueryRow(`SELECT status,tenant_id,issued_at FROM dispatcher_inbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, eventID).Scan(&status, &tenantID, &issuedAt); err != nil {
return fmt.Errorf("inspect confirmed call end: %w", err)
}
switch status {
case "dispatching", "unknown":
if err := insertExecuteAck(tx, dispatcherID, eventID, tenantID, issuedAt, map[string]any{"status": "dispatched"}); err != nil {
return err
}
case "dispatched", "finished":
if err := requireExecuteAck(tx, dispatcherID, eventID); err != nil {
return err
}
if status == "finished" {
return nil
}
default:
return fmt.Errorf("call %q cannot confirm an end from state %q", eventID, status)
}
result, err := tx.Exec(`UPDATE dispatcher_inbox SET status='finished' WHERE dispatcher_id=? AND event_id=? AND status=?`, dispatcherID, eventID, status)
if err != nil {
return fmt.Errorf("persist confirmed call end: %w", err)
}
return requireOneRow(result, "confirmed call end")
if err := requireOneRow(result, "confirmed call end"); err != nil {
return err
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit confirmed call end and acknowledgment: %w", err)
}
return nil
}
func (s *CurrentStore) RejectExecute(dispatcherID, eventID, reason string) error {
@@ -270,8 +321,28 @@ func (s *CurrentStore) closeExecuteWithAck(dispatcherID, eventID, expectedStatus
return fmt.Errorf("read call acknowledgment identity: %w", err)
}
if status != expectedStatus {
if expectedStatus == "dispatching" && nextStatus == "dispatched" && status == "finished" {
return requireExecuteAck(tx, dispatcherID, eventID)
}
return fmt.Errorf("call %q cannot transition from %s to %s", eventID, status, nextStatus)
}
result, err := tx.Exec(`UPDATE dispatcher_inbox SET status=? WHERE dispatcher_id=? AND event_id=? AND status=?`, nextStatus, dispatcherID, eventID, expectedStatus)
if err != nil {
return fmt.Errorf("persist call acknowledgment status: %w", err)
}
if err := requireOneRow(result, "call acknowledgment status"); err != nil {
return err
}
if err := insertExecuteAck(tx, dispatcherID, eventID, tenantID, issuedAt, payload); err != nil {
return err
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit call status and acknowledgment: %w", err)
}
return nil
}
func insertExecuteAck(tx *sql.Tx, dispatcherID, eventID string, tenantID int64, issuedAt string, payload any) error {
body, err := json.Marshal(struct {
EventID string `json:"event_id"`
EventType string `json:"event_type"`
@@ -290,14 +361,19 @@ func (s *CurrentStore) closeExecuteWithAck(dispatcherID, eventID, expectedStatus
if err != nil {
return err
}
if _, err := tx.Exec(`UPDATE dispatcher_inbox SET status=? WHERE dispatcher_id=? AND event_id=? AND status=?`, nextStatus, dispatcherID, eventID, expectedStatus); err != nil {
return fmt.Errorf("persist call acknowledgment status: %w", err)
}
if _, err := tx.Exec(`INSERT INTO dispatcher_outbox(dispatcher_id,event_id,event_type,routing_key,body) VALUES(?,?,?,?,?)`, dispatcherID, eventID, "call.execute", route.BindingKey, body); err != nil {
return fmt.Errorf("persist call acknowledgment outbox: %w", err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit call status and acknowledgment: %w", err)
return nil
}
func requireExecuteAck(tx *sql.Tx, dispatcherID, eventID string) error {
var eventType string
if err := tx.QueryRow(`SELECT event_type FROM dispatcher_outbox WHERE dispatcher_id=? AND event_id=?`, dispatcherID, eventID).Scan(&eventType); err != nil {
return fmt.Errorf("confirmed call %q has no durable acknowledgment: %w", eventID, err)
}
if eventType != "call.execute" {
return fmt.Errorf("confirmed call %q has conflicting acknowledgment type %q", eventID, eventType)
}
return nil
}
+115
View File
@@ -137,6 +137,121 @@ func TestCurrentExecuteQuotaCountsUnknownAndReleasesOnlyConfirmedEnd(t *testing.
}
}
func TestCurrentConfirmedEndBeforeOriginateAckIsDurableAndIdempotent(t *testing.T) {
s := preparedCurrentCallStore(t)
cmd := currentCall("call-ended-before-ack")
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.FinishExecute(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatalf("confirmed Agent end raced with originate acknowledgment: %v", err)
}
if err := s.MarkExecuteDispatched(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatalf("late originate acknowledgment was not idempotent: %v", err)
}
if err := s.MarkExecuteUnknown(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatalf("late RPC timeout obscured a confirmed call end: %v", err)
}
if err := s.FinishExecute(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatalf("Agent end-report retry was not idempotent: %v", err)
}
outbox, err := s.ListPendingOutbox(cmd.DispatcherID)
if err != nil || len(outbox) != 1 || outbox[0].EventID != cmd.EventID || !strings.Contains(string(outbox[0].Body), `"status":"dispatched"`) {
t.Fatalf("early end lost or duplicated original acknowledgment: outbox=%+v err=%v", outbox, err)
}
occupied, err := s.TrunkOccupancy(cmd.DispatcherID)
if err != nil || occupied["trunk-mock"] != 0 {
t.Fatalf("confirmed end held occupancy after acknowledgment: occupied=%+v err=%v", occupied, err)
}
if _, created, err := s.RecordExecute(cmd); err != nil || created {
t.Fatalf("restarted finished call could originate again: created=%t err=%v", created, err)
}
}
func TestCurrentConfirmedUnknownCallEndReleasesOnlyAfterDurableAck(t *testing.T) {
s := preparedCurrentCallStore(t)
cmd := currentCall("call-ended-after-unknown")
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.MarkExecuteUnknown(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatal(err)
}
if err := s.FinishExecute(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatalf("confirmed unknown call end did not release: %v", err)
}
outbox, err := s.ListPendingOutbox(cmd.DispatcherID)
if err != nil || len(outbox) != 1 || outbox[0].EventID != cmd.EventID {
t.Fatalf("confirmed unknown execution lacked one durable acknowledgment: outbox=%+v err=%v", outbox, err)
}
occupied, err := s.TrunkOccupancy(cmd.DispatcherID)
if err != nil || occupied["trunk-mock"] != 0 {
t.Fatalf("confirmed unknown call still occupied trunk: occupied=%+v err=%v", occupied, err)
}
}
func TestCurrentEarlyEndOutboxFailureRetainsUnknownOccupancy(t *testing.T) {
s := preparedCurrentCallStore(t)
cmd := currentCall("call-early-end-fault")
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.db.Exec(`CREATE TRIGGER fail_early_ack BEFORE INSERT ON dispatcher_outbox BEGIN SELECT RAISE(ABORT,'injected acknowledgment fault'); END`); err != nil {
t.Fatal(err)
}
if err := s.FinishExecute(cmd.DispatcherID, cmd.EventID); err == nil {
t.Fatal("acknowledgment write failure released the call")
}
occupied, err := s.TrunkOccupancy(cmd.DispatcherID)
if err != nil || occupied["trunk-mock"] != 1 {
t.Fatalf("failed acknowledgment lost the unknown occupancy: occupied=%+v err=%v", occupied, err)
}
if _, err := s.db.Exec(`DROP TRIGGER fail_early_ack`); err != nil {
t.Fatal(err)
}
if err := s.FinishExecute(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatalf("original confirmed end could not resume: %v", err)
}
}
func TestCurrentEndAndOriginateAckRaceNeverDuplicate(t *testing.T) {
for attempt := 0; attempt < 30; attempt++ {
s := preparedCurrentCallStore(t)
cmd := currentCall("call-concurrent-end")
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)
}
errorsCh := make(chan error, 2)
go func() { errorsCh <- s.FinishExecute(cmd.DispatcherID, cmd.EventID) }()
go func() { errorsCh <- s.MarkExecuteDispatched(cmd.DispatcherID, cmd.EventID) }()
for range 2 {
if err := <-errorsCh; err != nil {
t.Fatal(err)
}
}
outbox, err := s.ListPendingOutbox(cmd.DispatcherID)
if err != nil || len(outbox) != 1 {
t.Fatalf("race produced multiple acknowledgments: outbox=%+v err=%v", outbox, err)
}
occupied, err := s.TrunkOccupancy(cmd.DispatcherID)
if err != nil || occupied["trunk-mock"] != 0 {
t.Fatalf("race failed to release confirmed end: occupied=%+v err=%v", occupied, err)
}
}
}
func TestCurrentExecuteOutboxFailureDoesNotMarkOriginatedCallDelivered(t *testing.T) {
s := preparedCurrentCallStore(t)
cmd := currentCall("call-fault")
+26 -1
View File
@@ -103,11 +103,36 @@ func TestCurrentFinalResultNeedsConfirmedEndAndRemainsExactlyOne(t *testing.T) {
}
}
func TestCurrentFinalResultAfterEndRacedOriginateAck(t *testing.T) {
s := preparedCurrentCallStore(t)
cmd := currentCall("result-ended-before-ack")
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.FinishExecute(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatal(err)
}
result, created, err := s.RecordCallResult(cmd.DispatcherID, cmd.EventID, currentResultPayload(t))
if err != nil || !created {
t.Fatalf("confirmed early end could not produce its sole final result: created=%t err=%v", created, err)
}
if err := s.MarkExecuteDispatched(cmd.DispatcherID, cmd.EventID); err != nil {
t.Fatalf("late originate response conflicted with the confirmed result: %v", err)
}
outbox, err := s.ListPendingOutbox(cmd.DispatcherID)
if err != nil || len(outbox) != 2 || outbox[0].EventID != cmd.EventID || outbox[1].EventID != result.EventID {
t.Fatalf("early end lost its acknowledgment or produced duplicate results: outbox=%+v err=%v", outbox, err)
}
}
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"},
{"call_id", "other-call"}, {"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 {