231 lines
6.5 KiB
Go
231 lines
6.5 KiB
Go
package asterisk
|
|
|
|
import (
|
|
"context"
|
|
"encoding/binary"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
"unicode/utf8"
|
|
|
|
"github.com/emiago/sipgo/sip"
|
|
)
|
|
|
|
// SIPResponse contains original, unmodified bytes from one observed SIP
|
|
// response. Never put Raw in logs or retained test evidence.
|
|
type SIPResponse struct {
|
|
Code int
|
|
StatusLine string
|
|
Raw string
|
|
}
|
|
|
|
type hepTransaction struct {
|
|
invited bool
|
|
response *SIPResponse
|
|
}
|
|
|
|
type hepCall struct {
|
|
callID string
|
|
last string
|
|
txns map[string]*hepTransaction
|
|
ready chan struct{}
|
|
}
|
|
|
|
// HEPMirror binds Asterisk's pre-Dial Call-ID to a single execution. HEP UUIDs
|
|
// for early SIP failures are the SIP Call-ID, not ARI channel IDs or names.
|
|
type HEPMirror struct {
|
|
mu sync.Mutex
|
|
byExecution map[string]*hepCall
|
|
byCallID map[string]*hepCall
|
|
}
|
|
|
|
func NewHEPMirror() *HEPMirror {
|
|
return &HEPMirror{byExecution: make(map[string]*hepCall), byCallID: make(map[string]*hepCall)}
|
|
}
|
|
|
|
func (m *HEPMirror) Bind(executionID, callID string) error {
|
|
if m == nil || executionID == "" || callID == "" || strings.TrimSpace(callID) != callID {
|
|
return errors.New("HEP requires a unique execution and exact SIP Call-ID")
|
|
}
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if m.byExecution[executionID] != nil || m.byCallID[callID] != nil {
|
|
return errors.New("HEP execution or SIP Call-ID is already bound")
|
|
}
|
|
call := &hepCall{callID: callID, txns: make(map[string]*hepTransaction), ready: make(chan struct{}, 1)}
|
|
m.byExecution[executionID], m.byCallID[callID] = call, call
|
|
return nil
|
|
}
|
|
|
|
func (m *HEPMirror) Release(executionID string) {
|
|
if m == nil {
|
|
return
|
|
}
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if c := m.byExecution[executionID]; c != nil {
|
|
delete(m.byCallID, c.callID)
|
|
delete(m.byExecution, executionID)
|
|
}
|
|
}
|
|
|
|
// Final waits only as long as its caller allows. A missing mirror is explicitly
|
|
// reported; it never delays a confirmed call end or invents a SIP status.
|
|
func (m *HEPMirror) Final(ctx context.Context, executionID string) (*SIPResponse, error) {
|
|
if ctx == nil || m == nil {
|
|
return nil, errors.New("HEP receiver unavailable")
|
|
}
|
|
defer m.Release(executionID)
|
|
for {
|
|
m.mu.Lock()
|
|
call := m.byExecution[executionID]
|
|
if call == nil {
|
|
m.mu.Unlock()
|
|
return nil, errors.New("HEP execution was never bound")
|
|
}
|
|
if tx := call.txns[call.last]; tx != nil && tx.invited && tx.response != nil {
|
|
resp := *tx.response
|
|
m.mu.Unlock()
|
|
return &resp, nil
|
|
}
|
|
ready := call.ready
|
|
m.mu.Unlock()
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, errors.New("original SIP INVITE response was not captured or correlated")
|
|
case <-ready:
|
|
}
|
|
}
|
|
}
|
|
|
|
// Observe accepts one HEPv3 UDP datagram. It decodes only the envelope chunks
|
|
// Asterisk emits; SIP itself is parsed by sipgo rather than a custom parser.
|
|
func (m *HEPMirror) Observe(frame []byte) error {
|
|
if m == nil {
|
|
return errors.New("HEP receiver unavailable")
|
|
}
|
|
cid, payload, err := unpackHEP3(frame)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
m.mu.Lock()
|
|
call := m.byCallID[cid]
|
|
m.mu.Unlock()
|
|
if call == nil { // other SIP traffic is not part of an Agent execution
|
|
return nil
|
|
}
|
|
if !utf8.Valid(payload) || len(payload) > 65535 {
|
|
return errors.New("HEP SIP message cannot be preserved verbatim as a result string")
|
|
}
|
|
message, err := sip.NewParser().ParseSIP(payload)
|
|
if err != nil {
|
|
return errors.New("invalid SIP payload in HEP mirror")
|
|
}
|
|
if message.CallID() == nil || message.CallID().Value() != cid || message.CSeq() == nil || message.Via() == nil {
|
|
return errors.New("HEP CID and SIP transaction headers disagree")
|
|
}
|
|
branch, ok := message.Via().Params.Get("branch")
|
|
if !ok || branch == "" {
|
|
return errors.New("HEP SIP transaction has no Via branch")
|
|
}
|
|
if message.CSeq().MethodName != sip.INVITE {
|
|
return nil
|
|
}
|
|
txID := fmt.Sprintf("%d/%s", message.CSeq().SeqNo, branch)
|
|
var response *SIPResponse
|
|
switch msg := message.(type) {
|
|
case *sip.Request:
|
|
if msg.Method != sip.INVITE {
|
|
return nil
|
|
}
|
|
case *sip.Response:
|
|
if msg.StatusCode < 200 || msg.StatusCode > 699 {
|
|
return nil
|
|
}
|
|
line, _, ok := strings.Cut(string(payload), "\r\n")
|
|
if !ok || !strings.HasPrefix(line, "SIP/2.0 ") || !strings.Contains(string(payload), "\r\n\r\n") {
|
|
return errors.New("HEP SIP response is missing its original status line or headers")
|
|
}
|
|
response = &SIPResponse{Code: msg.StatusCode, StatusLine: line, Raw: string(payload)}
|
|
default:
|
|
return nil
|
|
}
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
if m.byCallID[cid] != call { // execution finished while this packet was parsed
|
|
return nil
|
|
}
|
|
if call.last != "" && call.last != txID {
|
|
return nil // later reINVITEs cannot replace the original dial transaction
|
|
}
|
|
tx := call.txns[txID]
|
|
if tx == nil {
|
|
if len(call.txns) >= 32 {
|
|
return errors.New("HEP original SIP transaction exceeded the bounded pending response limit")
|
|
}
|
|
tx = &hepTransaction{}
|
|
call.txns[txID] = tx
|
|
}
|
|
if response == nil {
|
|
tx.invited = true
|
|
if call.last == "" {
|
|
call.last = txID
|
|
}
|
|
} else {
|
|
tx.response = response
|
|
}
|
|
if current := call.txns[call.last]; current != nil && current.invited && current.response != nil {
|
|
select {
|
|
case call.ready <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func unpackHEP3(frame []byte) (string, []byte, error) {
|
|
if len(frame) < 6 || len(frame) > 65535 || string(frame[:4]) != "HEP3" || int(binary.BigEndian.Uint16(frame[4:6])) != len(frame) {
|
|
return "", nil, errors.New("invalid HEP3 packet size or version")
|
|
}
|
|
var cid string
|
|
var payload []byte
|
|
proto := byte(0)
|
|
for offset := 6; offset < len(frame); {
|
|
if len(frame)-offset < 6 {
|
|
return "", nil, errors.New("short HEP3 chunk header")
|
|
}
|
|
vendor := binary.BigEndian.Uint16(frame[offset : offset+2])
|
|
kind := binary.BigEndian.Uint16(frame[offset+2 : offset+4])
|
|
length := int(binary.BigEndian.Uint16(frame[offset+4 : offset+6]))
|
|
if length < 6 || length > len(frame)-offset {
|
|
return "", nil, errors.New("invalid HEP3 chunk length")
|
|
}
|
|
data := frame[offset+6 : offset+length]
|
|
if vendor == 0 {
|
|
switch kind {
|
|
case 11: // HEP protocol type: 1 is SIP
|
|
if len(data) != 1 {
|
|
return "", nil, errors.New("invalid HEP3 protocol chunk")
|
|
}
|
|
proto = data[0]
|
|
case 15: // unmodified SIP payload
|
|
if payload != nil {
|
|
return "", nil, errors.New("duplicate HEP3 SIP payload")
|
|
}
|
|
payload = data
|
|
case 17: // Asterisk UUID (actual SIP Call-ID for early responses)
|
|
if cid != "" {
|
|
return "", nil, errors.New("duplicate HEP3 correlation ID")
|
|
}
|
|
cid = string(data)
|
|
}
|
|
}
|
|
offset += length
|
|
}
|
|
if proto != 1 || cid == "" || len(payload) == 0 {
|
|
return "", nil, errors.New("HEP3 packet lacks SIP payload or correlation ID")
|
|
}
|
|
return cid, payload, nil
|
|
}
|