bind SIP Call-ID before native ARI Dial and await real answer
This commit is contained in:
@@ -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
|
||||
|
||||
+31
-33
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
+34
-16
@@ -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")
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user