release capacity after verified pre-answer channel termination

This commit is contained in:
2026-10-05 17:10:43 +08:00
parent 04e8383e2c
commit 65cef9ee30
6 changed files with 146 additions and 8 deletions
+1 -1
View File
@@ -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
}
@@ -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 前中止;**不能声称完整门禁通过**。此处代码尚未部署到非生产主机,**不会追认或自动核销上述两条历史未知事件**。现存未知占用仍需针对事件逐条明确授权、备份和核查后人工处理;不得以新逻辑冒充过去的终结回执。
+38 -5
View File
@@ -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())
}
+46 -1
View File
@@ -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"},
+15 -1
View File
@@ -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
@@ -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