246 lines
8.4 KiB
Go
246 lines
8.4 KiB
Go
package asterisk
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"net"
|
|
"net/netip"
|
|
"strconv"
|
|
"sync"
|
|
"sync/atomic"
|
|
"time"
|
|
|
|
"git.ipao.vip/rogee/go-sip/internal/media"
|
|
"github.com/CyCoreSystems/ari/v5"
|
|
"github.com/CyCoreSystems/ari/v5/client/native"
|
|
)
|
|
|
|
// NativeDial contains only the Dispatcher-selected SIP identity and the
|
|
// Agent's explicitly configured SLIN16 RTP payload profile.
|
|
type NativeDial struct {
|
|
ExecutionID, TrunkID, DialedCallee, CallerID string
|
|
AnswerTimeout time.Duration
|
|
MediaPayloadType uint8
|
|
}
|
|
|
|
// 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
|
|
}
|
|
|
|
func (e *NativeCallFailure) Error() string {
|
|
if errors.Is(e.Cause, context.DeadlineExceeded) {
|
|
return fmt.Sprintf("native ARI %s deadline: %T", e.Phase, e.Cause)
|
|
}
|
|
return fmt.Sprintf("native ARI %s outcome unknown: %T", e.Phase, e.Cause)
|
|
}
|
|
func (e *NativeCallFailure) Unwrap() error { return e.Cause }
|
|
func NativeCallPhase(err error) string {
|
|
var failure *NativeCallFailure
|
|
if errors.As(err, &failure) {
|
|
return failure.Phase
|
|
}
|
|
return "unclassified"
|
|
}
|
|
|
|
// 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.
|
|
func NativeCallNeverSubmitted(err error) bool {
|
|
switch NativeCallPhase(err) {
|
|
case "validate", "ari_open", "rtp_listen", "subscribe":
|
|
return true
|
|
default:
|
|
return false
|
|
}
|
|
}
|
|
|
|
// NativeCallHTTPStatus extracts a numeric ARI HTTP status without retaining
|
|
// 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()
|
|
}
|
|
return 0
|
|
}
|
|
|
|
// NativeCall owns one real ARI channel, bridge, ExternalMedia channel and RTP
|
|
// socket; closing it cannot originate or replay a second call.
|
|
type NativeCall struct {
|
|
Media *media.RTPStream
|
|
Context context.Context
|
|
cancel context.CancelCauseFunc
|
|
monitorDone chan struct{}
|
|
destroyed atomic.Bool
|
|
client ari.Client
|
|
subscription ari.Subscription
|
|
outbound *ari.ChannelHandle
|
|
external *ari.ChannelHandle
|
|
bridge *ari.BridgeHandle
|
|
closeOnce sync.Once
|
|
closeErr error
|
|
}
|
|
|
|
// Originate uses the issued Cell-local ARI credential. The caller must have
|
|
// already persisted the signed instruction and passed the host capture gate.
|
|
func (l Loader) Originate(ctx context.Context, request NativeDial) (*NativeCall, error) {
|
|
if _, err := approvedOriginateRequest(request.ExecutionID, request.TrunkID, request.DialedCallee, request.CallerID, request.AnswerTimeout); err != nil {
|
|
return nil, &NativeCallFailure{Phase: "validate", Cause: err}
|
|
}
|
|
client, err := l.OpenARI()
|
|
if err != nil {
|
|
return nil, &NativeCallFailure{Phase: "ari_open", Cause: err}
|
|
}
|
|
return dialWithClient(ctx, client, request)
|
|
}
|
|
|
|
func dialWithClient(ctx context.Context, client ari.Client, request NativeDial) (_ *NativeCall, err error) {
|
|
call := &NativeCall{client: client}
|
|
defer func() {
|
|
if err != nil {
|
|
err = errors.Join(err, call.Close())
|
|
}
|
|
}()
|
|
if ctx == nil || client == nil {
|
|
return nil, errors.New("native ARI call context and client required")
|
|
}
|
|
call.Context, call.cancel = context.WithCancelCause(ctx)
|
|
originate, err := approvedOriginateRequest(request.ExecutionID, request.TrunkID, request.DialedCallee, request.CallerID, request.AnswerTimeout)
|
|
if err != nil {
|
|
return nil, &NativeCallFailure{Phase: "validate", Cause: err}
|
|
}
|
|
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")
|
|
if call.subscription == nil {
|
|
return nil, &NativeCallFailure{Phase: "subscribe", Cause: errors.New("native ARI channel subscription unavailable before origination")}
|
|
}
|
|
call.outbound, err = client.Channel().Originate(key, originate)
|
|
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}
|
|
}
|
|
answerCtx, cancel := context.WithTimeout(ctx, request.AnswerTimeout)
|
|
defer cancel()
|
|
if err := awaitStasisStart(answerCtx, call.subscription, request.ExecutionID); err != nil {
|
|
return nil, &NativeCallFailure{Phase: "answer_wait", Cause: err}
|
|
}
|
|
bridgeKey := ari.NewKey(ari.BridgeKey, request.ExecutionID+"-bridge")
|
|
call.bridge, err = client.Bridge().Create(bridgeKey, "mixing", "go-sip-agent")
|
|
if err != nil {
|
|
// A lost reply does not prove that the bridge was not created.
|
|
call.bridge = client.Bridge().Get(bridgeKey)
|
|
return nil, &NativeCallFailure{Phase: "bridge_create", Cause: err}
|
|
}
|
|
if err := call.bridge.AddChannel(call.outbound.ID()); err != nil {
|
|
return nil, &NativeCallFailure{Phase: "bridge_add_outbound", Cause: err}
|
|
}
|
|
mediaKey := ari.NewKey(ari.ChannelKey, request.ExecutionID+"-media")
|
|
call.external, err = client.Channel().ExternalMedia(mediaKey, ari.ExternalMediaOptions{
|
|
App: "go-sip-agent", ExternalHost: net.JoinHostPort("127.0.0.1", strconv.Itoa(call.Media.LocalAddr().(*net.UDPAddr).Port)),
|
|
Format: "slin16", Encapsulation: "rtp", Transport: "udp", ConnectionType: "client", Direction: "both",
|
|
})
|
|
if err != nil {
|
|
// Use the same identity to stop an ExternalMedia channel created
|
|
// before the HTTP outcome became unknown; never create another.
|
|
call.external = client.Channel().Get(mediaKey)
|
|
return nil, &NativeCallFailure{Phase: "external_media_create", Cause: err}
|
|
}
|
|
address, err := call.external.GetVariable("UNICASTRTP_LOCAL_ADDRESS")
|
|
if err != nil {
|
|
return nil, &NativeCallFailure{Phase: "rtp_peer_address", Cause: err}
|
|
}
|
|
port, err := call.external.GetVariable("UNICASTRTP_LOCAL_PORT")
|
|
if err != nil {
|
|
return nil, &NativeCallFailure{Phase: "rtp_peer_port", Cause: err}
|
|
}
|
|
peer, err := netip.ParseAddr(address)
|
|
if err != nil || !peer.IsLoopback() {
|
|
return nil, &NativeCallFailure{Phase: "rtp_peer_validation", Cause: errors.New("Agent ARI RTP peer is not a local Cell address")}
|
|
}
|
|
if err := call.Media.SetPeer(net.JoinHostPort(address, port)); err != nil {
|
|
return nil, &NativeCallFailure{Phase: "rtp_peer_config", Cause: err}
|
|
}
|
|
if err := call.bridge.AddChannel(call.external.ID()); err != nil {
|
|
return nil, &NativeCallFailure{Phase: "bridge_add_media", Cause: err}
|
|
}
|
|
call.monitorDone = make(chan struct{})
|
|
go call.watch(request.ExecutionID)
|
|
return call, nil
|
|
}
|
|
|
|
func (c *NativeCall) watch(channelID string) {
|
|
defer close(c.monitorDone)
|
|
for {
|
|
select {
|
|
case <-c.Context.Done():
|
|
return
|
|
case event, open := <-c.subscription.Events():
|
|
if !open || event == nil {
|
|
c.cancel(errors.New("native ARI event subscription lost during call"))
|
|
return
|
|
}
|
|
if !eventBelongsToChannel(event, channelID) {
|
|
continue
|
|
}
|
|
switch event.GetType() {
|
|
case "ChannelDestroyed":
|
|
if ended, ok := event.(*ari.ChannelDestroyed); ok {
|
|
c.destroyed.Store(true)
|
|
c.cancel(fmt.Errorf("native channel ended: cause=%d", ended.Cause))
|
|
} else {
|
|
c.cancel(errors.New("native channel ended: invalid destroy event"))
|
|
}
|
|
return
|
|
case "ChannelHangupRequest", "StasisEnd":
|
|
c.cancel(fmt.Errorf("native channel ended: event=%s", event.GetType()))
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// Close releases each resource once. Errors remain visible to the caller; a
|
|
// failed Hangup never silently turns an uncertain call into a success.
|
|
func (c *NativeCall) Close() error {
|
|
if c == nil {
|
|
return nil
|
|
}
|
|
c.closeOnce.Do(func() {
|
|
if c.cancel != nil {
|
|
c.cancel(context.Canceled)
|
|
}
|
|
if c.subscription != nil {
|
|
c.subscription.Cancel()
|
|
}
|
|
if c.external != nil {
|
|
c.closeErr = errors.Join(c.closeErr, c.external.Hangup())
|
|
}
|
|
if c.outbound != nil && !c.destroyed.Load() {
|
|
c.closeErr = errors.Join(c.closeErr, c.outbound.Hangup())
|
|
}
|
|
if c.bridge != nil {
|
|
c.closeErr = errors.Join(c.closeErr, c.bridge.Delete())
|
|
}
|
|
if c.Media != nil {
|
|
c.closeErr = errors.Join(c.closeErr, c.Media.Close())
|
|
}
|
|
if c.client != nil {
|
|
c.client.Close()
|
|
}
|
|
if c.monitorDone != nil {
|
|
<-c.monitorDone
|
|
}
|
|
})
|
|
return c.closeErr
|
|
}
|