diff --git a/internal/agent/recording_delivery_test.go b/internal/agent/recording_delivery_test.go index a0dfa89..08f78c3 100644 --- a/internal/agent/recording_delivery_test.go +++ b/internal/agent/recording_delivery_test.go @@ -72,7 +72,7 @@ func testRecordingDelivery(stub *recordingDeliveryRPC) *RecordingDelivery { func TestRecordingDeliveryNoRecordingReportsOnlyAfterConfirmedEnd(t *testing.T) { stub := &recordingDeliveryRPC{} delivery := testRecordingDelivery(stub) - payload := []byte(`{"task_id":"task-asr","caller_profile_id":"caller-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:05Z","duration_ms":5000,"outcome":"no_answer","reason_code":480,"reason_message":"no answer","transcript":[],"opt_out":false,"recording":{}}`) + payload := []byte(`{"task_id":"task-asr","caller_profile_id":"caller-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:05Z","duration_ms":5000,"outcome":"no_answer","reason_code":480,"reason_message":"no answer","status_line":null,"raw":null,"sip_capture_error":null,"transcript":[],"opt_out":false,"recording":{}}`) if err := delivery.Complete(context.Background(), CompletedRecording{ResultPayload: payload}); err != nil { t.Fatal(err) } @@ -209,7 +209,7 @@ func testDirectRecordingDelivery(t *testing.T, putStatus int) (*RecordingDeliver } delivery.Recovery = &RecordingRecovery{Root: root, Upload: UploadClient{AllowInsecureHTTP: true}} completed := CompletedRecording{ - ResultPayload: []byte(`{"task_id":"task-asr","caller_profile_id":"caller-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:20Z","duration_ms":20000,"outcome":"answered","reason_code":null,"reason_message":"answered","transcript":[],"opt_out":false,"recording":{}}`), + ResultPayload: []byte(`{"task_id":"task-asr","caller_profile_id":"caller-mock","callee":"15003164745","trunk_id":"trunk-mock","started_at":"2026-09-21T01:30:00Z","ended_at":"2026-09-21T01:30:20Z","duration_ms":20000,"outcome":"answered","reason_code":null,"reason_message":"answered","status_line":null,"raw":null,"sip_capture_error":null,"transcript":[],"opt_out":false,"recording":{}}`), Expected: true, RecordingID: "recording-mock", UploadID: "upload-mock", WAV: wav, DurationMS: durationMS, } return delivery, stub, completed, &puts, root diff --git a/internal/asterisk/answer.go b/internal/asterisk/answer.go index d44806d..57408dd 100644 --- a/internal/asterisk/answer.go +++ b/internal/asterisk/answer.go @@ -4,54 +4,52 @@ import ( "context" "errors" "fmt" - "strings" "github.com/CyCoreSystems/ari/v5" ) -// awaitStasisStart only marks the exact originated channel answered after it -// enters our ARI application. A generic Up event is not proof of media setup. -func awaitStasisStart(ctx context.Context, subscription ari.Subscription, channelID string) error { - if ctx == nil || subscription == nil || channelID == "" { - return errors.New("ARI answer subscription and channel identity required") +// awaitDialUp only runs after Dial was submitted. Create enters Stasis before +// any SIP request is sent, so StasisStart alone can never prove an answer. +func awaitDialUp(ctx context.Context, sub ari.Subscription, channelID string) error { + if ctx == nil || sub == nil || channelID == "" { + return errors.New("native answer wait has no channel subscription") } - recent := make([]string, 0, 8) + stasis, up := false, false for { select { case <-ctx.Done(): - return fmt.Errorf("ARI answer deadline or cancellation (recent_events=%s): %w", strings.Join(recent, ","), ctx.Err()) - case event, open := <-subscription.Events(): - if !open { - return errors.New("ARI event subscription closed before answer") + return ctx.Err() + case event, ok := <-sub.Events(): + if !ok { + return errors.New("native answer subscription closed before confirmed answer") } if event == nil { - return errors.New("ARI answer subscription delivered an empty event") - } - typ := event.GetType() - if len(recent) == cap(recent) { - recent = recent[1:] - } - recent = append(recent, typ) - if !eventBelongsToChannel(event, channelID) { continue } - switch typ { - case "StasisStart": - if _, ok := event.(*ari.StasisStart); ok { - return nil + switch e := event.(type) { + case *ari.StasisStart: + if e.Channel.ID == channelID { + stasis = true } - case "StasisEnd": - return errors.New("ARI channel left the application before StasisStart") - case "ChannelHangupRequest": - if hangup, ok := event.(*ari.ChannelHangupRequest); ok { - return fmt.Errorf("ARI channel ended before StasisStart: hangup_cause=%d", hangup.Cause) + case *ari.ChannelStateChange: + if e.Channel.ID == channelID && e.Channel.State == "Up" { + up = true } - return errors.New("ARI channel ended before StasisStart: hangup request") - case "ChannelDestroyed": - if destroyed, ok := event.(*ari.ChannelDestroyed); ok { - return fmt.Errorf("ARI channel ended before StasisStart: destroy_cause=%d", destroyed.Cause) + case *ari.ChannelDestroyed: + if e.Channel.ID == channelID { + return fmt.Errorf("native outbound channel destroyed before verified answer: cause=%d", e.Cause) } - return errors.New("ARI channel ended before StasisStart: destroyed") + case *ari.ChannelHangupRequest: + if e.Channel.ID == channelID { + return fmt.Errorf("native outbound channel requested hangup before verified answer: cause=%d", e.Cause) + } + case *ari.StasisEnd: + if e.Channel.ID == channelID { + return errors.New("native outbound channel left Stasis before verified answer") + } + } + if stasis && up { + return nil } } } diff --git a/internal/asterisk/answer_test.go b/internal/asterisk/answer_test.go index 8498ada..ee31965 100644 --- a/internal/asterisk/answer_test.go +++ b/internal/asterisk/answer_test.go @@ -2,6 +2,7 @@ package asterisk import ( "context" + "errors" "strings" "testing" "time" @@ -17,61 +18,44 @@ type answerEvents struct { func (f answerEvents) Events() <-chan ari.Event { return f.events } func (f answerEvents) Cancel() {} -func TestAwaitNativeAnswerRequiresMatchingStasisStart(t *testing.T) { - sub := answerEvents{events: make(chan ari.Event, 3)} - sub.events <- &ari.StasisStart{EventData: ari.EventData{Type: "StasisStart"}, Channel: ari.ChannelData{ID: "someone-else"}} - sub.events <- &ari.ChannelStateChange{EventData: ari.EventData{Type: "ChannelStateChange"}, Channel: ari.ChannelData{ID: "exec-1", State: "Up"}} - sub.events <- &ari.StasisStart{EventData: ari.EventData{Type: "StasisStart"}, Channel: ari.ChannelData{ID: "exec-1"}} - ctx, cancel := context.WithTimeout(context.Background(), time.Second) - defer cancel() - if err := awaitStasisStart(ctx, sub, "exec-1"); err != nil { - t.Fatalf("actual answered channel must enter the approved Stasis application: %v", err) - } -} - -func TestAwaitNativeAnswerFailsClosedOnHangupOrEventLoss(t *testing.T) { +func TestAwaitDialUpRequiresStasisAndPostDialUp(t *testing.T) { + const id = "execution-1" for _, tc := range []struct { - name string - event ari.Event + name string + events []ari.Event + answer bool }{ - {"hangup", &ari.ChannelHangupRequest{EventData: ari.EventData{Type: "ChannelHangupRequest"}, Channel: ari.ChannelData{ID: "exec-1"}, Cause: 17}}, - {"destroyed", &ari.ChannelDestroyed{EventData: ari.EventData{Type: "ChannelDestroyed"}, Channel: ari.ChannelData{ID: "exec-1"}, Cause: 17}}, - {"stasis-ended", &ari.StasisEnd{EventData: ari.EventData{Type: "StasisEnd"}, Channel: ari.ChannelData{ID: "exec-1"}}}, + {"created_not_answered", []ari.Event{&ari.StasisStart{Channel: ari.ChannelData{ID: id}}}, false}, + {"up_without_stasis", []ari.Event{&ari.ChannelStateChange{Channel: ari.ChannelData{ID: id, State: "Up"}}}, false}, + {"wrong_channel", []ari.Event{&ari.StasisStart{Channel: ari.ChannelData{ID: "other"}}, &ari.ChannelStateChange{Channel: ari.ChannelData{ID: "other", State: "Up"}}}, false}, + {"answer", []ari.Event{&ari.StasisStart{Channel: ari.ChannelData{ID: id}}, &ari.ChannelStateChange{Channel: ari.ChannelData{ID: id, State: "Up"}}}, true}, + {"answer_reordered", []ari.Event{&ari.ChannelStateChange{Channel: ari.ChannelData{ID: id, State: "Up"}}, &ari.StasisStart{Channel: ari.ChannelData{ID: id}}}, true}, + {"destroyed_before_answer", []ari.Event{&ari.StasisStart{Channel: ari.ChannelData{ID: id}}, &ari.ChannelDestroyed{Channel: ari.ChannelData{ID: id}}}, false}, } { t.Run(tc.name, func(t *testing.T) { - sub := answerEvents{events: make(chan ari.Event, 1)} - sub.events <- tc.event - ctx, cancel := context.WithTimeout(context.Background(), time.Second) + ch := make(chan ari.Event, len(tc.events)) + for _, event := range tc.events { + ch <- event + } + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond) defer cancel() - if err := awaitStasisStart(ctx, sub, "exec-1"); err == nil || !strings.Contains(err.Error(), "before StasisStart") { - t.Fatalf("failed origination must not count as answered: %v", err) + err := awaitDialUp(ctx, answerEvents{events: ch}, id) + if (err == nil) != tc.answer { + t.Fatalf("answer proof mismatch: %v", err) } }) } - t.Run("empty-event", func(t *testing.T) { - sub := answerEvents{events: make(chan ari.Event, 1)} - sub.events <- nil - ctx, cancel := context.WithTimeout(context.Background(), time.Second) - defer cancel() - if err := awaitStasisStart(ctx, sub, "exec-1"); err == nil || !strings.Contains(err.Error(), "empty event") { - t.Fatalf("empty ARI event cannot be ignored as if the stream was healthy: %v", err) - } - }) - t.Run("subscription-lost", func(t *testing.T) { - sub := answerEvents{events: make(chan ari.Event)} - close(sub.events) - ctx, cancel := context.WithTimeout(context.Background(), time.Second) - defer cancel() - if err := awaitStasisStart(ctx, sub, "exec-1"); err == nil || !strings.Contains(err.Error(), "subscription closed") { - t.Fatalf("event loss cannot be interpreted as answer: %v", err) - } - }) - t.Run("deadline", func(t *testing.T) { - sub := answerEvents{events: make(chan ari.Event)} - ctx, cancel := context.WithCancel(context.Background()) - cancel() - if err := awaitStasisStart(ctx, sub, "exec-1"); err == nil || !strings.Contains(err.Error(), "deadline or cancellation") { - t.Fatalf("expired answer window cannot be interpreted as answer: %v", err) - } - }) +} + +func TestAwaitDialUpFailsClosedOnLostEventsOrCancellation(t *testing.T) { + ch := make(chan ari.Event) + close(ch) + if err := awaitDialUp(context.Background(), answerEvents{events: ch}, "execution-1"); err == nil || !strings.Contains(err.Error(), "closed") { + t.Fatalf("subscription loss was hidden: %v", err) + } + ctx, cancel := context.WithCancel(context.Background()) + cancel() + if err := awaitDialUp(ctx, answerEvents{events: make(chan ari.Event)}, "execution-1"); !errors.Is(err, context.Canceled) { + t.Fatalf("cancellation was hidden: %v", err) + } } diff --git a/internal/asterisk/call.go b/internal/asterisk/call.go index 238c808..34bdc4e 100644 --- a/internal/asterisk/call.go +++ b/internal/asterisk/call.go @@ -7,6 +7,7 @@ import ( "net" "net/netip" "strconv" + "strings" "sync" "sync/atomic" "time" @@ -22,6 +23,7 @@ type NativeDial struct { ExecutionID, TrunkID, DialedCallee, CallerID string AnswerTimeout time.Duration MediaPayloadType uint8 + BindSIPCall func(executionID, callID string) error // before Dial; no number/time matching } // NativeCallFailure records a secret-free stage for a possibly ambiguous ARI @@ -29,7 +31,7 @@ type NativeDial struct { type NativeCallFailure struct { Phase string Cause error - ConfirmedEnd bool // only after an originated pre-answer channel is verified absent from Asterisk + ConfirmedEnd bool // only after a dialed pre-answer channel is verified absent from Asterisk } func (e *NativeCallFailure) Error() string { @@ -47,7 +49,7 @@ func NativeCallPhase(err error) string { return "unclassified" } -// NativeCallConfirmedPreAnswerEnd requires a successful ARI originate followed +// NativeCallConfirmedPreAnswerEnd requires a submitted native Dial 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 { @@ -55,12 +57,11 @@ func NativeCallConfirmedPreAnswerEnd(err error) bool { 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. +// NativeCallNeverSubmitted is true only when no ARI Dial request was sent. +// A failed Dial request is uncertain even if the subsequent cleanup succeeded. func NativeCallNeverSubmitted(err error) bool { switch NativeCallPhase(err) { - case "validate", "ari_open", "rtp_listen", "subscribe": + case "validate", "ari_open", "rtp_listen", "subscribe", "create", "caller_id", "sip_call_id", "sip_bind": return true default: return false @@ -137,28 +138,45 @@ func dialWithClient(ctx context.Context, client ari.Client, request NativeDial) if err != nil { return nil, &NativeCallFailure{Phase: "validate", Cause: err} } + if request.BindSIPCall == nil { + return nil, &NativeCallFailure{Phase: "validate", Cause: errors.New("HEP SIP call binding is required before native Dial")} + } call.Media, err = media.ListenRTPWithFormat("127.0.0.1:0", request.MediaPayloadType, media.FormatSLIN16, 16000) if err != nil { return nil, &NativeCallFailure{Phase: "rtp_listen", Cause: err} } key := ari.NewKey(ari.ChannelKey, request.ExecutionID) - call.subscription = client.Bus().Subscribe(key, "StasisStart", "StasisEnd", "ChannelHangupRequest", "ChannelDestroyed") + call.subscription = client.Bus().Subscribe(key, "StasisStart", "ChannelStateChange", "StasisEnd", "ChannelHangupRequest", "ChannelDestroyed") if call.subscription == nil { return nil, &NativeCallFailure{Phase: "subscribe", Cause: errors.New("native ARI channel subscription unavailable before origination")} } - // The reference key is an existing originator, not the new channel ID. - // ChannelID in originate already fixes this call's identity. - call.outbound, err = client.Channel().Originate(nil, originate) + // Create joins Stasis without sending SIP. The Asterisk channel's own + // PJSIP Call-ID is available before Dial, even for immediate SIP 480. + call.outbound, err = client.Channel().Create(nil, ari.ChannelCreateRequest{ + ChannelID: originate.ChannelID, Endpoint: originate.Endpoint, App: originate.App, + }) if err != nil { - // The HTTP response can be lost after Asterisk has created the - // channel. Do not originate again: try to stop the same identity. - call.outbound = client.Channel().Get(key) - return nil, &NativeCallFailure{Phase: "originate", Cause: err} + call.outbound = client.Channel().Get(key) // HTTP outcome may be lost; no Dial was submitted. + return nil, &NativeCallFailure{Phase: "create", Cause: err} + } + if err := call.outbound.SetVariable("CALLERID(num)", originate.CallerID); err != nil { + return nil, &NativeCallFailure{Phase: "caller_id", Cause: err} + } + callID, err := call.outbound.GetVariable("CHANNEL(pjsip,call-id)") + if err != nil || strings.TrimSpace(callID) == "" { + return nil, &NativeCallFailure{Phase: "sip_call_id", Cause: errors.Join(err, errors.New("Asterisk did not provide a pre-Dial SIP Call-ID"))} + } + if err := request.BindSIPCall(request.ExecutionID, callID); err != nil { + return nil, &NativeCallFailure{Phase: "sip_bind", Cause: err} + } + if err := call.outbound.Dial("", time.Duration(originate.Timeout)*time.Second); err != nil { + // The Dial HTTP reply may be lost after SIP was sent. Never Dial twice. + return nil, &NativeCallFailure{Phase: "dial", Cause: err} } 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. + if err := awaitDialUp(answerCtx, call.subscription, request.ExecutionID); err != nil { + call.verifyEnd = true // Created channel has no ExternalMedia or bridge yet. return nil, &NativeCallFailure{Phase: "answer_wait", Cause: err} } bridgeKey := ari.NewKey(ari.BridgeKey, request.ExecutionID+"-bridge") diff --git a/internal/asterisk/call_test.go b/internal/asterisk/call_test.go index 4ba553f..6112fd2 100644 --- a/internal/asterisk/call_test.go +++ b/internal/asterisk/call_test.go @@ -34,32 +34,57 @@ func (b *testBus) Subscribe(_ *ari.Key, _ ...string) ari.Subscription { return b type testChannels struct { ari.Channel - events chan ari.Event - originate ari.OriginateRequest - originateRef *ari.Key - media ari.ExternalMediaOptions - issued int - originateErr error - externalErr error - dataErr error - hangupErr error - dataChecks int - hungup []string - mediaPeerIP string + events chan ari.Event + create ari.ChannelCreateRequest + createRef *ari.Key + media ari.ExternalMediaOptions + issued int + createErr error + dialErr error + variableErr error + externalErr error + dataErr error + hangupErr error + dataChecks int + hungup []string + mediaPeerIP string + callerID string + noAnswer bool } -func (c *testChannels) Originate(reference *ari.Key, request ari.OriginateRequest) (*ari.ChannelHandle, error) { - c.issued++ - c.originateRef = reference - c.originate = request - if c.originateErr != nil { - return nil, c.originateErr +func (c *testChannels) Create(reference *ari.Key, request ari.ChannelCreateRequest) (*ari.ChannelHandle, error) { + c.createRef, c.create = reference, request + if c.createErr != nil { + return nil, c.createErr } if c.events != nil { - c.events <- &ari.StasisStart{EventData: ari.EventData{Type: "StasisStart"}, Channel: ari.ChannelData{ID: request.ChannelID}} + select { + case c.events <- &ari.StasisStart{EventData: ari.EventData{Type: "StasisStart"}, Channel: ari.ChannelData{ID: request.ChannelID, State: "Down"}}: + default: + } } return ari.NewChannelHandle(ari.NewKey(ari.ChannelKey, request.ChannelID), c, nil), nil } +func (c *testChannels) Dial(key *ari.Key, _ string, _ time.Duration) error { + c.issued++ + if c.dialErr != nil { + return c.dialErr + } + if c.events != nil && !c.noAnswer { + select { + case c.events <- &ari.ChannelStateChange{EventData: ari.EventData{Type: "ChannelStateChange"}, Channel: ari.ChannelData{ID: key.ID, State: "Up"}}: + default: + } + } + return nil +} +func (c *testChannels) SetVariable(_ *ari.Key, name, value string) error { + if name != "CALLERID(num)" { + return errors.New("unexpected caller variable") + } + c.callerID = value + return nil +} func (c *testChannels) Get(key *ari.Key) *ari.ChannelHandle { return ari.NewChannelHandle(key, c, nil) } func (c *testChannels) ExternalMedia(key *ari.Key, options ari.ExternalMediaOptions) (*ari.ChannelHandle, error) { c.media = options @@ -70,6 +95,11 @@ func (c *testChannels) ExternalMedia(key *ari.Key, options ari.ExternalMediaOpti } func (c *testChannels) GetVariable(_ *ari.Key, name string) (string, error) { switch name { + case "CHANNEL(pjsip,call-id)": + if c.variableErr != nil { + return "", c.variableErr + } + return "isolated-sip-call-id", nil case "UNICASTRTP_LOCAL_ADDRESS": return c.mediaPeerIP, nil case "UNICASTRTP_LOCAL_PORT": @@ -137,7 +167,7 @@ func TestNativeCallFailureReportsStageAndHTTPStatusWithoutResponseBody(t *testin } } -func TestNativeCallUsesOneARIOriginateAndActualRTPBridge(t *testing.T) { +func TestNativeCallBindsSIPCallIDBeforeOneDialAndActualRTPBridge(t *testing.T) { events := make(chan ari.Event, 2) client := &testARIClient{ channels: &testChannels{events: events, mediaPeerIP: "127.0.0.1"}, @@ -146,16 +176,21 @@ func TestNativeCallUsesOneARIOriginateAndActualRTPBridge(t *testing.T) { } ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() + boundBeforeDial := false call, err := dialWithClient(ctx, client, NativeDial{ ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118, + BindSIPCall: func(executionID, callID string) error { + boundBeforeDial = executionID == "exec-1" && callID == "isolated-sip-call-id" && client.channels.issued == 0 + return nil + }, }) if err != nil { t.Fatal(err) } - if call.Media == nil || client.channels.issued != 1 || client.channels.originate.Endpoint != "PJSIP/708915000000001@shuqi" || - client.channels.originate.CallerID != "BD1234" || client.channels.originate.ChannelID != "exec-1" || - client.channels.originateRef != nil || client.channels.originate.Originator != "" || + if call.Media == nil || !boundBeforeDial || client.channels.issued != 1 || client.channels.create.Endpoint != "PJSIP/708915000000001@shuqi" || + client.channels.callerID != "BD1234" || client.channels.create.ChannelID != "exec-1" || + client.channels.createRef != nil || client.channels.create.Originator != "" || client.channels.media.Format != "slin16" || client.channels.media.App != "go-sip-agent" || len(client.bridges.attached) != 2 || client.bridges.attached[0] != "exec-1" { t.Fatal("native channel/ExternalMedia was not attached exactly once to the approved bridge") @@ -168,12 +203,14 @@ func TestNativeCallUsesOneARIOriginateAndActualRTPBridge(t *testing.T) { } } +func testBindSIPCall(string, string) error { return nil } + func TestNativeCallStopsMediaOnMatchingCarrierHangup(t *testing.T) { events := make(chan ari.Event, 3) client := &testARIClient{channels: &testChannels{events: events, mediaPeerIP: "127.0.0.1"}, bridges: &testBridges{}, bus: &testBus{sub: answerEvents{events: events}}} ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() - call, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118}) + call, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118, BindSIPCall: testBindSIPCall}) if err != nil { t.Fatal(err) } @@ -200,13 +237,13 @@ func TestNativeCallStopsMediaOnMatchingCarrierHangup(t *testing.T) { func TestNativeCallUnknownOriginationMustAttemptHangupWithoutRetry(t *testing.T) { client := &testARIClient{ - channels: &testChannels{originateErr: errors.New("connection lost during originate")}, + channels: &testChannels{dialErr: errors.New("connection lost during dial")}, bridges: &testBridges{}, bus: &testBus{sub: answerEvents{events: make(chan ari.Event)}}, } ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() - if _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118}); err == nil || !strings.Contains(err.Error(), "outcome unknown") || NativeCallPhase(err) != "originate" { + if _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118, BindSIPCall: testBindSIPCall}); err == nil || !strings.Contains(err.Error(), "outcome unknown") || NativeCallPhase(err) != "dial" { t.Fatalf("unknown originate outcome must never be reported as a completed call: %v", err) } if client.channels.issued != 1 || len(client.channels.hungup) != 1 || client.channels.hungup[0] != "exec-1" || client.closed != 1 { @@ -217,7 +254,7 @@ func TestNativeCallUnknownOriginationMustAttemptHangupWithoutRetry(t *testing.T) func TestNativeCallUnknownBridgeOrMediaCreationCleansKnownIdentities(t *testing.T) { for _, step := range []string{"bridge", "external-media"} { t.Run(step, func(t *testing.T) { - events := make(chan ari.Event, 1) + events := make(chan ari.Event, 2) channels := &testChannels{events: events, mediaPeerIP: "127.0.0.1"} bridges := &testBridges{} if step == "bridge" { @@ -228,7 +265,7 @@ func TestNativeCallUnknownBridgeOrMediaCreationCleansKnownIdentities(t *testing. client := &testARIClient{channels: channels, bridges: bridges, bus: &testBus{sub: answerEvents{events: events}}} ctx, cancel := context.WithTimeout(context.Background(), time.Second) defer cancel() - if _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118}); err == nil || NativeCallPhase(err) != map[string]string{"bridge": "bridge_create", "external-media": "external_media_create"}[step] { + if _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118, BindSIPCall: testBindSIPCall}); err == nil || NativeCallPhase(err) != map[string]string{"bridge": "bridge_create", "external-media": "external_media_create"}[step] { t.Fatalf("unknown bridge/media outcome needs its exact failure phase: %v", err) } wantHungup := 1 @@ -265,11 +302,11 @@ func TestNativeCallPreAnswerReleaseRequiresMissingOriginatedChannel(t *testing.T {"query failed", nativeHTTPStatusError{code: 503}, false}, } { t.Run(tc.name, func(t *testing.T) { - channels := &testChannels{dataErr: tc.dataErr, hangupErr: nativeHTTPStatusError{code: 404}} + channels := &testChannels{dataErr: tc.dataErr, hangupErr: nativeHTTPStatusError{code: 404}, noAnswer: true} 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}) + _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118, BindSIPCall: testBindSIPCall}) 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) } @@ -279,13 +316,13 @@ func TestNativeCallPreAnswerReleaseRequiresMissingOriginatedChannel(t *testing.T func TestNativeCallAnswerTimeoutNeverCreatesFakeMedia(t *testing.T) { client := &testARIClient{ - channels: &testChannels{mediaPeerIP: "127.0.0.1"}, + channels: &testChannels{mediaPeerIP: "127.0.0.1", noAnswer: true}, bridges: &testBridges{}, bus: &testBus{sub: answerEvents{events: make(chan ari.Event)}}, } ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond) defer cancel() - if _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118}); err == nil || !strings.Contains(err.Error(), "deadline") { + if _, err := dialWithClient(ctx, client, NativeDial{ExecutionID: "exec-1", TrunkID: "shuqi", DialedCallee: "708915000000001", CallerID: "BD1234", AnswerTimeout: time.Second, MediaPayloadType: 118, BindSIPCall: testBindSIPCall}); err == nil || !strings.Contains(err.Error(), "deadline") { t.Fatalf("unanswered channel must fail rather than synthesize media: %v", err) } if client.channels.issued != 1 || len(client.channels.hungup) == 0 || client.closed != 1 { diff --git a/internal/rpc/recording_server_flow_test.go b/internal/rpc/recording_server_flow_test.go index a74467f..7efe8f1 100644 --- a/internal/rpc/recording_server_flow_test.go +++ b/internal/rpc/recording_server_flow_test.go @@ -112,6 +112,7 @@ func recordingResultPayload(t *testing.T, task configread.Snapshot, grant *agent "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", + "status_line": nil, "raw": nil, "sip_capture_error": nil, "transcript": []any{}, "opt_out": false, "recording": map[string]any{}, } if grant != nil { diff --git a/internal/store/result_test.go b/internal/store/result_test.go index 121776a..0bc76c5 100644 --- a/internal/store/result_test.go +++ b/internal/store/result_test.go @@ -19,6 +19,7 @@ func currentResultPayload(t *testing.T) []byte { "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", + "status_line": nil, "raw": nil, "sip_capture_error": nil, "transcript": []any{}, "opt_out": false, "recording": map[string]any{}, } raw, err := json.Marshal(payload)