Files
go-sip/internal/sipcall/capture.go
T

266 lines
8.9 KiB
Go

package sipcall
import (
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"io"
"os"
"os/exec"
"path/filepath"
"regexp"
"strconv"
"strings"
"syscall"
"time"
)
type Evidence struct {
PCAP string `json:"pcap"`
PCAPSHA256 string `json:"pcap_sha256"`
SIPTimeline string `json:"sip_timeline,omitempty"`
Packets int `json:"sip_packets"`
FinalINVITECode int `json:"final_invite_code,omitempty"`
FinalINVITEStatus string `json:"final_invite_status,omitempty"`
RTPFromServer int `json:"rtp_from_server"`
RTPToServer int `json:"rtp_to_server"`
RTPBytes int `json:"rtp_bytes"`
MediaScope string `json:"media_scope"`
RTPStatus string `json:"rtp_status"`
}
type packetCapture struct {
cmd *exec.Cmd
done chan error
log *os.File
dir, tshark string
config Config
}
var sipCallID = regexp.MustCompile(`^[A-Za-z0-9_.@:+-]{1,256}$`)
func startCapture(ctx context.Context, c Config, dir string) (captureSession, error) {
tcpdump, err := exec.LookPath("tcpdump")
if err != nil {
return nil, errors.New("--debug requires tcpdump; no call was dialed")
}
tshark, err := exec.LookPath("tshark")
if err != nil {
return nil, errors.New("--debug requires tshark for call-linked evidence; no call was dialed")
}
log, err := os.OpenFile(filepath.Join(dir, "tcpdump.log"), os.O_CREATE|os.O_EXCL|os.O_WRONLY, 0600)
if err != nil {
return nil, err
}
// The narrow filter avoids copying unrelated UDP traffic. RTP sent from a
// different provider media IP is explicitly outside this evidence's scope.
cmd := exec.Command(tcpdump, "-i", c.Interface, "-nn", "-U", "-s", "0", "-w", filepath.Join(dir, "capture.pcap"), "udp and host "+c.Server)
cmd.Stdout = log
cmd.Stderr = log
if err = cmd.Start(); err != nil {
return nil, errors.Join(fmt.Errorf("start tcpdump: %w", err), log.Close())
}
cap := &packetCapture{cmd: cmd, done: make(chan error, 1), log: log, dir: dir, tshark: tshark, config: c}
go func() { cap.done <- cmd.Wait(); close(cap.done) }()
deadline := time.NewTimer(5 * time.Second)
defer deadline.Stop()
tick := time.NewTicker(25 * time.Millisecond)
defer tick.Stop()
for {
select {
case e := <-cap.done:
return nil, errors.Join(errors.New("tcpdump stopped before readiness; inspect private tcpdump.log"), e, log.Close())
case <-ctx.Done():
_, stopErr := cap.Stop(context.Background(), "")
return nil, errors.Join(ctx.Err(), stopErr)
case <-deadline.C:
_, stopErr := cap.Stop(context.Background(), "")
return nil, errors.Join(errors.New("tcpdump readiness timed out; no call was dialed"), stopErr)
case <-tick.C:
text, e := os.ReadFile(filepath.Join(dir, "tcpdump.log"))
if e != nil {
continue
}
info, e := os.Stat(filepath.Join(dir, "capture.pcap"))
if e == nil && info.Size() >= 24 && bytes.Contains(text, []byte("listening on ")) {
if e = os.Chmod(filepath.Join(dir, "capture.pcap"), 0600); e != nil {
_, stopErr := cap.Stop(context.Background(), "")
return nil, errors.Join(e, stopErr)
}
return cap, nil
}
}
}
}
func (c *packetCapture) Check() error {
select {
case <-c.done:
return errors.New("tcpdump is no longer active; no new Dial is permitted")
default:
return nil
}
}
func (c *packetCapture) Stop(ctx context.Context, callID string) (e Evidence, err error) {
e = Evidence{PCAP: filepath.Join(c.dir, "capture.pcap"), MediaScope: "UDP to/from configured SIP server IP only; different RTP media IPs are not covered", RTPStatus: "not_observed_within_capture_scope"}
select {
case waitErr := <-c.done:
err = errors.Join(errors.New("tcpdump exited before capture finalization"), waitErr)
default:
signalErr := c.cmd.Process.Signal(syscall.SIGINT)
if signalErr != nil {
err = errors.Join(err, signalErr)
}
timeout := time.NewTimer(5 * time.Second)
defer timeout.Stop()
select {
case waitErr := <-c.done:
err = errors.Join(err, waitErr)
case <-ctx.Done():
err = errors.Join(err, ctx.Err(), c.cmd.Process.Kill(), <-c.done)
case <-timeout.C:
err = errors.Join(err, errors.New("tcpdump did not stop within 5 seconds"), c.cmd.Process.Kill(), <-c.done)
}
}
err = errors.Join(err, c.log.Close())
file, readErr := os.Open(e.PCAP)
if readErr == nil {
hash := sha256.New()
_, readErr = io.Copy(hash, file)
readErr = errors.Join(readErr, file.Close())
if readErr == nil {
e.PCAPSHA256 = hex.EncodeToString(hash.Sum(nil))
}
}
err = errors.Join(err, readErr)
if callID == "" {
return e, err
}
if !sipCallID.MatchString(callID) {
return e, errors.Join(err, errors.New("native SIP Call-ID cannot be safely correlated"))
}
args := []string{"-n", "-r", e.PCAP, "-d", fmt.Sprintf("udp.port==%d,sip", c.config.Port), "-Y", `sip.Call-ID == "` + callID + `"`, "-T", "fields", "-E", "separator=/t", "-E", "occurrence=f"}
for _, field := range []string{"frame.number", "frame.time_epoch", "ip.src", "ip.dst", "sip.Request-Line", "sip.Status-Code", "sip.CSeq.method", "sip.CSeq.seq", "sip.Via.branch", "sip.Status-Line", "sdp.connection_info.address", "sdp.media.media", "sdp.media.port"} {
args = append(args, "-e", field)
}
timeline, decodeErr := runTshark(ctx, c.tshark, args...)
err = errors.Join(err, decodeErr)
e.SIPTimeline = filepath.Join(c.dir, "sip.tsv")
err = errors.Join(err, os.WriteFile(e.SIPTimeline, timeline, 0600))
if decodeErr != nil {
return e, err
}
parsed, ports, parseErr := parseTimeline(timeline)
e.Packets = parsed.Packets
e.FinalINVITECode = parsed.FinalINVITECode
e.FinalINVITEStatus = parsed.FinalINVITEStatus
err = errors.Join(err, parseErr)
if len(ports) == 0 {
return e, err
}
filters := []string{}
for _, port := range ports {
filters = append(filters, fmt.Sprintf("udp.port == %d", port))
}
rtpArgs := []string{"-n", "-r", e.PCAP, "-d", fmt.Sprintf("udp.port==%d,sip", c.config.Port)}
for _, port := range ports {
rtpArgs = append(rtpArgs, "-d", fmt.Sprintf("udp.port==%d,rtp", port))
}
rtpArgs = append(rtpArgs, "-Y", "rtp && ("+strings.Join(filters, " || ")+")", "-T", "fields", "-E", "separator=/t", "-e", "frame.len", "-e", "ip.src", "-e", "ip.dst")
rtp, rtpErr := runTshark(ctx, c.tshark, rtpArgs...)
err = errors.Join(err, rtpErr, os.WriteFile(filepath.Join(c.dir, "rtp.tsv"), rtp, 0600))
for _, line := range strings.Split(strings.TrimSpace(string(rtp)), "\n") {
if line == "" {
continue
}
fields := strings.Split(line, "\t")
if len(fields) != 3 {
return e, errors.Join(err, errors.New("unexpected tshark RTP fields"))
}
n, convErr := strconv.Atoi(fields[0])
if convErr != nil {
return e, errors.Join(err, convErr)
}
e.RTPBytes += n
if fields[1] == c.config.Server {
e.RTPFromServer++
}
if fields[2] == c.config.Server {
e.RTPToServer++
}
}
if e.RTPFromServer+e.RTPToServer > 0 {
e.RTPStatus = "observed_packets_only; no audio intelligibility or AI test"
}
return e, err
}
func runTshark(ctx context.Context, binary string, args ...string) ([]byte, error) {
bounded, cancel := context.WithTimeout(ctx, 10*time.Second)
defer cancel()
cmd := exec.CommandContext(bounded, binary, args...)
var diagnostic bytes.Buffer
cmd.Stderr = &diagnostic
out, err := cmd.Output()
if err != nil {
for i, arg := range args {
if arg == "-r" && i+1 < len(args) {
err = errors.Join(err, os.WriteFile(filepath.Join(filepath.Dir(args[i+1]), "tshark.log"), diagnostic.Bytes(), 0600))
break
}
}
return nil, fmt.Errorf("tshark evidence extraction failed; inspect private tshark.log: %w", err)
}
return out, nil
}
// parseTimeline consumes decoded tshark fields, not a hand-written SIP parser.
// Final responses are tied to the original INVITE's CSeq and Via transaction.
func parseTimeline(data []byte) (Evidence, []int, error) {
e := Evidence{}
var rows [][]string
var seq, branch string
ports := []int{}
seen := map[int]bool{}
for _, line := range strings.Split(strings.TrimRight(string(data), "\n"), "\n") {
if line == "" {
continue
}
fields := strings.Split(line, "\t")
if len(fields) != 13 {
return e, nil, errors.New("unexpected tshark SIP fields")
}
rows = append(rows, fields)
e.Packets++
if seq == "" && strings.HasPrefix(fields[4], "INVITE ") {
seq, branch = fields[7], fields[8]
}
if fields[11] == "audio" {
p, err := strconv.Atoi(fields[12])
if err == nil && p > 0 && p <= 65535 && !seen[p] {
seen[p] = true
ports = append(ports, p)
}
}
}
if seq == "" || branch == "" {
return e, ports, errors.New("capture contains no identifiable original INVITE for this SIP Call-ID")
}
for _, fields := range rows {
if fields[6] != "INVITE" || fields[7] != seq || fields[8] != branch {
continue
}
code, convErr := strconv.Atoi(fields[5])
if convErr != nil || code < 200 {
continue
}
if e.FinalINVITECode != 0 && e.FinalINVITECode != code {
return e, ports, errors.New("conflicting final responses for original INVITE transaction")
}
e.FinalINVITECode = code
e.FinalINVITEStatus = fields[9]
}
return e, ports, nil
}