diff --git a/internal/asterisk/answer.go b/internal/asterisk/answer.go new file mode 100644 index 0000000..1e31481 --- /dev/null +++ b/internal/asterisk/answer.go @@ -0,0 +1,65 @@ +package asterisk + +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") + } + recent := make([]string, 0, 8) + 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") + } + 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) + matched := false + for _, key := range event.Keys() { + if key != nil && key.Kind == "channel" && key.ID == channelID { + matched = true + break + } + } + if !matched { + continue + } + switch typ { + case "StasisStart": + if _, ok := event.(*ari.StasisStart); ok { + return nil + } + 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) + } + 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) + } + return errors.New("ARI channel ended before StasisStart: destroyed") + } + } + } +} diff --git a/internal/asterisk/answer_test.go b/internal/asterisk/answer_test.go new file mode 100644 index 0000000..7c638c3 --- /dev/null +++ b/internal/asterisk/answer_test.go @@ -0,0 +1,76 @@ +package asterisk + +import ( + "context" + "strings" + "testing" + "time" + + "github.com/CyCoreSystems/ari/v5" +) + +type answerEvents struct { + ari.Subscription + events chan ari.Event +} + +func (f answerEvents) Events() <-chan ari.Event { return f.events } + +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) { + for _, tc := range []struct { + name string + event ari.Event + }{ + {"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"}}}, + } { + 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) + 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) + } + }) + } + 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) + } + }) +}