From 65cef9ee30a6c490acfd3c9ed0575ad72ba274ec Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 5 Oct 2026 17:10:43 +0800 Subject: [PATCH] release capacity after verified pre-answer channel termination --- cmd/sip-go-agent/agent_real.go | 2 +- .../evidence/nonprod-formal-calls-20261005.md | 8 ++++ internal/asterisk/call.go | 43 +++++++++++++++-- internal/asterisk/call_test.go | 47 ++++++++++++++++++- internal/rpc/approved_recorded_real.go | 16 ++++++- internal/rpc/approved_recorded_real_test.go | 38 +++++++++++++++ 6 files changed, 146 insertions(+), 8 deletions(-) diff --git a/cmd/sip-go-agent/agent_real.go b/cmd/sip-go-agent/agent_real.go index f4fc50f..2ce2ccc 100644 --- a/cmd/sip-go-agent/agent_real.go +++ b/cmd/sip-go-agent/agent_real.go @@ -90,7 +90,7 @@ func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, d }).Prepare(execution) }, OnFailure: func(execution rpc.ApprovedExecution, cause error) error { - log.Printf("Agent real call requires inspection: event_id=%q task_id=%q native_phase=%q ari_http_status=%d cause_type=%T", execution.SourceEventID, execution.TaskID, asterisk.NativeCallPhase(cause), asterisk.NativeCallHTTPStatus(cause), cause) + log.Printf("Agent real call failure: event_id=%q task_id=%q native_phase=%q ari_http_status=%d confirmed_preanswer_end=%t cause_type=%T", execution.SourceEventID, execution.TaskID, asterisk.NativeCallPhase(cause), asterisk.NativeCallHTTPStatus(cause), asterisk.NativeCallConfirmedPreAnswerEnd(cause), cause) if !asterisk.NativeCallNeverSubmitted(cause) { return nil } diff --git a/docs/evidence/nonprod-formal-calls-20261005.md b/docs/evidence/nonprod-formal-calls-20261005.md index 20b48a1..720a0db 100644 --- a/docs/evidence/nonprod-formal-calls-20261005.md +++ b/docs/evidence/nonprod-formal-calls-20261005.md @@ -62,3 +62,11 @@ | `call-8bad49c7-5394-4576-b94f-baa9605ac20d` | 百应 / `1047` | INVITE→供应商 SIP 480→ACK;3 包、抓包退出 0;没有最终结果 | `8cbf1d04dc8dd2bcf47184ea95d3c9e1bfbf7da6d6b8cb274b613c729bcbf750` | `57e2ae884a5fb4a3beb38e51a5b89ed47a836ab8a486a23bcf4d99db2bd4a945` | **截至 2026-10-05 13:58 CST,矩阵仍缺一格:数企 / `4745`。** 它此前只经历 ARI 400 或任务过期,未发出真实 SIP;但当日该线路/号码的保守尝试登记已达到 **3/3**,不能清零计数或绕开门禁补测。其余五格均有真实 INVITE 和供应商 480,均未接通。新中鼎与百应事件仍为 `dispatched` 未知占用,两格再次用满额度 2/2;Agent 原有 7 条 `unknown` 日记保持原样,最终结果 outbox 仍为 0。四项服务 active、Asterisk 0 活动通话、抓包凭证 0。没有得到对**新事件**的人工释放授权,当前不再外呼。禁止自动等到次日拨号;剩余格须在新的有效窗口由使用者明确安排,并先按规则处理未知占用。六格未全部实测,**不作六格总判定**。 + +## 17:10 CST 不拨号定位:额度为何再次占满 + +最近两条 `call-644430bc-862a-47d7-9d50-cda36edc947a`、`call-8bad49c7-5394-4576-b94f-baa9605ac20d` 的 Agent 日志均在原生 ARI `answer_wait` 阶段报 HTTP 404;两份原始抓包分别显示真实 INVITE→480→ACK。只读向同一台 Asterisk 查询两条事件对应的通道均得 **HTTP 404**;当时活动通话为 0。这能证明当前无对应活动通道,但旧 Agent 未将结束事实写入可恢复日记,也未发送 `ReportEnded`,Dispatcher 不得凭后台单次查看自动清除历史未知占用。两次真实尝试各占一个 `dispatched` 名额,达到容量上限 2/2;“未知”是**无可核实的终结报告**,不是已确认仍在通话。 + +代码根因是原生通道在进入 Stasis 前结束时,拨号器只返回 `answer_wait` 错误,清理时 ARI Hangup 又可能得到 404;上层一律将该错误视作不确定,既不记录无录音的真实失败,也不提交 Agent 终结报告。新代码按使用者选择,只对**已经成功收到 ARI originate 响应且尚未接通**的同一通道,在清理之后再次查询、明确获得 Asterisk `GET channel` 的 404 时,记真实的未接通失败、写 Agent 恢复日记并报告 Dispatcher,随后由既有确认流程释放容量;查询失败或通道仍存在则保留占用,不通过定时过期或 SIP 480 本身释放。日志增加确认标志,不写密钥或音频。隔离测试同时覆盖查询返回 404、仍存在、查询失败、报告成功与无法证明时不报告;未进行任何新的拨号。 + +本机 `go vet ./...`、`go test -race ./...`(19 个通过的包)、`go build ./...`、业务覆盖率 69.1%、当前合同 35 项与 Proto 7 项来源哈希通过。`make check` 因本机缺少 `buf` 在 Proto lint 前中止;**不能声称完整门禁通过**。此处代码尚未部署到非生产主机,**不会追认或自动核销上述两条历史未知事件**。现存未知占用仍需针对事件逐条明确授权、备份和核查后人工处理;不得以新逻辑冒充过去的终结回执。 diff --git a/internal/asterisk/call.go b/internal/asterisk/call.go index a9c9bde..238c808 100644 --- a/internal/asterisk/call.go +++ b/internal/asterisk/call.go @@ -27,8 +27,9 @@ type NativeDial struct { // NativeCallFailure records a secret-free stage for a possibly ambiguous ARI // attempt. The wrapped cause remains available for programmatic inspection. type NativeCallFailure struct { - Phase string - Cause error + Phase string + Cause error + ConfirmedEnd bool // only after an originated pre-answer channel is verified absent from Asterisk } func (e *NativeCallFailure) Error() string { @@ -46,6 +47,14 @@ func NativeCallPhase(err error) string { return "unclassified" } +// NativeCallConfirmedPreAnswerEnd requires a successful ARI originate followed +// by a read of the same channel returning 404 after cleanup. A hangup request, +// timeout, or failed status read alone is not termination evidence. +func NativeCallConfirmedPreAnswerEnd(err error) bool { + var failure *NativeCallFailure + return errors.As(err, &failure) && failure.Phase == "answer_wait" && failure.ConfirmedEnd +} + // NativeCallNeverSubmitted is true only when the error occurred before any // ARI originate request could have been sent. After submission, even a local // error or successful cleanup does not prove the carrier never received SIP. @@ -62,8 +71,16 @@ func NativeCallNeverSubmitted(err error) bool { // a provider response body, credential or request URL in diagnostic logs. func NativeCallHTTPStatus(err error) int { var response native.RequestError - if errors.As(err, &response) { - return response.Code() + for err != nil { + if errors.As(err, &response) { + return response.Code() + } + // Channel.Data wraps the HTTP error with a Cause() wrapper in ari/v5. + var caused interface{ Cause() error } + if !errors.As(err, &caused) { + break + } + err = caused.Cause() } return 0 } @@ -76,6 +93,8 @@ type NativeCall struct { cancel context.CancelCauseFunc monitorDone chan struct{} destroyed atomic.Bool + verifyEnd bool + confirmedEnd bool client ari.Client subscription ari.Subscription outbound *ari.ChannelHandle @@ -102,7 +121,12 @@ func dialWithClient(ctx context.Context, client ari.Client, request NativeDial) call := &NativeCall{client: client} defer func() { if err != nil { - err = errors.Join(err, call.Close()) + cleanupErr := call.Close() + var failure *NativeCallFailure + if errors.As(err, &failure) && failure.Phase == "answer_wait" && call.confirmedEnd { + failure.ConfirmedEnd = true + } + err = errors.Join(err, cleanupErr) } }() if ctx == nil || client == nil { @@ -134,6 +158,7 @@ func dialWithClient(ctx context.Context, client ari.Client, request NativeDial) answerCtx, cancel := context.WithTimeout(ctx, request.AnswerTimeout) defer cancel() if err := awaitStasisStart(answerCtx, call.subscription, request.ExecutionID); err != nil { + call.verifyEnd = true // Originate returned successfully; no bridge or media channel exists yet. return nil, &NativeCallFailure{Phase: "answer_wait", Cause: err} } bridgeKey := ari.NewKey(ari.BridgeKey, request.ExecutionID+"-bridge") @@ -230,6 +255,14 @@ func (c *NativeCall) Close() error { if c.outbound != nil && !c.destroyed.Load() { c.closeErr = errors.Join(c.closeErr, c.outbound.Hangup()) } + if c.verifyEnd && c.outbound != nil { + _, err := c.outbound.Data() + if NativeCallHTTPStatus(err) == 404 { + c.confirmedEnd = true + } else if err != nil { + c.closeErr = errors.Join(c.closeErr, fmt.Errorf("native outbound channel end confirmation unavailable: %w", err)) + } + } if c.bridge != nil { c.closeErr = errors.Join(c.closeErr, c.bridge.Delete()) } diff --git a/internal/asterisk/call_test.go b/internal/asterisk/call_test.go index e3d51bc..4ba553f 100644 --- a/internal/asterisk/call_test.go +++ b/internal/asterisk/call_test.go @@ -41,6 +41,9 @@ type testChannels struct { issued int originateErr error externalErr error + dataErr error + hangupErr error + dataChecks int hungup []string mediaPeerIP string } @@ -76,7 +79,14 @@ func (c *testChannels) GetVariable(_ *ari.Key, name string) (string, error) { } func (c *testChannels) Hangup(key *ari.Key, _ string) error { c.hungup = append(c.hungup, key.ID) - return nil + return c.hangupErr +} +func (c *testChannels) Data(key *ari.Key) (*ari.ChannelData, error) { + c.dataChecks++ + if c.dataErr != nil { + return nil, c.dataErr + } + return &ari.ChannelData{ID: key.ID}, nil } type testBridges struct { @@ -232,6 +242,41 @@ func TestNativeCallUnknownBridgeOrMediaCreationCleansKnownIdentities(t *testing. } } +type channelDataCause struct{ err error } + +func (e channelDataCause) Error() string { return "channel read failed" } +func (e channelDataCause) Cause() error { return e.err } + +func TestNativeCallHTTPStatusReadsARICauseWrapper(t *testing.T) { + err := fmt.Errorf("channel read: %w", channelDataCause{err: nativeHTTPStatusError{code: 404}}) + if got := NativeCallHTTPStatus(err); got != 404 { + t.Fatalf("ARI channel data 404 was hidden by Cause wrapper: got %d", got) + } +} + +func TestNativeCallPreAnswerReleaseRequiresMissingOriginatedChannel(t *testing.T) { + for _, tc := range []struct { + name string + dataErr error + confirmed bool + }{ + {"absent", channelDataCause{err: nativeHTTPStatusError{code: 404}}, true}, + {"still active", nil, false}, + {"query failed", nativeHTTPStatusError{code: 503}, false}, + } { + t.Run(tc.name, func(t *testing.T) { + channels := &testChannels{dataErr: tc.dataErr, hangupErr: nativeHTTPStatusError{code: 404}} + client := &testARIClient{channels: channels, bridges: &testBridges{}, bus: &testBus{sub: answerEvents{events: make(chan ari.Event)}}} + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond) + defer cancel() + _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118}) + if err == nil || NativeCallPhase(err) != "answer_wait" || NativeCallConfirmedPreAnswerEnd(err) != tc.confirmed || channels.dataChecks != 1 || channels.issued != 1 || client.closed != 1 { + t.Fatalf("only the confirmed missing channel can release capacity: err=%v confirmed=%v checks=%d issued=%d closed=%d", err, NativeCallConfirmedPreAnswerEnd(err), channels.dataChecks, channels.issued, client.closed) + } + }) + } +} + func TestNativeCallAnswerTimeoutNeverCreatesFakeMedia(t *testing.T) { client := &testARIClient{ channels: &testChannels{mediaPeerIP: "127.0.0.1"}, diff --git a/internal/rpc/approved_recorded_real.go b/internal/rpc/approved_recorded_real.go index 2cb6211..4bbb6ca 100644 --- a/internal/rpc/approved_recorded_real.go +++ b/internal/rpc/approved_recorded_real.go @@ -95,9 +95,23 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con return err } } + attemptedAt := time.Now().UTC() call, err := originate(ctx, dial) if err != nil { - return err // unknown origination is never retried or reported as a completed call + if !asterisk.NativeCallConfirmedPreAnswerEnd(err) { + return err // uncertain origination or end is never retried or reported as completed + } + payload, payloadErr := callflow.FinalResultPayload(callflow.FinalCallFacts{ + TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID, + Callee: approved.Callee, TrunkID: approved.SelectedTrunkID, + StartedAt: attemptedAt, EndedAt: time.Now().UTC(), Outcome: "failed", ReasonMessage: "call ended before answer", + }, callflow.Result{}) + if payloadErr != nil { + return errors.Join(err, payloadErr) + } + reportCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), r.ReportTimeout) + defer cancel() + return errors.Join(err, r.Delivery.Complete(reportCtx, agent.CompletedRecording{ResultPayload: payload})) } startedAt := time.Now().UTC() // Originate returns only after the actual channel entered Stasis. // A real channel may have answered even if media/AI setup subsequently diff --git a/internal/rpc/approved_recorded_real_test.go b/internal/rpc/approved_recorded_real_test.go index 7445a05..7cb37c4 100644 --- a/internal/rpc/approved_recorded_real_test.go +++ b/internal/rpc/approved_recorded_real_test.go @@ -154,6 +154,44 @@ func TestApprovedRecordedRealCallDoesNotReportUnknownHangup(t *testing.T) { } } +func TestApprovedRecordedRealCallReportsVerifiedPreAnswerEndWithoutRecording(t *testing.T) { + fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second) + approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second + originate := 0 + runner := &ApprovedRecordedRealCall{ + MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second, Delivery: fixture.Delivery, + originator: func(context.Context, asterisk.NativeDial) (realMedia, error) { + originate++ + return realMedia{}, &asterisk.NativeCallFailure{Phase: "answer_wait", ConfirmedEnd: true, Cause: errors.New("no answer")} + }, + } + if err := runner.Run(context.Background(), approved); err == nil || originate != 1 || strings.Join(stub.calls, ",") != "end,result" || puts.Load() != 0 { + t.Fatalf("confirmed pre-answer end must report actual failure once without recording: err=%v originate=%d calls=%v puts=%d", err, originate, stub.calls, puts.Load()) + } + var result struct { + Outcome string `json:"outcome"` + Recording map[string]any `json:"recording"` + ReasonMessage string `json:"reason_message"` + } + if err := json.Unmarshal(stub.result, &result); err != nil || result.Outcome != "failed" || len(result.Recording) != 0 || result.ReasonMessage != "call ended before answer" { + t.Fatalf("pre-answer SIP rejection must not be recorded as answered: result=%+v err=%v", result, err) + } +} + +func TestApprovedRecordedRealCallDoesNotReportUnverifiedPreAnswerEnd(t *testing.T) { + fixture, stub, _, _, approved := recordedMockFixture(t, nil, time.Second) + approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second + runner := &ApprovedRecordedRealCall{ + MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second, Delivery: fixture.Delivery, + originator: func(context.Context, asterisk.NativeDial) (realMedia, error) { + return realMedia{}, &asterisk.NativeCallFailure{Phase: "answer_wait", Cause: errors.New("channel state unavailable")} + }, + } + if err := runner.Run(context.Background(), approved); err == nil || len(stub.calls) != 0 { + t.Fatalf("unverified call end must remain occupied: err=%v calls=%v", err, stub.calls) + } +} + func TestApprovedRecordedRealCallCannotOriginateTwiceOrReportUnknownOrigination(t *testing.T) { fixture, stub, _, _, approved := recordedMockFixture(t, nil, time.Second) approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second