From 95eeee8400808115061377608725a3217cd635e9 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 5 Oct 2026 21:59:46 +0800 Subject: [PATCH] Capture original INVITE responses through local HEP mirror --- cmd/sip-go-agent/agent_command.go | 36 ++- cmd/sip-go-agent/agent_real.go | 6 +- cmd/sip-go-agent/agent_real_test.go | 17 +- go.mod | 6 + go.sum | 11 + internal/asterisk/hep.go | 230 ++++++++++++++++++++ internal/asterisk/hep_listener.go | 52 +++++ internal/asterisk/hep_listener_test.go | 59 +++++ internal/asterisk/hep_test.go | 140 ++++++++++++ internal/config/agent_runtime.go | 7 + internal/config/agent_runtime_test.go | 14 +- internal/rpc/approved_recorded_real.go | 37 +++- internal/rpc/approved_recorded_real_test.go | 78 ++++++- 13 files changed, 663 insertions(+), 30 deletions(-) create mode 100644 internal/asterisk/hep.go create mode 100644 internal/asterisk/hep_listener.go create mode 100644 internal/asterisk/hep_listener_test.go create mode 100644 internal/asterisk/hep_test.go diff --git a/cmd/sip-go-agent/agent_command.go b/cmd/sip-go-agent/agent_command.go index 5b3d1ab..f3e7278 100644 --- a/cmd/sip-go-agent/agent_command.go +++ b/cmd/sip-go-agent/agent_command.go @@ -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) } diff --git a/cmd/sip-go-agent/agent_real.go b/cmd/sip-go-agent/agent_real.go index 2ce2ccc..1538fb3 100644 --- a/cmd/sip-go-agent/agent_real.go +++ b/cmd/sip-go-agent/agent_real.go @@ -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 { diff --git a/cmd/sip-go-agent/agent_real_test.go b/cmd/sip-go-agent/agent_real_test.go index a7b6bc2..b243b5d 100644 --- a/cmd/sip-go-agent/agent_real_test.go +++ b/cmd/sip-go-agent/agent_real_test.go @@ -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") } } diff --git a/go.mod b/go.mod index 5b06d66..8ad1c08 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 93ade87..c7b7728 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/internal/asterisk/hep.go b/internal/asterisk/hep.go new file mode 100644 index 0000000..785d9df --- /dev/null +++ b/internal/asterisk/hep.go @@ -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 +} diff --git a/internal/asterisk/hep_listener.go b/internal/asterisk/hep_listener.go new file mode 100644 index 0000000..c006f8d --- /dev/null +++ b/internal/asterisk/hep_listener.go @@ -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) + } + } +} diff --git a/internal/asterisk/hep_listener_test.go b/internal/asterisk/hep_listener_test.go new file mode 100644 index 0000000..48aa405 --- /dev/null +++ b/internal/asterisk/hep_listener_test.go @@ -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") + } +} diff --git a/internal/asterisk/hep_test.go b/internal/asterisk/hep_test.go new file mode 100644 index 0000000..a1a3c7a --- /dev/null +++ b/internal/asterisk/hep_test.go @@ -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: ;tag=a\r\nTo: \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) + } +} diff --git a/internal/config/agent_runtime.go b/internal/config/agent_runtime.go index d248d4f..2ea6707 100644 --- a/internal/config/agent_runtime.go +++ b/internal/config/agent_runtime.go @@ -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") } diff --git a/internal/config/agent_runtime_test.go b/internal/config/agent_runtime_test.go index bf0ddd8..98c4f1b 100644 --- a/internal/config/agent_runtime_test.go +++ b/internal/config/agent_runtime_test.go @@ -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]) }) diff --git a/internal/rpc/approved_recorded_real.go b/internal/rpc/approved_recorded_real.go index 4bbb6ca..5e3cb6d 100644 --- a/internal/rpc/approved_recorded_real.go +++ b/internal/rpc/approved_recorded_real.go @@ -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) } diff --git a/internal/rpc/approved_recorded_real_test.go b/internal/rpc/approved_recorded_real_test.go index 7cb37c4..88d8cd1 100644 --- a/internal/rpc/approved_recorded_real_test.go +++ b/internal/rpc/approved_recorded_real_test.go @@ -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: ;tag=x\r\nTo: \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) {