Capture original INVITE responses through local HEP mirror
This commit is contained in:
@@ -7,11 +7,13 @@ import (
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net"
|
||||
"os"
|
||||
"strings"
|
||||
|
||||
agentpb "git.ipao.vip/rogee/go-sip/gen/agent"
|
||||
"git.ipao.vip/rogee/go-sip/internal/asterisk"
|
||||
"git.ipao.vip/rogee/go-sip/internal/config"
|
||||
"git.ipao.vip/rogee/go-sip/internal/rpc"
|
||||
"github.com/spf13/cobra"
|
||||
@@ -47,6 +49,17 @@ func newAgentCommand() *cobra.Command {
|
||||
}
|
||||
var handler *rpc.Server
|
||||
var recoveryClient agentpb.AgentControlServiceClient
|
||||
var mirror *asterisk.HEPMirror
|
||||
var hepConn net.PacketConn
|
||||
if mode == "nonprod-real" {
|
||||
hepConn, err = asterisk.ListenHEP(settings.HEPListen)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Agent SIP mirror listener unavailable: %w", err)
|
||||
}
|
||||
defer hepConn.Close()
|
||||
mirror = asterisk.NewHEPMirror()
|
||||
}
|
||||
|
||||
if mode == "sip-only" {
|
||||
handler, err = newSIPOnlyAgentServer(settings)
|
||||
} else {
|
||||
@@ -76,7 +89,7 @@ func newAgentCommand() *cobra.Command {
|
||||
if mode == "mock" {
|
||||
handler, err = newAgentServer(cmd.Context(), settings, scenario, applied, client)
|
||||
} else {
|
||||
handler, err = newRealAgentServer(cmd.Context(), settings, client)
|
||||
handler, err = newRealAgentServer(cmd.Context(), settings, client, mirror)
|
||||
recoveryClient = client
|
||||
}
|
||||
}
|
||||
@@ -93,6 +106,18 @@ func newAgentCommand() *cobra.Command {
|
||||
if mode == "nonprod-real" {
|
||||
startRealAgentRecovery(cmd.Context(), settings, recoveryClient, handler)
|
||||
}
|
||||
var hepErrors <-chan error
|
||||
if mirror != nil {
|
||||
faults := make(chan error, 1)
|
||||
hepErrors = faults
|
||||
go func() {
|
||||
fault := mirror.ServeHEP(cmd.Context(), hepConn, func(err error) { log.Printf("Agent HEP frame rejected: %v", err) })
|
||||
faults <- fault
|
||||
if fault != nil {
|
||||
server.Stop()
|
||||
}
|
||||
}()
|
||||
}
|
||||
stopped := make(chan struct{})
|
||||
go func() {
|
||||
select {
|
||||
@@ -103,6 +128,15 @@ func newAgentCommand() *cobra.Command {
|
||||
}()
|
||||
err = server.Serve(listener)
|
||||
close(stopped)
|
||||
if hepErrors != nil {
|
||||
select {
|
||||
case fault := <-hepErrors:
|
||||
if fault != nil {
|
||||
return fmt.Errorf("Agent SIP mirror listener stopped: %w", fault)
|
||||
}
|
||||
default:
|
||||
}
|
||||
}
|
||||
if err != nil && !errors.Is(err, grpc.ErrServerStopped) {
|
||||
return fmt.Errorf("Agent gRPC listener stopped: %w", err)
|
||||
}
|
||||
|
||||
@@ -30,10 +30,10 @@ func reportDefiniteNonDialFailure(ctx context.Context, cause error, client agent
|
||||
|
||||
// newRealAgentServer only exposes the signed Dispatcher execution path. There
|
||||
// is no scenario, synthetic recording, local dial entry, or Mock upload host.
|
||||
func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, dispatcher agentpb.AgentControlServiceClient) (*rpc.Server, error) {
|
||||
func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, dispatcher agentpb.AgentControlServiceClient, mirror *asterisk.HEPMirror) (*rpc.Server, error) {
|
||||
if ctx == nil || ctx.Err() != nil || settings.Mode != "nonprod-real" || settings.AgentID == "" || settings.CellID == "" ||
|
||||
settings.SessionPath == "" || settings.RecoveryRoot == "" || len(settings.PeerFingerprints) == 0 ||
|
||||
settings.OSSAllowedHost == "" || dispatcher == nil || tenant.ValidateDispatcherID(settings.DispatcherID) != nil {
|
||||
settings.OSSAllowedHost == "" || dispatcher == nil || mirror == nil || tenant.ValidateDispatcherID(settings.DispatcherID) != nil {
|
||||
return nil, errors.New("real Agent requires approved nonproduction identity, pinned Dispatcher and private recovery")
|
||||
}
|
||||
value := reflect.ValueOf(dispatcher)
|
||||
@@ -86,7 +86,7 @@ func newRealAgentServer(ctx context.Context, settings config.AgentEnvironment, d
|
||||
}
|
||||
return (&rpc.ApprovedRecordedRealCall{
|
||||
Loader: loader, MediaPayloadType: 118, EvidenceRoot: settings.EvidenceRoot,
|
||||
MaxWAVBytes: 64 << 20, ReportTimeout: 15 * time.Minute, Delivery: delivery,
|
||||
MaxWAVBytes: 64 << 20, ReportTimeout: 15 * time.Minute, Delivery: delivery, Mirror: mirror,
|
||||
}).Prepare(execution)
|
||||
},
|
||||
OnFailure: func(execution rpc.ApprovedExecution, cause error) error {
|
||||
|
||||
@@ -65,6 +65,12 @@ func TestRealAgentCommandStartsPinnedServerWithoutMockFixturesOrDial(t *testing.
|
||||
}
|
||||
listen := listener.Addr().String()
|
||||
listener.Close()
|
||||
hepPort, err := net.ListenPacket("udp4", "127.0.0.1:0")
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
hepListen := hepPort.LocalAddr().String()
|
||||
hepPort.Close()
|
||||
for name, value := range map[string]string{
|
||||
"AGENT_ID": settings.AgentID, "CELL_ID": settings.CellID,
|
||||
"AGENT_GRPC_LISTEN": listen, "AGENT_SESSION_PATH": settings.SessionPath,
|
||||
@@ -76,6 +82,7 @@ func TestRealAgentCommandStartsPinnedServerWithoutMockFixturesOrDial(t *testing.
|
||||
"ASTERISK_BIN": "/usr/sbin/asterisk", "ASTERISK_LIBRARY_DIR": "/usr/lib/asterisk/modules",
|
||||
"AGENT_EVIDENCE_ROOT": filepath.Join(settings.RecoveryRoot, "evidence"),
|
||||
"AGENT_OSS_ALLOWED_HOST": "bucket.oss-cn-beijing.aliyuncs.com",
|
||||
"AGENT_HEP_LISTEN_ADDR": hepListen,
|
||||
"AGENT_MOCK_SCENARIO_FILE": "", "AGENT_MOCK_APPLIED_SIP_FILE": "",
|
||||
} {
|
||||
t.Setenv(name, value)
|
||||
@@ -133,7 +140,8 @@ func TestRealAgentServerRequiresBoundNativeSIPAndNoMock(t *testing.T) {
|
||||
settings.EvidenceRoot = filepath.Join(settings.RecoveryRoot, "evidence")
|
||||
settings.OSSAllowedHost = "bucket.oss-cn-beijing.aliyuncs.com"
|
||||
client := &isolatedAgentRecordingClient{}
|
||||
server, err := newRealAgentServer(context.Background(), settings, client)
|
||||
mirror := asterisk.NewHEPMirror()
|
||||
server, err := newRealAgentServer(context.Background(), settings, client, mirror)
|
||||
if err != nil || server == nil {
|
||||
t.Fatalf("real Agent must use the native SIP and non-Mock delivery boundary: %v", err)
|
||||
}
|
||||
@@ -142,11 +150,14 @@ func TestRealAgentServerRequiresBoundNativeSIPAndNoMock(t *testing.T) {
|
||||
}
|
||||
var typedNil *isolatedAgentRecordingClient
|
||||
var nilClient agentpb.AgentControlServiceClient = typedNil
|
||||
if server, err := newRealAgentServer(context.Background(), settings, nilClient); err == nil || server != nil {
|
||||
if server, err := newRealAgentServer(context.Background(), settings, nilClient, mirror); err == nil || server != nil {
|
||||
t.Fatal("typed-nil Dispatcher must block real Agent startup")
|
||||
}
|
||||
if server, err := newRealAgentServer(context.Background(), settings, client, nil); err == nil || server != nil {
|
||||
t.Fatal("real Agent must require an initialized SIP mirror")
|
||||
}
|
||||
settings.MockScenarioFile = "old-mock.json"
|
||||
if server, err := newRealAgentServer(context.Background(), settings, client); err == nil || server != nil {
|
||||
if server, err := newRealAgentServer(context.Background(), settings, client, mirror); err == nil || server != nil {
|
||||
t.Fatal("real Agent must not accept isolated Mock fixtures")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -26,10 +26,15 @@ require (
|
||||
github.com/cyberphone/json-canonicalization v0.0.0-20241213102144-19d51d7fe467 // indirect
|
||||
github.com/dustin/go-humanize v1.0.1 // indirect
|
||||
github.com/ebitengine/purego v0.10.2 // indirect
|
||||
github.com/emiago/sipgo v1.6.0 // indirect
|
||||
github.com/go-ole/go-ole v1.2.6 // indirect
|
||||
github.com/go-stack/stack v1.8.0 // indirect
|
||||
github.com/gobwas/httphead v0.1.0 // indirect
|
||||
github.com/gobwas/pool v0.2.1 // indirect
|
||||
github.com/gobwas/ws v1.3.2 // indirect
|
||||
github.com/gogo/protobuf v1.3.2 // indirect
|
||||
github.com/gorilla/websocket v1.5.3 // indirect
|
||||
github.com/icholy/digest v1.1.0 // indirect
|
||||
github.com/inconshreveable/log15 v0.0.0-20201112154412-8562bdadbbac // indirect
|
||||
github.com/inconshreveable/mousetrap v1.1.0 // indirect
|
||||
github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 // indirect
|
||||
@@ -50,6 +55,7 @@ require (
|
||||
github.com/tklauser/numcpus v0.11.0 // indirect
|
||||
github.com/yusufpapurcu/wmi v1.2.4 // indirect
|
||||
golang.org/x/net v0.58.0 // indirect
|
||||
golang.org/x/sync v0.22.0 // indirect
|
||||
golang.org/x/sys v0.47.0 // indirect
|
||||
golang.org/x/text v0.41.0 // indirect
|
||||
google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect
|
||||
|
||||
@@ -21,6 +21,8 @@ github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkp
|
||||
github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto=
|
||||
github.com/ebitengine/purego v0.10.2 h1:W809HbnvzAxgdm+aOvlSekrM16wGCdT/e76+9tS7gzE=
|
||||
github.com/ebitengine/purego v0.10.2/go.mod h1:iIjxzd6CiRiOG0UyXP+V1+jWqUXVjPKLAI0mRfJZTmQ=
|
||||
github.com/emiago/sipgo v1.6.0 h1:6EuOP7c6f0VRatKYTPEYNezt4hslBEsaCzZZOhT2n3s=
|
||||
github.com/emiago/sipgo v1.6.0/go.mod h1:DuwAxBZhKMqIzQFPGZb1MVAGU6Wuxj64oTOhd5dx/FY=
|
||||
github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI=
|
||||
github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY=
|
||||
github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag=
|
||||
@@ -29,6 +31,12 @@ github.com/go-ole/go-ole v1.2.6 h1:/Fpf6oFPoeFik9ty7siob0G6Ke8QvQEuVcuChpwXzpY=
|
||||
github.com/go-ole/go-ole v1.2.6/go.mod h1:pprOEPIfldk/42T2oK7lQ4v4JSDwmV0As9GaiUsvbm0=
|
||||
github.com/go-stack/stack v1.8.0 h1:5SgMzNM5HxrEjV0ww2lTmX6E2Izsfxas4+YHWRs3Lsk=
|
||||
github.com/go-stack/stack v1.8.0/go.mod h1:v0f6uXyyMGvRgIKkXu+yp6POWl0qKG85gN/melR3HDY=
|
||||
github.com/gobwas/httphead v0.1.0 h1:exrUm0f4YX0L7EBwZHuCF4GDp8aJfVeBrlLQrs6NqWU=
|
||||
github.com/gobwas/httphead v0.1.0/go.mod h1:O/RXo79gxV8G+RqlR/otEwx4Q36zl9rqC5u12GKvMCM=
|
||||
github.com/gobwas/pool v0.2.1 h1:xfeeEhW7pwmX8nuLVlqbzVc7udMDrwetjEv+TZIz1og=
|
||||
github.com/gobwas/pool v0.2.1/go.mod h1:q8bcK0KcYlCgd9e7WYLm9LpyS+YeLd8JVDW6WezmKEw=
|
||||
github.com/gobwas/ws v1.3.2 h1:zlnbNHxumkRvfPWgfXu8RBwyNR1x8wh9cf5PTOCqs9Q=
|
||||
github.com/gobwas/ws v1.3.2/go.mod h1:hRKAFb8wOxFROYNsT1bqfWnhX+b5MFeJM9r2ZSwg/KY=
|
||||
github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q=
|
||||
github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q=
|
||||
github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek=
|
||||
@@ -44,6 +52,8 @@ github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aN
|
||||
github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
|
||||
github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k=
|
||||
github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM=
|
||||
github.com/icholy/digest v1.1.0 h1:HfGg9Irj7i+IX1o1QAmPfIBNu/Q5A5Tu3n/MED9k9H4=
|
||||
github.com/icholy/digest v1.1.0/go.mod h1:QNrsSGQ5v7v9cReDI0+eyjsXGUoRSUZQHeQ5C4XLa0Y=
|
||||
github.com/inconshreveable/log15 v0.0.0-20201112154412-8562bdadbbac h1:n1DqxAo4oWPMvH1+v+DLYlMCecgumhhgnxAPdqDIFHI=
|
||||
github.com/inconshreveable/log15 v0.0.0-20201112154412-8562bdadbbac/go.mod h1:cOaXtrgN4ScfRrD9Bre7U1thNq5RtJ8ZoP4iXVGRj6o=
|
||||
github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8=
|
||||
@@ -167,6 +177,7 @@ golang.org/x/sys v0.0.0-20201204225414-ed752295db88/go.mod h1:h1NjWce9XRLGQEsW7w
|
||||
golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.6.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
|
||||
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
|
||||
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
|
||||
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
|
||||
|
||||
@@ -0,0 +1,230 @@
|
||||
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
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
package asterisk
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"net"
|
||||
"time"
|
||||
)
|
||||
|
||||
// ListenHEP opens the already-approved on-host SIP mirror before real call
|
||||
// admission. No traffic is sent to any trunk by this listener.
|
||||
func ListenHEP(address string) (net.PacketConn, error) {
|
||||
host, _, err := net.SplitHostPort(address)
|
||||
if err != nil || host != "127.0.0.1" {
|
||||
return nil, errors.New("HEP receiver must bind an explicit IPv4 loopback address")
|
||||
}
|
||||
return net.ListenPacket("udp4", address)
|
||||
}
|
||||
|
||||
// ServeHEP reports only static parse errors; it never logs SIP bytes, numbers,
|
||||
// Call-IDs, addresses from HEP payloads, or any provider response headers.
|
||||
func (m *HEPMirror) ServeHEP(ctx context.Context, conn net.PacketConn, report func(error)) error {
|
||||
if ctx == nil || conn == nil || m == nil {
|
||||
return errors.New("HEP receiver requires a context, socket, and mirror")
|
||||
}
|
||||
buf := make([]byte, 65535)
|
||||
for {
|
||||
if err := conn.SetReadDeadline(time.Now().Add(time.Second)); err != nil {
|
||||
return err
|
||||
}
|
||||
n, sender, err := conn.ReadFrom(buf)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return nil
|
||||
}
|
||||
if timeout, ok := err.(net.Error); ok && timeout.Timeout() {
|
||||
continue
|
||||
}
|
||||
return err
|
||||
}
|
||||
peer, ok := sender.(*net.UDPAddr)
|
||||
if !ok || !peer.IP.IsLoopback() {
|
||||
if report != nil {
|
||||
report(errors.New("HEP mirror frame was not sent by the local Asterisk host"))
|
||||
}
|
||||
continue
|
||||
}
|
||||
if err := m.Observe(buf[:n]); err != nil && report != nil {
|
||||
report(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
package asterisk
|
||||
|
||||
import (
|
||||
"context"
|
||||
"net"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestHEPListenerOnlyBindsLoopbackAndReceivesSIP(t *testing.T) {
|
||||
for _, address := range []string{"0.0.0.0:9060", "[::]:9060", "localhost:9060", "127.0.0.1:not-a-port"} {
|
||||
if conn, err := ListenHEP(address); err == nil {
|
||||
conn.Close()
|
||||
t.Fatalf("uncontrolled HEP binding accepted: %q", address)
|
||||
}
|
||||
}
|
||||
conn, err := ListenHEP("127.0.0.1:0") // test mode uses an ephemeral port; production requires a fixed port.
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer conn.Close()
|
||||
mirror := NewHEPMirror()
|
||||
if err := mirror.Bind("exec", "sip-id"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
errors := make(chan error, 1)
|
||||
go func() { errors <- mirror.ServeHEP(ctx, conn, nil) }()
|
||||
sender, err := net.Dial("udp", conn.LocalAddr().String())
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
defer sender.Close()
|
||||
for _, raw := range []string{
|
||||
hepTestSIP("sip-id", "INVITE", "z9hG4bK-a", "INVITE sip:callee@localhost SIP/2.0"),
|
||||
hepTestSIP("sip-id", "INVITE", "z9hG4bK-a", "SIP/2.0 480 Temporarily Unavailable"),
|
||||
} {
|
||||
if _, err := sender.Write(hepTestFrame("sip-id", raw)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
finalCtx, stop := context.WithTimeout(context.Background(), time.Second)
|
||||
result, err := mirror.Final(finalCtx, "exec")
|
||||
stop()
|
||||
if err != nil || result.Code != 480 {
|
||||
t.Fatalf("HEP loopback result missing: %+v err=%v", result, err)
|
||||
}
|
||||
cancel()
|
||||
conn.Close()
|
||||
select {
|
||||
case err := <-errors:
|
||||
if err != nil {
|
||||
t.Fatalf("HEP listener did not stop cleanly: %v", err)
|
||||
}
|
||||
case <-time.After(time.Second):
|
||||
t.Fatal("HEP listener did not stop")
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,140 @@
|
||||
package asterisk
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func hepTestFrame(cid, sip string) []byte {
|
||||
packet := []byte("HEP3\x00\x00")
|
||||
for _, chunk := range []struct {
|
||||
kind uint16
|
||||
data []byte
|
||||
}{{11, []byte{1}}, {17, []byte(cid)}, {15, []byte(sip)}} {
|
||||
part := make([]byte, 6+len(chunk.data))
|
||||
binary.BigEndian.PutUint16(part[2:4], chunk.kind)
|
||||
binary.BigEndian.PutUint16(part[4:6], uint16(len(part)))
|
||||
copy(part[6:], chunk.data)
|
||||
packet = append(packet, part...)
|
||||
}
|
||||
binary.BigEndian.PutUint16(packet[4:6], uint16(len(packet)))
|
||||
return packet
|
||||
}
|
||||
|
||||
func hepTestSIP(callID, method, branch, firstLine string) string {
|
||||
return firstLine + "\r\nVia: SIP/2.0/UDP 127.0.0.1:5060;branch=" + branch +
|
||||
"\r\nFrom: <sip:caller@localhost>;tag=a\r\nTo: <sip:callee@localhost>\r\nCall-ID: " + callID +
|
||||
"\r\nCSeq: 1 " + method + "\r\nContent-Length: 0\r\n\r\n"
|
||||
}
|
||||
|
||||
func TestHEPMirrorKeepsConcurrentSIPTransactionsSeparate(t *testing.T) {
|
||||
mirror := NewHEPMirror()
|
||||
if err := mirror.Bind("exec-a", "sip-id-a"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := mirror.Bind("exec-b", "sip-id-b"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if err := mirror.Bind("exec-duplicate", "sip-id-a"); err == nil {
|
||||
t.Fatal("one SIP Call-ID cannot belong to two executions")
|
||||
}
|
||||
var wg sync.WaitGroup
|
||||
for _, tc := range []struct {
|
||||
id, branch, code string
|
||||
}{{"sip-id-a", "z9hG4bK-a", "480 Temporarily Unavailable"}, {"sip-id-b", "z9hG4bK-b", "486 Busy Here"}} {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
invite := hepTestSIP(tc.id, "INVITE", tc.branch, "INVITE sip:callee@localhost SIP/2.0")
|
||||
response := hepTestSIP(tc.id, "INVITE", tc.branch, "SIP/2.0 "+tc.code)
|
||||
if err := mirror.Observe(hepTestFrame(tc.id, response)); err != nil { // UDP may reorder the response.
|
||||
t.Error(err)
|
||||
}
|
||||
if err := mirror.Observe(hepTestFrame(tc.id, invite)); err != nil {
|
||||
t.Error(err)
|
||||
}
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
for _, tc := range []struct {
|
||||
execID, want string
|
||||
}{{"exec-a", "SIP/2.0 480 Temporarily Unavailable"}, {"exec-b", "SIP/2.0 486 Busy Here"}} {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
resp, err := mirror.Final(ctx, tc.execID)
|
||||
cancel()
|
||||
if err != nil || resp.StatusLine != tc.want || !strings.HasPrefix(resp.Raw, tc.want+"\r\n") {
|
||||
t.Fatalf("wrong execution response: %+v err=%v", resp, err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestHEPMirrorRejectsUnmatchedAndMalformedPacketsWithoutGuessing(t *testing.T) {
|
||||
mirror := NewHEPMirror()
|
||||
if err := mirror.Bind("exec-a", "sip-id-a"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, bad := range [][]byte{nil, []byte("HEP3"), []byte("HEP3\x00\x06"), []byte("HEP3\x00\x0c\x00\x00\x00\x0f\x00\x00")} {
|
||||
if err := mirror.Observe(bad); err == nil {
|
||||
t.Fatalf("malformed frame accepted: len=%d", len(bad))
|
||||
}
|
||||
}
|
||||
unrelated := hepTestSIP("sip-id-b", "INVITE", "z9hG4bK-b", "SIP/2.0 480 Temporarily Unavailable")
|
||||
if err := mirror.Observe(hepTestFrame("sip-id-b", unrelated)); err != nil {
|
||||
t.Fatalf("unrelated traffic is not a parser error: %v", err)
|
||||
}
|
||||
incorrectCID := hepTestSIP("sip-id-b", "INVITE", "z9hG4bK-b", "SIP/2.0 480 Temporarily Unavailable")
|
||||
if err := mirror.Observe(hepTestFrame("sip-id-a", incorrectCID)); err == nil {
|
||||
t.Fatal("payload Call-ID and HEP CID disagreed")
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond)
|
||||
defer cancel()
|
||||
if resp, err := mirror.Final(ctx, "exec-a"); err == nil || resp != nil {
|
||||
t.Fatalf("missing observed response was fabricated: %+v err=%v", resp, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHEPMirrorKeepsOriginalInviteWhenLaterReinviteArrives(t *testing.T) {
|
||||
mirror := NewHEPMirror()
|
||||
if err := mirror.Bind("exec", "sip-id"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, sip := range []string{
|
||||
hepTestSIP("sip-id", "INVITE", "z9hG4bK-first", "INVITE sip:callee@localhost SIP/2.0"),
|
||||
hepTestSIP("sip-id", "INVITE", "z9hG4bK-first", "SIP/2.0 200 OK"),
|
||||
hepTestSIP("sip-id", "INVITE", "z9hG4bK-later", "INVITE sip:callee@localhost SIP/2.0"),
|
||||
hepTestSIP("sip-id", "INVITE", "z9hG4bK-later", "SIP/2.0 488 Not Acceptable Here"),
|
||||
} {
|
||||
if err := mirror.Observe(hepTestFrame("sip-id", sip)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
resp, err := mirror.Final(context.Background(), "exec")
|
||||
if err != nil || resp.Code != 200 {
|
||||
t.Fatalf("later reINVITE must not replace original INVITE outcome: %+v %v", resp, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestHEPMirrorNeverTreatsProvisionalOrOtherMethodsAsFinalINVITE(t *testing.T) {
|
||||
mirror := NewHEPMirror()
|
||||
if err := mirror.Bind("exec-a", "sip-id-a"); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
for _, sip := range []string{
|
||||
hepTestSIP("sip-id-a", "INVITE", "z9hG4bK-a", "INVITE sip:callee@localhost SIP/2.0"),
|
||||
hepTestSIP("sip-id-a", "INVITE", "z9hG4bK-a", "SIP/2.0 180 Ringing"),
|
||||
hepTestSIP("sip-id-a", "BYE", "z9hG4bK-a", "SIP/2.0 200 OK"),
|
||||
} {
|
||||
if err := mirror.Observe(hepTestFrame("sip-id-a", sip)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 10*time.Millisecond)
|
||||
defer cancel()
|
||||
if resp, err := mirror.Final(ctx, "exec-a"); err == nil || resp != nil {
|
||||
t.Fatalf("non-final or non-INVITE response became an outcome: %+v err=%v", resp, err)
|
||||
}
|
||||
}
|
||||
@@ -35,6 +35,7 @@ type AgentEnvironment struct {
|
||||
AsteriskLibraryDir string
|
||||
EvidenceRoot string
|
||||
OSSAllowedHost string
|
||||
HEPListen string
|
||||
}
|
||||
|
||||
func LoadAgentEnvironment(mode string) (AgentEnvironment, error) {
|
||||
@@ -104,6 +105,9 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) {
|
||||
if settings.OSSAllowedHost, err = get("AGENT_OSS_ALLOWED_HOST"); err != nil {
|
||||
return AgentEnvironment{}, err
|
||||
}
|
||||
if settings.HEPListen, err = get("AGENT_HEP_LISTEN_ADDR"); err != nil {
|
||||
return AgentEnvironment{}, err
|
||||
}
|
||||
}
|
||||
if tenant.ValidateDispatcherID(settings.DispatcherID) != nil {
|
||||
return AgentEnvironment{}, errors.New("DISPATCHER_ID must be a canonical UUID v4")
|
||||
@@ -115,6 +119,9 @@ func LoadAgentEnvironment(mode string) (AgentEnvironment, error) {
|
||||
return AgentEnvironment{}, errors.New("DISPATCHER_GRPC_ENDPOINT must be an isolated local Mock address")
|
||||
}
|
||||
if mode == "nonprod-real" {
|
||||
if !strings.HasPrefix(settings.HEPListen, "127.0.0.1:") || !localGRPCAddress(settings.HEPListen, false) {
|
||||
return AgentEnvironment{}, errors.New("AGENT_HEP_LISTEN_ADDR must be an explicit IPv4 loopback address with a nonzero port")
|
||||
}
|
||||
if !filepath.IsAbs(settings.EvidenceRoot) {
|
||||
return AgentEnvironment{}, errors.New("AGENT_EVIDENCE_ROOT must be absolute")
|
||||
}
|
||||
|
||||
@@ -68,17 +68,25 @@ func TestLoadNonprodRealAgentRequiresNativeCallAndEvidenceSettings(t *testing.T)
|
||||
"ASTERISK_LIBRARY_DIR": "/tmp/asterisk-libraries",
|
||||
"AGENT_EVIDENCE_ROOT": "/tmp/nonprod-call-evidence",
|
||||
"AGENT_OSS_ALLOWED_HOST": "test.oss-cn-beijing.aliyuncs.com",
|
||||
"AGENT_HEP_LISTEN_ADDR": "127.0.0.1:19060",
|
||||
} {
|
||||
t.Setenv(name, value)
|
||||
}
|
||||
settings, err := LoadAgentEnvironment("nonprod-real")
|
||||
if err != nil || settings.AsteriskConfigDir == "" || settings.DispatcherEndpoint == "" || settings.EvidenceRoot != "/tmp/nonprod-call-evidence" || settings.OSSAllowedHost != "test.oss-cn-beijing.aliyuncs.com" || settings.MockScenarioFile != "" {
|
||||
if err != nil || settings.AsteriskConfigDir == "" || settings.DispatcherEndpoint == "" || settings.EvidenceRoot != "/tmp/nonprod-call-evidence" || settings.OSSAllowedHost != "test.oss-cn-beijing.aliyuncs.com" || settings.HEPListen != "127.0.0.1:19060" || settings.MockScenarioFile != "" {
|
||||
t.Fatalf("nonproduction real mode must use explicit native, capture and OSS settings: %+v %v", settings, err)
|
||||
}
|
||||
if _, err := os.Stat(state); !os.IsNotExist(err) {
|
||||
t.Fatalf("configuration inspection opened durable session: %v", err)
|
||||
}
|
||||
for _, name := range []string{"ASTERISK_CONFIG_DIR", "ASTERISK_BIN", "ASTERISK_LIBRARY_DIR", "AGENT_EVIDENCE_ROOT", "AGENT_OSS_ALLOWED_HOST", "DISPATCHER_GRPC_ENDPOINT", "DISPATCHER_GRPC_SERVER_NAME"} {
|
||||
for _, address := range []string{"0.0.0.0:19060", "localhost:19060", "127.0.0.1:0", "127.0.0.1:bad"} {
|
||||
t.Setenv("AGENT_HEP_LISTEN_ADDR", address)
|
||||
if _, err := LoadAgentEnvironment("nonprod-real"); err == nil {
|
||||
t.Fatalf("uncontrolled or unusable HEP listener admitted: %q", address)
|
||||
}
|
||||
}
|
||||
t.Setenv("AGENT_HEP_LISTEN_ADDR", "127.0.0.1:19060")
|
||||
for _, name := range []string{"ASTERISK_CONFIG_DIR", "ASTERISK_BIN", "ASTERISK_LIBRARY_DIR", "AGENT_EVIDENCE_ROOT", "AGENT_OSS_ALLOWED_HOST", "AGENT_HEP_LISTEN_ADDR", "DISPATCHER_GRPC_ENDPOINT", "DISPATCHER_GRPC_SERVER_NAME"} {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
t.Setenv(name, "")
|
||||
if _, err := LoadAgentEnvironment("nonprod-real"); err == nil || !strings.Contains(err.Error(), name) {
|
||||
@@ -87,7 +95,7 @@ func TestLoadNonprodRealAgentRequiresNativeCallAndEvidenceSettings(t *testing.T)
|
||||
t.Setenv(name, map[string]string{
|
||||
"ASTERISK_CONFIG_DIR": "/tmp/asterisk-config", "ASTERISK_BIN": "/tmp/asterisk",
|
||||
"ASTERISK_LIBRARY_DIR": "/tmp/asterisk-libraries", "AGENT_EVIDENCE_ROOT": "/tmp/nonprod-call-evidence",
|
||||
"AGENT_OSS_ALLOWED_HOST": "test.oss-cn-beijing.aliyuncs.com", "DISPATCHER_GRPC_ENDPOINT": "127.0.0.1:39443",
|
||||
"AGENT_OSS_ALLOWED_HOST": "test.oss-cn-beijing.aliyuncs.com", "AGENT_HEP_LISTEN_ADDR": "127.0.0.1:19060", "DISPATCHER_GRPC_ENDPOINT": "127.0.0.1:39443",
|
||||
"DISPATCHER_GRPC_SERVER_NAME": "dispatcher.local",
|
||||
}[name])
|
||||
})
|
||||
|
||||
@@ -27,6 +27,7 @@ type ApprovedRecordedRealCall struct {
|
||||
MaxWAVBytes int64
|
||||
ReportTimeout time.Duration
|
||||
Delivery *agent.RecordingDelivery
|
||||
Mirror *asterisk.HEPMirror
|
||||
originator func(context.Context, asterisk.NativeDial) (realMedia, error) // isolated tests only
|
||||
}
|
||||
|
||||
@@ -39,7 +40,7 @@ type realMedia struct {
|
||||
// Prepare rejects missing signed identity and recovery capacity before the
|
||||
// Agent acknowledges the call. No SIP channel is created during preparation.
|
||||
func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(context.Context) error, error) {
|
||||
if r == nil || r.Delivery == nil || r.ReportTimeout <= 0 || r.MaxWAVBytes <= 44 || r.MediaPayloadType < 96 || r.MediaPayloadType > 127 {
|
||||
if r == nil || r.Delivery == nil || r.Mirror == nil || r.ReportTimeout <= 0 || r.MaxWAVBytes <= 44 || r.MediaPayloadType < 96 || r.MediaPayloadType > 127 {
|
||||
return nil, errors.New("real call requires bounded media and per-call delivery")
|
||||
}
|
||||
if approved.DispatcherID == "" || approved.TenantID <= 0 || approved.SourceEventID == "" || approved.CallID != approved.SourceEventID || approved.TaskID == "" ||
|
||||
@@ -81,12 +82,24 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con
|
||||
ExecutionID: approved.SourceEventID, TrunkID: approved.SelectedTrunkID,
|
||||
DialedCallee: approved.DialedCallee, CallerID: approved.CallerID,
|
||||
AnswerTimeout: approved.RingTimeout, MediaPayloadType: r.MediaPayloadType,
|
||||
BindSIPCall: r.Mirror.Bind,
|
||||
}
|
||||
addSIPResponse := func(ctx context.Context, facts *callflow.FinalCallFacts) {
|
||||
mirrorCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), time.Second)
|
||||
defer cancel()
|
||||
response, err := r.Mirror.Final(mirrorCtx, approved.SourceEventID)
|
||||
if err != nil {
|
||||
facts.SIPCaptureError = err.Error() // static failure text; never raw SIP bytes
|
||||
return
|
||||
}
|
||||
facts.ReasonCode, facts.SIPStatusLine, facts.SIPRaw = &response.Code, response.StatusLine, response.Raw
|
||||
}
|
||||
var started atomic.Bool
|
||||
return func(ctx context.Context) error {
|
||||
if !started.CompareAndSwap(false, true) {
|
||||
return errors.New("approved real call already started; refusing a second origination")
|
||||
return errors.New("approved real call already started; refusing a second dial")
|
||||
}
|
||||
defer r.Mirror.Release(approved.SourceEventID)
|
||||
if ctx == nil || ctx.Err() != nil || !time.Now().Before(approved.DialBefore) {
|
||||
return errors.New("real call instruction expired or cancelled before origination")
|
||||
}
|
||||
@@ -101,11 +114,13 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con
|
||||
if !asterisk.NativeCallConfirmedPreAnswerEnd(err) {
|
||||
return err // uncertain origination or end is never retried or reported as completed
|
||||
}
|
||||
payload, payloadErr := callflow.FinalResultPayload(callflow.FinalCallFacts{
|
||||
facts := callflow.FinalCallFacts{
|
||||
TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID,
|
||||
Callee: approved.Callee, TrunkID: approved.SelectedTrunkID,
|
||||
StartedAt: attemptedAt, EndedAt: time.Now().UTC(), Outcome: "failed", ReasonMessage: "call ended before answer",
|
||||
}, callflow.Result{})
|
||||
}
|
||||
addSIPResponse(ctx, &facts)
|
||||
payload, payloadErr := callflow.FinalResultPayload(facts, callflow.Result{})
|
||||
if payloadErr != nil {
|
||||
return errors.Join(err, payloadErr)
|
||||
}
|
||||
@@ -113,7 +128,7 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con
|
||||
defer cancel()
|
||||
return errors.Join(err, r.Delivery.Complete(reportCtx, agent.CompletedRecording{ResultPayload: payload}))
|
||||
}
|
||||
startedAt := time.Now().UTC() // Originate returns only after the actual channel entered Stasis.
|
||||
startedAt := time.Now().UTC() // Dial returned and the channel entered Stasis and became Up.
|
||||
// A real channel may have answered even if media/AI setup subsequently
|
||||
// fails. Report a failed call only after its hangup is confirmed; never
|
||||
// keep a known-ended call occupying capacity or invent a recording.
|
||||
@@ -124,11 +139,13 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con
|
||||
if err := call.close(); err != nil {
|
||||
return errors.Join(cause, err)
|
||||
}
|
||||
payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{
|
||||
facts := callflow.FinalCallFacts{
|
||||
TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID,
|
||||
Callee: approved.Callee, TrunkID: approved.SelectedTrunkID,
|
||||
StartedAt: startedAt, EndedAt: time.Now().UTC(), Outcome: "failed", ReasonMessage: reason,
|
||||
}, callflow.Result{})
|
||||
}
|
||||
addSIPResponse(ctx, &facts)
|
||||
payload, err := callflow.FinalResultPayload(facts, callflow.Result{})
|
||||
if err != nil {
|
||||
return errors.Join(cause, err)
|
||||
}
|
||||
@@ -160,11 +177,13 @@ func (r *ApprovedRecordedRealCall) Prepare(approved ApprovedExecution) (func(con
|
||||
if runErr != nil {
|
||||
outcome, reason = "failed", "real AI or media flow failed"
|
||||
}
|
||||
payload, err := callflow.FinalResultPayload(callflow.FinalCallFacts{
|
||||
facts := callflow.FinalCallFacts{
|
||||
TaskID: approved.TaskID, CallerProfileID: approved.CallerProfileID,
|
||||
Callee: approved.Callee, TrunkID: approved.SelectedTrunkID,
|
||||
StartedAt: startedAt, EndedAt: endedAt, Outcome: outcome, ReasonMessage: reason,
|
||||
}, observed)
|
||||
}
|
||||
addSIPResponse(ctx, &facts)
|
||||
payload, err := callflow.FinalResultPayload(facts, observed)
|
||||
if err != nil {
|
||||
return errors.Join(runErr, err)
|
||||
}
|
||||
|
||||
@@ -2,6 +2,7 @@ package rpc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/binary"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
@@ -72,7 +73,7 @@ func TestApprovedRecordedRealCallReportsOnlyEndedObservedCall(t *testing.T) {
|
||||
fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
var originate, hangup int
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
runner := &ApprovedRecordedRealCall{Mirror: asterisk.NewHEPMirror(),
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(ctx context.Context, dial asterisk.NativeDial) (realMedia, error) {
|
||||
@@ -100,7 +101,7 @@ func TestApprovedRecordedRealCallReportsConfirmedMediaSetupFailure(t *testing.T)
|
||||
fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
var originate, hangup int
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
runner := &ApprovedRecordedRealCall{Mirror: asterisk.NewHEPMirror(),
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(ctx context.Context, _ asterisk.NativeDial) (realMedia, error) {
|
||||
@@ -125,7 +126,7 @@ func TestApprovedRecordedRealCallReportsConfirmedMediaSetupFailure(t *testing.T)
|
||||
func TestApprovedRecordedRealCallMissingMediaWithUnknownHangupDoesNotReport(t *testing.T) {
|
||||
fixture, stub, _, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
runner := &ApprovedRecordedRealCall{Mirror: asterisk.NewHEPMirror(),
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(ctx context.Context, _ asterisk.NativeDial) (realMedia, error) {
|
||||
@@ -141,7 +142,7 @@ func TestApprovedRecordedRealCallDoesNotReportUnknownHangup(t *testing.T) {
|
||||
fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
var originate, hangup int
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
runner := &ApprovedRecordedRealCall{Mirror: asterisk.NewHEPMirror(),
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(ctx context.Context, _ asterisk.NativeDial) (realMedia, error) {
|
||||
@@ -154,11 +155,63 @@ func TestApprovedRecordedRealCallDoesNotReportUnknownHangup(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func hepRealResultFrame(callID, payload string) []byte {
|
||||
b := []byte("HEP3\x00\x00")
|
||||
for _, chunk := range []struct {
|
||||
kind uint16
|
||||
body []byte
|
||||
}{{11, []byte{1}}, {17, []byte(callID)}, {15, []byte(payload)}} {
|
||||
p := make([]byte, 6+len(chunk.body))
|
||||
binary.BigEndian.PutUint16(p[2:4], chunk.kind)
|
||||
binary.BigEndian.PutUint16(p[4:6], uint16(len(p)))
|
||||
copy(p[6:], chunk.body)
|
||||
b = append(b, p...)
|
||||
}
|
||||
binary.BigEndian.PutUint16(b[4:6], uint16(len(b)))
|
||||
return b
|
||||
}
|
||||
|
||||
func TestApprovedRecordedRealCallReportsExactSIPResponseBeforeAnswer(t *testing.T) {
|
||||
fixture, stub, _, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
mirror := asterisk.NewHEPMirror()
|
||||
const callID = "isolated-480"
|
||||
const common = "Via: SIP/2.0/UDP 127.0.0.1;branch=z9hG4bK-real\r\nFrom: <sip:a@localhost>;tag=x\r\nTo: <sip:b@localhost>\r\nCall-ID: isolated-480\r\nCSeq: 1 INVITE\r\nContent-Length: 0\r\n\r\n"
|
||||
invite := "INVITE sip:b@localhost SIP/2.0\r\n" + common
|
||||
raw := "SIP/2.0 480 Temporarily Unavailable\r\n" + common
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Mirror: mirror, Delivery: fixture.Delivery,
|
||||
originator: func(_ context.Context, dial asterisk.NativeDial) (realMedia, error) {
|
||||
if dial.BindSIPCall == nil || dial.BindSIPCall(dial.ExecutionID, callID) != nil {
|
||||
t.Fatal("real dial did not bind the original SIP Call-ID")
|
||||
}
|
||||
for _, sip := range []string{invite, raw} {
|
||||
if err := mirror.Observe(hepRealResultFrame(callID, sip)); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
}
|
||||
return realMedia{}, &asterisk.NativeCallFailure{Phase: "answer_wait", ConfirmedEnd: true, Cause: errors.New("unanswered")}
|
||||
},
|
||||
}
|
||||
if err := runner.Run(context.Background(), approved); err == nil || strings.Join(stub.calls, ",") != "end,result" {
|
||||
t.Fatalf("verified SIP end did not reach real result: %v calls=%v", err, stub.calls)
|
||||
}
|
||||
var result struct {
|
||||
ReasonCode *int `json:"reason_code"`
|
||||
StatusLine string `json:"status_line"`
|
||||
Raw string `json:"raw"`
|
||||
}
|
||||
if err := json.Unmarshal(stub.result, &result); err != nil || result.ReasonCode == nil || *result.ReasonCode != 480 || result.StatusLine != "SIP/2.0 480 Temporarily Unavailable" || result.Raw != raw {
|
||||
t.Fatalf("original SIP 480 not delivered verbatim: %+v err=%v", result, err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestApprovedRecordedRealCallReportsVerifiedPreAnswerEndWithoutRecording(t *testing.T) {
|
||||
fixture, stub, puts, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
originate := 0
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
runner := &ApprovedRecordedRealCall{Mirror: asterisk.NewHEPMirror(),
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second, Delivery: fixture.Delivery,
|
||||
originator: func(context.Context, asterisk.NativeDial) (realMedia, error) {
|
||||
originate++
|
||||
@@ -169,11 +222,14 @@ func TestApprovedRecordedRealCallReportsVerifiedPreAnswerEndWithoutRecording(t *
|
||||
t.Fatalf("confirmed pre-answer end must report actual failure once without recording: err=%v originate=%d calls=%v puts=%d", err, originate, stub.calls, puts.Load())
|
||||
}
|
||||
var result struct {
|
||||
Outcome string `json:"outcome"`
|
||||
Recording map[string]any `json:"recording"`
|
||||
ReasonMessage string `json:"reason_message"`
|
||||
Outcome string `json:"outcome"`
|
||||
Recording map[string]any `json:"recording"`
|
||||
ReasonMessage string `json:"reason_message"`
|
||||
ReasonCode *int `json:"reason_code"`
|
||||
Raw *string `json:"raw"`
|
||||
SIPCaptureError *string `json:"sip_capture_error"`
|
||||
}
|
||||
if err := json.Unmarshal(stub.result, &result); err != nil || result.Outcome != "failed" || len(result.Recording) != 0 || result.ReasonMessage != "call ended before answer" {
|
||||
if err := json.Unmarshal(stub.result, &result); err != nil || result.Outcome != "failed" || len(result.Recording) != 0 || result.ReasonMessage != "call ended before answer" || result.ReasonCode != nil || result.Raw != nil || result.SIPCaptureError == nil || *result.SIPCaptureError == "" {
|
||||
t.Fatalf("pre-answer SIP rejection must not be recorded as answered: result=%+v err=%v", result, err)
|
||||
}
|
||||
}
|
||||
@@ -181,7 +237,7 @@ func TestApprovedRecordedRealCallReportsVerifiedPreAnswerEndWithoutRecording(t *
|
||||
func TestApprovedRecordedRealCallDoesNotReportUnverifiedPreAnswerEnd(t *testing.T) {
|
||||
fixture, stub, _, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
runner := &ApprovedRecordedRealCall{Mirror: asterisk.NewHEPMirror(),
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second, Delivery: fixture.Delivery,
|
||||
originator: func(context.Context, asterisk.NativeDial) (realMedia, error) {
|
||||
return realMedia{}, &asterisk.NativeCallFailure{Phase: "answer_wait", Cause: errors.New("channel state unavailable")}
|
||||
@@ -196,7 +252,7 @@ func TestApprovedRecordedRealCallCannotOriginateTwiceOrReportUnknownOrigination(
|
||||
fixture, stub, _, _, approved := recordedMockFixture(t, nil, time.Second)
|
||||
approved.CallerID, approved.DialedCallee, approved.RingTimeout = "BD93205882", "7089"+approved.Callee, time.Second
|
||||
originate := 0
|
||||
runner := &ApprovedRecordedRealCall{
|
||||
runner := &ApprovedRecordedRealCall{Mirror: asterisk.NewHEPMirror(),
|
||||
MaxWAVBytes: 4096, MediaPayloadType: 118, ReportTimeout: time.Second,
|
||||
Delivery: fixture.Delivery,
|
||||
originator: func(context.Context, asterisk.NativeDial) (realMedia, error) {
|
||||
|
||||
Reference in New Issue
Block a user