From 0f7415e366907c8d8943dfc1220a824c7200bd9a Mon Sep 17 00:00:00 2001 From: Rogee Date: Wed, 30 Sep 2026 00:14:09 +0800 Subject: [PATCH] Persist ended calls racing originate acknowledgment --- .../saas-dispatcher-implementation.md | 1 + internal/store/current_calls.go | 92 ++++++++++++-- internal/store/current_calls_test.go | 115 ++++++++++++++++++ internal/store/current_result_test.go | 27 +++- 4 files changed, 226 insertions(+), 9 deletions(-) diff --git a/docs/evidence/saas-dispatcher-implementation.md b/docs/evidence/saas-dispatcher-implementation.md index 174e62f..a00abce 100644 --- a/docs/evidence/saas-dispatcher-implementation.md +++ b/docs/evidence/saas-dispatcher-implementation.md @@ -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 通过。 ## 验收台账 diff --git a/internal/store/current_calls.go b/internal/store/current_calls.go index 7588565..091b2f0 100644 --- a/internal/store/current_calls.go +++ b/internal/store/current_calls.go @@ -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 } diff --git a/internal/store/current_calls_test.go b/internal/store/current_calls_test.go index fc8f23e..d93effe 100644 --- a/internal/store/current_calls_test.go +++ b/internal/store/current_calls_test.go @@ -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") diff --git a/internal/store/current_result_test.go b/internal/store/current_result_test.go index 20e2eed..3f10765 100644 --- a/internal/store/current_result_test.go +++ b/internal/store/current_result_test.go @@ -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 {