Files
go-sip/internal/asterisk/call.go
T

299 lines
11 KiB
Go

package asterisk
import (
"context"
"errors"
"fmt"
"net"
"net/netip"
"strconv"
"strings"
"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
BindSIPCall func(executionID, callID string) error // before Dial; no number/time matching
}
// 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
ConfirmedEnd bool // only after a dialed pre-answer channel is verified absent from Asterisk
}
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"
}
// 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 {
var failure *NativeCallFailure
return errors.As(err, &failure) && failure.Phase == "answer_wait" && failure.ConfirmedEnd
}
// 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", "create", "caller_id", "sip_call_id", "sip_bind":
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
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
}
// 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
verifyEnd bool
confirmedEnd 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 {
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 {
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}
}
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", "ChannelStateChange", "StasisEnd", "ChannelHangupRequest", "ChannelDestroyed")
if call.subscription == nil {
return nil, &NativeCallFailure{Phase: "subscribe", Cause: errors.New("native ARI channel subscription unavailable before origination")}
}
// 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 {
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 := 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")
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.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())
}
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
}