From 74c7b9317086782d5678b077e686dbd0cef6f72f Mon Sep 17 00:00:00 2001 From: Denys Smirnov Date: Mon, 17 Jun 2024 21:49:51 +0300 Subject: [PATCH] Support new SIP Trunk API. Improve Redis tests. (#2799) --- go.mod | 25 +- go.sum | 76 +++- pkg/service/docker_test.go | 79 ++++ pkg/service/interfaces.go | 6 + pkg/service/ioservice_sip.go | 10 +- pkg/service/redisstore.go | 121 ++---- pkg/service/redisstore_sip.go | 199 +++++++++ pkg/service/redisstore_sip_test.go | 263 ++++++++++++ pkg/service/redisstore_test.go | 15 +- pkg/service/servicefakes/fake_sipstore.go | 472 ++++++++++++++++++++++ pkg/service/sip.go | 91 ++++- pkg/service/utils_test.go | 34 +- 12 files changed, 1288 insertions(+), 103 deletions(-) create mode 100644 pkg/service/docker_test.go create mode 100644 pkg/service/redisstore_sip.go create mode 100644 pkg/service/redisstore_sip_test.go diff --git a/go.mod b/go.mod index b85a91cda..d80c3a640 100644 --- a/go.mod +++ b/go.mod @@ -20,13 +20,14 @@ require ( github.com/jxskiss/base62 v1.1.0 github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 github.com/livekit/mediatransportutil v0.0.0-20240613015318-84b69facfb75 - github.com/livekit/protocol v1.17.1-0.20240614060801-425cb974f7a4 + github.com/livekit/protocol v1.17.1-0.20240617184219-32c577d805ed github.com/livekit/psrpc v0.5.3-0.20240526192918-fbdaf10e6aa5 github.com/mackerelio/go-osstat v0.2.5 github.com/magefile/mage v1.15.0 github.com/maxbrunsfeld/counterfeiter/v6 v6.8.1 github.com/mitchellh/go-homedir v1.1.0 github.com/olekukonko/tablewriter v0.0.5 + github.com/ory/dockertest/v3 v3.10.0 github.com/pion/dtls/v2 v2.2.11 github.com/pion/ice/v2 v2.3.24 github.com/pion/interceptor v0.1.29 @@ -57,18 +58,30 @@ require ( ) require ( + dario.cat/mergo v1.0.0 // indirect + github.com/Azure/go-ansiterm v0.0.0-20230124172434-306776ec8161 // indirect + github.com/Microsoft/go-winio v0.6.2 // indirect + github.com/Nvveen/Gotty v0.0.0-20120604004816-cd527374f1e5 // indirect github.com/benbjohnson/clock v1.3.5 // indirect github.com/beorn7/perks v1.0.1 // indirect + github.com/cenkalti/backoff/v4 v4.3.0 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/containerd/continuity v0.4.3 // indirect github.com/cpuguy83/go-md2man/v2 v2.0.4 // indirect github.com/davecgh/go-spew v1.1.1 // indirect github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect + github.com/docker/cli v26.1.4+incompatible // indirect + github.com/docker/docker v27.0.0+incompatible // indirect + github.com/docker/go-connections v0.5.0 // indirect + github.com/docker/go-units v0.5.0 // indirect github.com/eapache/channels v1.1.0 // indirect github.com/eapache/queue v1.1.0 // indirect github.com/fsnotify/fsnotify v1.7.0 // indirect github.com/go-jose/go-jose/v3 v3.0.3 // indirect github.com/go-logr/logr v1.4.2 // indirect + github.com/gogo/protobuf v1.3.2 // indirect github.com/google/go-cmp v0.6.0 // indirect + github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510 // indirect github.com/google/subcommands v1.2.0 // indirect github.com/google/uuid v1.6.0 // indirect github.com/hashicorp/go-cleanhttp v0.5.2 // indirect @@ -81,9 +94,15 @@ require ( github.com/mattn/go-runewidth v0.0.9 // indirect github.com/mdlayher/netlink v1.7.1 // indirect github.com/mdlayher/socket v0.4.0 // indirect + github.com/mitchellh/mapstructure v1.5.0 // indirect + github.com/moby/docker-image-spec v1.3.1 // indirect + github.com/moby/term v0.5.0 // indirect github.com/nats-io/nats.go v1.35.0 // indirect github.com/nats-io/nkeys v0.4.7 // indirect github.com/nats-io/nuid v1.0.1 // indirect + github.com/opencontainers/go-digest v1.0.0 // indirect + github.com/opencontainers/image-spec v1.1.0 // indirect + github.com/opencontainers/runc v1.1.13 // indirect github.com/pion/datachannel v1.5.5 // indirect github.com/pion/logging v0.2.2 // indirect github.com/pion/mdns v0.0.12 // indirect @@ -96,6 +115,10 @@ require ( github.com/prometheus/procfs v0.12.0 // indirect github.com/puzpuzpuz/xsync/v3 v3.1.0 // indirect github.com/russross/blackfriday/v2 v2.1.0 // indirect + github.com/sirupsen/logrus v1.9.3 // indirect + github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb // indirect + github.com/xeipuuv/gojsonreference v0.0.0-20180127040603-bd5ef7bd5415 // indirect + github.com/xeipuuv/gojsonschema v1.2.0 // indirect github.com/xrash/smetrics v0.0.0-20240312152122-5f08fbb34913 // indirect github.com/zeebo/xxh3 v1.0.2 // indirect go.uber.org/zap/exp v0.2.0 // indirect diff --git a/go.sum b/go.sum index f42fa94c8..061c3305a 100644 --- a/go.sum +++ b/go.sum @@ -1,3 +1,11 @@ +dario.cat/mergo v1.0.0 h1:AGCNq9Evsj31mOgNPcLyXc+4PNABt905YmuqPYYpBWk= +dario.cat/mergo v1.0.0/go.mod h1:uNxQE+84aUszobStD9th8a29P2fMDhsBdgRYvZOxGmk= +github.com/Azure/go-ansiterm v0.0.0-20230124172434-306776ec8161 h1:L/gRVlceqvL25UVaW/CKtUDjefjrs0SPonmDGUVOYP0= +github.com/Azure/go-ansiterm v0.0.0-20230124172434-306776ec8161/go.mod h1:xomTg63KZ2rFqZQzSB4Vz2SUXa1BpHTVz9L5PTmPC4E= +github.com/Microsoft/go-winio v0.6.2 h1:F2VQgta7ecxGYO8k3ZZz3RS8fVIXVxONVUPlNERoyfY= +github.com/Microsoft/go-winio v0.6.2/go.mod h1:yd8OoFMLzJbo9gZq8j5qaps8bJ9aShtEA8Ipt1oGCvU= +github.com/Nvveen/Gotty v0.0.0-20120604004816-cd527374f1e5 h1:TngWCqHvy9oXAN6lEVMRuU21PR1EtLVZJmdB18Gu3Rw= +github.com/Nvveen/Gotty v0.0.0-20120604004816-cd527374f1e5/go.mod h1:lmUJ/7eu/Q8D7ML55dXQrVaamCz2vxCfdQBasLZfHKk= github.com/avast/retry-go/v4 v4.6.0 h1:K9xNA+KeB8HHc2aWFuLb25Offp+0iVRXEvFx8IinRJA= github.com/avast/retry-go/v4 v4.6.0/go.mod h1:gvWlPhBVsvBbLkVGDg/KwvBv0bEkCOLRRSHKIr2PyOE= github.com/benbjohnson/clock v1.3.5 h1:VvXlSJBzZpA/zum6Sj74hxwYI2DIxRWuNIoXAzHZz5o= @@ -10,15 +18,21 @@ github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= +github.com/cenkalti/backoff/v4 v4.3.0 h1:MyRJ/UdXutAwSAT+s3wNd7MfTIcy71VQueUuFK343L8= +github.com/cenkalti/backoff/v4 v4.3.0/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= github.com/cilium/ebpf v0.5.0/go.mod h1:4tRaxcgiL706VnOzHOdBlY8IEAIdxINsQBcU4xJJXRs= github.com/cilium/ebpf v0.7.0/go.mod h1:/oI2+1shJiTGAMgl6/RgJr36Eo1jzrRcAWbcXO2usCA= github.com/cilium/ebpf v0.8.1 h1:bLSSEbBLqGPXxls55pGr5qWZaTqcmfDJHhou7t254ao= github.com/cilium/ebpf v0.8.1/go.mod h1:f5zLIM0FSNuAkSyLAN7X+Hy6yznlF1mNiWUMfxMtrgk= +github.com/containerd/continuity v0.4.3 h1:6HVkalIp+2u1ZLH1J/pYX2oBVXlJZvh1X1A7bEZ9Su8= +github.com/containerd/continuity v0.4.3/go.mod h1:F6PTNCKepoxEaXLQp3wDAjygEnImnZ/7o4JzpodfroQ= github.com/cpuguy83/go-md2man/v2 v2.0.4 h1:wfIWP927BUkWJb2NmU/kNDYIBTh/ziUX91+lVfRxZq4= github.com/cpuguy83/go-md2man/v2 v2.0.4/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= +github.com/creack/pty v1.1.18 h1:n56/Zwd5o6whRC5PMGretI4IdRLlmBXYNjScPaBgsbY= +github.com/creack/pty v1.1.18/go.mod h1:MOBLtS5ELjhRRrroQr9kyvTxUAFNvYEK993ew/Vr4O4= github.com/d5/tengo/v2 v2.17.0 h1:BWUN9NoJzw48jZKiYDXDIF3QrIVZRm1uV1gTzeZ2lqM= github.com/d5/tengo/v2 v2.17.0/go.mod h1:XRGjEs5I9jYIKTxly6HCF8oiiilk5E/RYXOZ5b0DZC8= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -26,6 +40,14 @@ github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78= github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= +github.com/docker/cli v26.1.4+incompatible h1:I8PHdc0MtxEADqYJZvhBrW9bo8gawKwwenxRM7/rLu8= +github.com/docker/cli v26.1.4+incompatible/go.mod h1:JLrzqnKDaYBop7H2jaqPtU4hHvMKP+vjCwu2uszcLI8= +github.com/docker/docker v27.0.0+incompatible h1:JRugTYuelmWlW0M3jakcIadDx2HUoUO6+Tf2C5jVfwA= +github.com/docker/docker v27.0.0+incompatible/go.mod h1:eEKB0N0r5NX/I1kEveEz05bcu8tLC/8azJZsviup8Sk= +github.com/docker/go-connections v0.5.0 h1:USnMq7hx7gwdVZq1L49hLXaFtUdTADjXGp+uj1Br63c= +github.com/docker/go-connections v0.5.0/go.mod h1:ov60Kzw0kKElRwhNs9UlUHAE/F9Fe6GLaXnqyDdmEXc= +github.com/docker/go-units v0.5.0 h1:69rxXcBk27SvSaaxTtLh/8llcHD8vYHT7WSdRZ/jvr4= +github.com/docker/go-units v0.5.0/go.mod h1:fgPhTUdO+D/Jk86RDLlptpiXQzgHJF7gydDDbaIK4Dk= github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= github.com/eapache/channels v1.1.0 h1:F1taHcn7/F0i8DYqKXJnyhJcVpp2kgFcNePxXtnyu4k= @@ -50,6 +72,10 @@ github.com/go-jose/go-jose/v3 v3.0.3 h1:fFKWeig/irsp7XD2zBxvnmA/XaRWp5V3CBsZXJF7 github.com/go-jose/go-jose/v3 v3.0.3/go.mod h1:5b+7YgP7ZICgJDBdfjZaIt+H/9L9T/YQrVfLAMboGkQ= github.com/go-logr/logr v1.4.2 h1:6pFjapn8bFcIbiKo3XT4j/BhANplGihG6tvd+8rYgrY= github.com/go-logr/logr v1.4.2/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-sql-driver/mysql v1.6.0 h1:BCTh4TKNUYmOmMUcQ3IipzF5prigylS7XXjEkfCHuOE= +github.com/go-sql-driver/mysql v1.6.0/go.mod h1:DCzpHaOWr8IXmIStZouvnhqoel9Qv2LBy8hT2VhHyBg= +github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= +github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M= github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU= github.com/google/go-cmp v0.4.0/go.mod h1:v8dTdLbMG2kIc/vJvl+f65V22dbkXbowE6jgT/gNBxE= @@ -61,6 +87,8 @@ github.com/google/go-cmp v0.5.7/go.mod h1:n+brtR0CgQNWTVd5ZUFpTBC8YFBDLK/h/bpaJ8 github.com/google/go-cmp v0.5.9/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= github.com/google/go-cmp v0.6.0 h1:ofyhxvXcZhMsU5ulbFiLKl/XBFqE1GSq7atu8tAmTRI= github.com/google/go-cmp v0.6.0/go.mod h1:17dUlkBOakJ0+DkrSSNjCkIjxS6bF9zb3elmeNGIjoY= +github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510 h1:El6M4kTTCOh6aBiKaUGG7oYTSPP8MxqL4YI3kZKwcP4= +github.com/google/shlex v0.0.0-20191202100458-e7afc7fbc510/go.mod h1:pupxD2MaaD3pAXIBCelhxNneeOaAeabZDe5s4K6zSpQ= github.com/google/subcommands v1.2.0 h1:vWQspBTo2nEqTUFita5/KeEWlUL8kQObDFbub/EN9oE= github.com/google/subcommands v1.2.0/go.mod h1:ZjhPrFU+Olkh9WazFPsl27BQ4UPiG37m3yTrtFlrHVk= github.com/google/uuid v1.3.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= @@ -101,6 +129,8 @@ github.com/jsimonetti/rtnetlink v0.0.0-20211022192332-93da33804786 h1:N527AHMa79 github.com/jsimonetti/rtnetlink v0.0.0-20211022192332-93da33804786/go.mod h1:v4hqbTdfQngbVSZJVWUhGE/lbTFf9jb+ygmNUDQMuOs= github.com/jxskiss/base62 v1.1.0 h1:A5zbF8v8WXx2xixnAKD2w+abC+sIzYJX+nxmhA6HWFw= github.com/jxskiss/base62 v1.1.0/go.mod h1:HhWAlUXvxKThfOlZbcuFzsqwtF5TcqS9ru3y5GfjWAc= +github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= +github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= github.com/klauspost/compress v1.17.8 h1:YcnTYrq7MikUT7k0Yb5eceMmALQPYBW/Xltxn0NAMnU= github.com/klauspost/compress v1.17.8/go.mod h1:Di0epgTjJY877eYKx5yC51cX2A2Vl2ibi7bDH9ttBbw= github.com/klauspost/cpuid/v2 v2.2.6 h1:ndNyv040zDGIDh8thGkXYjnFtiN02M1PVVF+JE/48xc= @@ -114,14 +144,16 @@ github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/lib/pq v0.0.0-20180327071824-d34b9ff171c2 h1:hRGSmZu7j271trc9sneMrpOW7GN5ngLm8YUZIPzf394= +github.com/lib/pq v0.0.0-20180327071824-d34b9ff171c2/go.mod h1:5WUZQaWbwv1U+lTReE5YruASi9Al49XbQIvNi/34Woo= github.com/lithammer/shortuuid/v4 v4.0.0 h1:QRbbVkfgNippHOS8PXDkti4NaWeyYfcBTHtw7k08o4c= github.com/lithammer/shortuuid/v4 v4.0.0/go.mod h1:Zs8puNcrvf2rV9rTH51ZLLcj7ZXqQI3lv67aw4KiB1Y= github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1 h1:jm09419p0lqTkDaKb5iXdynYrzB84ErPPO4LbRASk58= github.com/livekit/mageutil v0.0.0-20230125210925-54e8a70427c1/go.mod h1:Rs3MhFwutWhGwmY1VQsygw28z5bWcnEYmS1OG9OxjOQ= github.com/livekit/mediatransportutil v0.0.0-20240613015318-84b69facfb75 h1:p60OjeixzXnhGFQL8wmdUwWPxijEDe9ZJFMosq+byec= github.com/livekit/mediatransportutil v0.0.0-20240613015318-84b69facfb75/go.mod h1:jwKUCmObuiEDH0iiuJHaGMXwRs3RjrB4G6qqgkr/5oE= -github.com/livekit/protocol v1.17.1-0.20240614060801-425cb974f7a4 h1:5O3/wahQIMnw+PO3O1kgxPSrpJf7PkPcD3GvWmb1TJQ= -github.com/livekit/protocol v1.17.1-0.20240614060801-425cb974f7a4/go.mod h1:cN8WmGQR+kWz1+UWcAQdFFUcbW76PnfZDdkLAbYIqd4= +github.com/livekit/protocol v1.17.1-0.20240617184219-32c577d805ed h1:S4avs1NKG6bBgHYuBOrQWnNxJSOdunGOB84BQfGzKmQ= +github.com/livekit/protocol v1.17.1-0.20240617184219-32c577d805ed/go.mod h1:cN8WmGQR+kWz1+UWcAQdFFUcbW76PnfZDdkLAbYIqd4= github.com/livekit/psrpc v0.5.3-0.20240526192918-fbdaf10e6aa5 h1:mTZyrjk5WEWMsvaYtJ42pG7DuxysKj21DKPINpGSIto= github.com/livekit/psrpc v0.5.3-0.20240526192918-fbdaf10e6aa5/go.mod h1:CQUBSPfYYAaevg1TNCc6/aYsa8DJH4jSRFdCeSZk5u0= github.com/mackerelio/go-osstat v0.2.5 h1:+MqTbZUhoIt4m8qzkVoXUJg1EuifwlAJSk4Yl2GXh+o= @@ -153,6 +185,12 @@ github.com/mdlayher/socket v0.4.0 h1:280wsy40IC9M9q1uPGcLBwXpcTQDtoGwVt+BNoITxIw github.com/mdlayher/socket v0.4.0/go.mod h1:xxFqz5GRCUN3UEOm9CZqEJsAbe1C8OwSK46NlmWuVoc= github.com/mitchellh/go-homedir v1.1.0 h1:lukF9ziXFxDFPkA1vsr5zpc1XuPDn/wFntq5mG+4E0Y= github.com/mitchellh/go-homedir v1.1.0/go.mod h1:SfyaCUpYCn1Vlf4IUYiD9fPX4A5wJrkLzIz1N1q0pr0= +github.com/mitchellh/mapstructure v1.5.0 h1:jeMsZIYE/09sWLaz43PL7Gy6RuMjD2eJVyuac5Z2hdY= +github.com/mitchellh/mapstructure v1.5.0/go.mod h1:bFUtVrKA4DC2yAKiSyO/QUcy7e+RRV2QTWOzhPopBRo= +github.com/moby/docker-image-spec v1.3.1 h1:jMKff3w6PgbfSa69GfNg+zN/XLhfXJGnEx3Nl2EsFP0= +github.com/moby/docker-image-spec v1.3.1/go.mod h1:eKmb5VW8vQEh/BAr2yvVNvuiJuY6UIocYsFu/DxxRpo= +github.com/moby/term v0.5.0 h1:xt8Q1nalod/v7BqbG21f8mQPqH+xAaC9C3N3wfWbVP0= +github.com/moby/term v0.5.0/go.mod h1:8FzsFHVUBGZdbDsJw/ot+X+d5HLUbvklYLJ9uGfcI3Y= github.com/nats-io/nats.go v1.35.0 h1:XFNqNM7v5B+MQMKqVGAyHwYhyKb48jrenXNxIU20ULk= github.com/nats-io/nats.go v1.35.0/go.mod h1:Ubdu4Nh9exXdSz0RVWRFBbRfrbSxOYd26oF0wkWclB8= github.com/nats-io/nkeys v0.4.7 h1:RwNJbbIdYCoClSDNY7QVKZlyb/wfT6ugvFCiKy6vDvI= @@ -163,6 +201,14 @@ github.com/olekukonko/tablewriter v0.0.5 h1:P2Ga83D34wi1o9J6Wh1mRuqd4mF/x/lgBS7N github.com/olekukonko/tablewriter v0.0.5/go.mod h1:hPp6KlRPjbx+hW8ykQs1w3UBbZlj6HuIJcUGPhkA7kY= github.com/onsi/gomega v1.30.0 h1:hvMK7xYz4D3HapigLTeGdId/NcfQx1VHMJc60ew99+8= github.com/onsi/gomega v1.30.0/go.mod h1:9sxs+SwGrKI0+PWe4Fxa9tFQQBG5xSsSbMXOI8PPpoQ= +github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= +github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= +github.com/opencontainers/image-spec v1.1.0 h1:8SG7/vwALn54lVB/0yZ/MMwhFrPYtpEHQb2IpWsCzug= +github.com/opencontainers/image-spec v1.1.0/go.mod h1:W4s4sFTMaBeK1BQLXbG4AdM2szdn85PY75RI83NrTrM= +github.com/opencontainers/runc v1.1.13 h1:98S2srgG9vw0zWcDpFMn5TRrh8kLxa/5OFUstuUhmRs= +github.com/opencontainers/runc v1.1.13/go.mod h1:R016aXacfp/gwQBYw2FDGa9m+n6atbLWrYY8hNMT/sA= +github.com/ory/dockertest/v3 v3.10.0 h1:4K3z2VMe8Woe++invjaTB7VRyQXQy5UY+loujO4aNE4= +github.com/ory/dockertest/v3 v3.10.0/go.mod h1:nr57ZbRWMqfsdGdFNLHz5jjNdDb7VVFnzAeW1n5N1Lg= github.com/pion/datachannel v1.5.5 h1:10ef4kwdjije+M9d7Xm9im2Y3O6A6ccQb0zcqZcJew8= github.com/pion/datachannel v1.5.5/go.mod h1:iMz+lECmfdCMqFRhXhcA/219B0SQlbpoR2V118yimL0= github.com/pion/dtls/v2 v2.2.7/go.mod h1:8WiMkebSHFD0T+dIU+UeBaoV7kDhOW5oDCzZ7WZ/F9s= @@ -234,11 +280,14 @@ github.com/russross/blackfriday/v2 v2.1.0 h1:JIOH55/0cWyOuilr9/qlrm0BSXldqnqwMsf github.com/russross/blackfriday/v2 v2.1.0/go.mod h1:+Rmxgy9KzJVeS9/2gXHxylqXiyQDYRxCVz55jmeOWTM= github.com/sclevine/spec v1.4.0 h1:z/Q9idDcay5m5irkZ28M7PtQM4aOISzOpj4bUPkDee8= github.com/sclevine/spec v1.4.0/go.mod h1:LvpgJaFyvQzRvc1kaDs0bulYwzC70PbiYjC4QnFHkOM= +github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= +github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= github.com/stretchr/objx v0.5.2/go.mod h1:FRsXN1f5AsAjCGJKqEizvkpNtU+EGNCLh3NxZ/8L+MA= github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXfy6kDkUVs= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= github.com/stretchr/testify v1.7.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= @@ -258,8 +307,17 @@ github.com/urfave/cli/v2 v2.27.2 h1:6e0H+AkS+zDckwPCUrZkKX38mRaau4nL2uipkJpbkcI= github.com/urfave/cli/v2 v2.27.2/go.mod h1:g0+79LmHHATl7DAcHO99smiR/T7uGLw84w8Y42x+4eM= github.com/urfave/negroni/v3 v3.1.1 h1:6MS4nG9Jk/UuCACaUlNXCbiKa0ywF9LXz5dGu09v8hw= github.com/urfave/negroni/v3 v3.1.1/go.mod h1:jWvnX03kcSjDBl/ShB0iHvx5uOs7mAzZXW+JvJ5XYAs= +github.com/xeipuuv/gojsonpointer v0.0.0-20180127040702-4e3ac2762d5f/go.mod h1:N2zxlSyiKSe5eX1tZViRH5QA0qijqEDrYZiPEAiq3wU= +github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb h1:zGWFAtiMcyryUHoUjUJX0/lt1H2+i2Ka2n+D3DImSNo= +github.com/xeipuuv/gojsonpointer v0.0.0-20190905194746-02993c407bfb/go.mod h1:N2zxlSyiKSe5eX1tZViRH5QA0qijqEDrYZiPEAiq3wU= +github.com/xeipuuv/gojsonreference v0.0.0-20180127040603-bd5ef7bd5415 h1:EzJWgHovont7NscjpAxXsDA8S8BMYve8Y5+7cuRE7R0= +github.com/xeipuuv/gojsonreference v0.0.0-20180127040603-bd5ef7bd5415/go.mod h1:GwrjFmJcFw6At/Gs6z4yjiIwzuJ1/+UwLxMQDVQXShQ= +github.com/xeipuuv/gojsonschema v1.2.0 h1:LhYJRs+L4fBtjZUfuSZIKGeVu0QRy8e5Xi7D17UxZ74= +github.com/xeipuuv/gojsonschema v1.2.0/go.mod h1:anYRn/JVcOK2ZgGU+IjEV4nwlhoK5sQluxsYJ78Id3Y= github.com/xrash/smetrics v0.0.0-20240312152122-5f08fbb34913 h1:+qGGcbkzsfDQNPPe9UDgpxAWQrhbbBXOYJFQDq/dtJw= github.com/xrash/smetrics v0.0.0-20240312152122-5f08fbb34913/go.mod h1:4aEEwZQutDLsQv2Deui4iYQ6DWTxR14g6m8Wv88+Xqk= +github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= +github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= github.com/zeebo/assert v1.3.0 h1:g7C04CbJuIDKNPFHmsk4hwZDO5O+kntRxzaUoNXj+IQ= github.com/zeebo/assert v1.3.0/go.mod h1:Pq9JiuJQpG8JLJdtkwrJESF0Foym2/D9XMU5ciN/wJ0= @@ -276,6 +334,7 @@ go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= go.uber.org/zap/exp v0.2.0 h1:FtGenNNeCATRB3CmB/yEUnjEFeJWpB/pMcy7e2bKPYs= go.uber.org/zap/exp v0.2.0/go.mod h1:t0gqAIdh1MfKv9EwN/dLwfZnJxe9ITAZN78HEWPFWDQ= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= golang.org/x/crypto v0.8.0/go.mod h1:mRqEX+O9/h5TFCrQhkgjo2yKi0yYA+9ecGkdQoHrywE= @@ -288,6 +347,8 @@ golang.org/x/crypto v0.24.0 h1:mnl8DM0o513X8fdIkmyFE/5hTYxbwYOjDS/+rK6qpRI= golang.org/x/crypto v0.24.0/go.mod h1:Z1PMYSOR5nyMcyAVAIQSKCDwalqy85Aqn1x3Ws4L5DM= golang.org/x/exp v0.0.0-20240613232115-7f521ea00fb8 h1:yixxcjnhBmY0nkL253HFVIm0JsFHwrHdT3Yh6szTnfY= golang.org/x/exp v0.0.0-20240613232115-7f521ea00fb8/go.mod h1:jj3sYF3dwk5D+ghuXyeI3r5MFf+NT2An6/9dOA95KSI= +golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= +golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= @@ -300,7 +361,9 @@ golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLL golang.org/x/net v0.0.0-20190827160401-ba9fcec4b297/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20191007182048-72f939374954/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20200202094626-16171245cfb2/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20201010224723-4f7140c49acb/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= +golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= golang.org/x/net v0.0.0-20201110031124-69a78807bb2b/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= golang.org/x/net v0.0.0-20201216054612-986b41b23924/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= golang.org/x/net v0.0.0-20201224014010-6772e930b67b/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= @@ -321,6 +384,8 @@ golang.org/x/net v0.20.0/go.mod h1:z8BVo6PvndSri0LbOE3hAn0apkU+1YvI6E70E9jsnvY= golang.org/x/net v0.26.0 h1:soB7SVo0PWrY4vPW/+ay0jKDNScG2X9wFeYlXIvJsOQ= golang.org/x/net v0.26.0/go.mod h1:5YKkiSynbBIh3p6iOc/vibscux0x38BZDkn8sCUPxHE= golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= @@ -348,11 +413,13 @@ golang.org/x/sys v0.0.0-20210305230114-8fe3ee5dd75b/go.mod h1:h1NjWce9XRLGQEsW7w golang.org/x/sys v0.0.0-20210423082822-04245dca01da/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210525143221-35b2ab0089ea/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20210616094352-59db8d763f22/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210906170528-6f6e22806c34/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20210927094055-39ccf1dd6fa6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20211216021012-1d35b9e2eb4e/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220128215802-99c3d69c2c27/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-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.2.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= @@ -393,6 +460,8 @@ golang.org/x/text v0.16.0 h1:a94ExnEXNtEwYLGJSIUxnWoxoRz/ZcCsV63ROupILh4= golang.org/x/text v0.16.0/go.mod h1:GhwF1Be+LQoKShO3cGOHzqOgRrGaYc9AvblQOmPVHnI= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= +golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= +golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA= golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58= @@ -400,6 +469,7 @@ golang.org/x/tools v0.17.0/go.mod h1:xsh6VxdV005rRVaS6SSAf9oiAqljS7UZUacMZ8Bnsps golang.org/x/tools v0.22.0 h1:gqSGLZqv+AI9lIQzniJ0nZDRG5GBPsSi+DRNHWNz6yA= golang.org/x/tools v0.22.0/go.mod h1:aCwcsjqvq7Yqt6TNyX7QMU2enbQ/Gt0bo6krSeEri+c= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= +golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= google.golang.org/genproto/googleapis/rpc v0.0.0-20240521202816-d264139d666e h1:Elxv5MwEkCI9f5SkoL6afed6NTdxaGoAo39eANBwHL8= @@ -421,3 +491,5 @@ gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gotest.tools/v3 v3.3.0 h1:MfDY1b1/0xN1CyMlQDac0ziEy9zJQd9CXBRRDHw2jJo= +gotest.tools/v3 v3.3.0/go.mod h1:Mcr9QNxkg0uMvy/YElmo4SpXgJKWgQvYrT7Kw5RzJ1A= diff --git a/pkg/service/docker_test.go b/pkg/service/docker_test.go new file mode 100644 index 000000000..48937f3d4 --- /dev/null +++ b/pkg/service/docker_test.go @@ -0,0 +1,79 @@ +// Copyright 2024 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package service_test + +import ( + "fmt" + "log" + "net" + "os" + "sync/atomic" + "testing" + + "github.com/ory/dockertest/v3" +) + +var Docker *dockertest.Pool + +func TestMain(m *testing.M) { + pool, err := dockertest.NewPool("") + if err != nil { + log.Fatalf("Could not construct pool: %s", err) + } + + // uses pool to try to connect to Docker + err = pool.Client.Ping() + if err != nil { + log.Fatalf("Could not connect to Docker: %s", err) + } + Docker = pool + + code := m.Run() + os.Exit(code) +} + +func waitTCPPort(t testing.TB, addr string) { + if err := Docker.Retry(func() error { + conn, err := net.Dial("tcp", addr) + if err != nil { + t.Log(err) + return err + } + _ = conn.Close() + return nil + }); err != nil { + t.Fatal(err) + } +} + +var redisLast uint32 + +func runRedis(t testing.TB) string { + c, err := Docker.RunWithOptions(&dockertest.RunOptions{ + Name: fmt.Sprintf("lktest-redis-%d", atomic.AddUint32(&redisLast, 1)), + Repository: "redis", Tag: "latest", + }) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + _ = Docker.Purge(c) + }) + addr := c.GetHostPort("6379/tcp") + waitTCPPort(t, addr) + + t.Log("Redis running on", addr) + return addr +} diff --git a/pkg/service/interfaces.go b/pkg/service/interfaces.go index 73c7d4a67..9ebce8271 100644 --- a/pkg/service/interfaces.go +++ b/pkg/service/interfaces.go @@ -80,8 +80,14 @@ type RoomAllocator interface { //counterfeiter:generate . SIPStore type SIPStore interface { StoreSIPTrunk(ctx context.Context, info *livekit.SIPTrunkInfo) error + StoreSIPInboundTrunk(ctx context.Context, info *livekit.SIPInboundTrunkInfo) error + StoreSIPOutboundTrunk(ctx context.Context, info *livekit.SIPOutboundTrunkInfo) error LoadSIPTrunk(ctx context.Context, sipTrunkID string) (*livekit.SIPTrunkInfo, error) + LoadSIPInboundTrunk(ctx context.Context, sipTrunkID string) (*livekit.SIPInboundTrunkInfo, error) + LoadSIPOutboundTrunk(ctx context.Context, sipTrunkID string) (*livekit.SIPOutboundTrunkInfo, error) ListSIPTrunk(ctx context.Context) ([]*livekit.SIPTrunkInfo, error) + ListSIPInboundTrunk(ctx context.Context) ([]*livekit.SIPInboundTrunkInfo, error) + ListSIPOutboundTrunk(ctx context.Context) ([]*livekit.SIPOutboundTrunkInfo, error) DeleteSIPTrunk(ctx context.Context, info *livekit.SIPTrunkInfo) error StoreSIPDispatchRule(ctx context.Context, info *livekit.SIPDispatchRuleInfo) error diff --git a/pkg/service/ioservice_sip.go b/pkg/service/ioservice_sip.go index f3c686bfb..e974833ac 100644 --- a/pkg/service/ioservice_sip.go +++ b/pkg/service/ioservice_sip.go @@ -26,8 +26,8 @@ import ( // matchSIPTrunk finds a SIP Trunk definition matching the request. // Returns nil if no rules matched or an error if there are conflicting definitions. -func (s *IOInfoService) matchSIPTrunk(ctx context.Context, calling, called string) (*livekit.SIPTrunkInfo, error) { - trunks, err := s.ss.ListSIPTrunk(ctx) +func (s *IOInfoService) matchSIPTrunk(ctx context.Context, calling, called string) (*livekit.SIPInboundTrunkInfo, error) { + trunks, err := s.ss.ListSIPInboundTrunk(ctx) if err != nil { return nil, err } @@ -36,7 +36,7 @@ func (s *IOInfoService) matchSIPTrunk(ctx context.Context, calling, called strin // matchSIPDispatchRule finds the best dispatch rule matching the request parameters. Returns an error if no rule matched. // Trunk parameter can be nil, in which case only wildcard dispatch rules will be effective (ones without Trunk IDs). -func (s *IOInfoService) matchSIPDispatchRule(ctx context.Context, trunk *livekit.SIPTrunkInfo, req *rpc.EvaluateSIPDispatchRulesRequest) (*livekit.SIPDispatchRuleInfo, error) { +func (s *IOInfoService) matchSIPDispatchRule(ctx context.Context, trunk *livekit.SIPInboundTrunkInfo, req *rpc.EvaluateSIPDispatchRulesRequest) (*livekit.SIPDispatchRuleInfo, error) { // Trunk can still be nil here in case none matched or were defined. // This is still fine, but only in case we'll match exactly one wildcard dispatch rule. rules, err := s.ss.ListSIPDispatchRule(ctx) @@ -96,7 +96,7 @@ func (s *IOInfoService) GetSIPTrunkAuthentication(ctx context.Context, req *rpc. log.Debugw("SIP trunk matched for auth", "sipTrunk", trunk.SipTrunkId) return &rpc.GetSIPTrunkAuthenticationResponse{ SipTrunkId: trunk.SipTrunkId, - Username: trunk.InboundUsername, - Password: trunk.InboundPassword, + Username: trunk.AuthUsername, + Password: trunk.AuthPassword, }, nil } diff --git a/pkg/service/redisstore.go b/pkg/service/redisstore.go index 52dfcaa62..4da0ab360 100644 --- a/pkg/service/redisstore.go +++ b/pkg/service/redisstore.go @@ -52,9 +52,6 @@ const ( IngressStatePrefix = "{ingress}_state:" RoomIngressPrefix = "room_{ingress}:" - SIPTrunkKey = "sip_trunk" - SIPDispatchRuleKey = "sip_dispatch_rule" - // RoomParticipantsPrefix is hash of participant_name => ParticipantInfo RoomParticipantsPrefix = "room_participants:" @@ -825,94 +822,54 @@ func (s *RedisStore) DeleteIngress(_ context.Context, info *livekit.IngressInfo) return nil } -func (s *RedisStore) loadOne(ctx context.Context, key, id string, info proto.Message, notFoundErr error) error { +func redisStoreOne(ctx context.Context, s *RedisStore, key, id string, p proto.Message) error { + if id == "" { + return errors.New("id is not set") + } + data, err := proto.Marshal(p) + if err != nil { + return err + } + return s.rc.HSet(s.ctx, key, id, data).Err() +} + +func redisLoadOne[T any, P interface { + *T + proto.Message +}](ctx context.Context, s *RedisStore, key, id string, notFoundErr error) (P, error) { data, err := s.rc.HGet(s.ctx, key, id).Result() - switch err { - case nil: - return proto.Unmarshal([]byte(data), info) - case redis.Nil: - return notFoundErr - default: - return err + if err == redis.Nil { + return nil, notFoundErr + } else if err != nil { + return nil, err } + var p P = new(T) + err = proto.Unmarshal([]byte(data), p) + if err != nil { + return nil, err + } + return p, err } -func (s *RedisStore) loadMany(ctx context.Context, key string, onResult func() proto.Message) error { +func redisLoadMany[T any, P interface { + *T + proto.Message +}](ctx context.Context, s *RedisStore, key string) ([]P, error) { data, err := s.rc.HGetAll(s.ctx, key).Result() - if err != nil { - if err == redis.Nil { - return nil - } - return err + if err == redis.Nil { + return nil, nil + } else if err != nil { + return nil, err } + list := make([]P, 0, len(data)) for _, d := range data { - if err = proto.Unmarshal([]byte(d), onResult()); err != nil { - return err + var p P = new(T) + if err = proto.Unmarshal([]byte(d), p); err != nil { + return list, err } + list = append(list, p) } - return nil -} - -func (s *RedisStore) StoreSIPTrunk(ctx context.Context, info *livekit.SIPTrunkInfo) error { - data, err := proto.Marshal(info) - if err != nil { - return err - } - - return s.rc.HSet(s.ctx, SIPTrunkKey, info.SipTrunkId, data).Err() -} - -func (s *RedisStore) LoadSIPTrunk(ctx context.Context, sipTrunkId string) (*livekit.SIPTrunkInfo, error) { - info := &livekit.SIPTrunkInfo{} - if err := s.loadOne(ctx, SIPTrunkKey, sipTrunkId, info, ErrSIPTrunkNotFound); err != nil { - return nil, err - } - - return info, nil -} - -func (s *RedisStore) DeleteSIPTrunk(ctx context.Context, info *livekit.SIPTrunkInfo) error { - return s.rc.HDel(s.ctx, SIPTrunkKey, info.SipTrunkId).Err() -} - -func (s *RedisStore) ListSIPTrunk(ctx context.Context) (infos []*livekit.SIPTrunkInfo, err error) { - err = s.loadMany(ctx, SIPTrunkKey, func() proto.Message { - infos = append(infos, &livekit.SIPTrunkInfo{}) - return infos[len(infos)-1] - }) - - return infos, err -} - -func (s *RedisStore) StoreSIPDispatchRule(ctx context.Context, info *livekit.SIPDispatchRuleInfo) error { - data, err := proto.Marshal(info) - if err != nil { - return err - } - - return s.rc.HSet(s.ctx, SIPDispatchRuleKey, info.SipDispatchRuleId, data).Err() -} - -func (s *RedisStore) LoadSIPDispatchRule(ctx context.Context, sipDispatchRuleId string) (*livekit.SIPDispatchRuleInfo, error) { - info := &livekit.SIPDispatchRuleInfo{} - if err := s.loadOne(ctx, SIPDispatchRuleKey, sipDispatchRuleId, info, ErrSIPDispatchRuleNotFound); err != nil { - return nil, err - } - - return info, nil -} - -func (s *RedisStore) DeleteSIPDispatchRule(ctx context.Context, info *livekit.SIPDispatchRuleInfo) error { - return s.rc.HDel(s.ctx, SIPDispatchRuleKey, info.SipDispatchRuleId).Err() -} - -func (s *RedisStore) ListSIPDispatchRule(ctx context.Context) (infos []*livekit.SIPDispatchRuleInfo, err error) { - err = s.loadMany(ctx, SIPDispatchRuleKey, func() proto.Message { - infos = append(infos, &livekit.SIPDispatchRuleInfo{}) - return infos[len(infos)-1] - }) - - return infos, err + return list, nil } diff --git a/pkg/service/redisstore_sip.go b/pkg/service/redisstore_sip.go new file mode 100644 index 000000000..d4cfbb09e --- /dev/null +++ b/pkg/service/redisstore_sip.go @@ -0,0 +1,199 @@ +// Copyright 2023 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package service + +import ( + "context" + + "github.com/livekit/protocol/livekit" +) + +const ( + SIPTrunkKey = "sip_trunk" + SIPInboundTrunkKey = "sip_inbound_trunk" + SIPOutboundTrunkKey = "sip_outbound_trunk" + SIPDispatchRuleKey = "sip_dispatch_rule" +) + +func (s *RedisStore) StoreSIPTrunk(ctx context.Context, info *livekit.SIPTrunkInfo) error { + return redisStoreOne(s.ctx, s, SIPTrunkKey, info.SipTrunkId, info) +} + +func (s *RedisStore) StoreSIPInboundTrunk(ctx context.Context, info *livekit.SIPInboundTrunkInfo) error { + return redisStoreOne(s.ctx, s, SIPInboundTrunkKey, info.SipTrunkId, info) +} + +func (s *RedisStore) StoreSIPOutboundTrunk(ctx context.Context, info *livekit.SIPOutboundTrunkInfo) error { + return redisStoreOne(s.ctx, s, SIPOutboundTrunkKey, info.SipTrunkId, info) +} + +func (s *RedisStore) loadSIPLegacyTrunk(ctx context.Context, id string) (*livekit.SIPTrunkInfo, error) { + return redisLoadOne[livekit.SIPTrunkInfo](ctx, s, SIPTrunkKey, id, ErrSIPTrunkNotFound) +} + +func (s *RedisStore) loadSIPInboundTrunk(ctx context.Context, id string) (*livekit.SIPInboundTrunkInfo, error) { + return redisLoadOne[livekit.SIPInboundTrunkInfo](ctx, s, SIPInboundTrunkKey, id, ErrSIPTrunkNotFound) +} + +func (s *RedisStore) loadSIPOutboundTrunk(ctx context.Context, id string) (*livekit.SIPOutboundTrunkInfo, error) { + return redisLoadOne[livekit.SIPOutboundTrunkInfo](ctx, s, SIPOutboundTrunkKey, id, ErrSIPTrunkNotFound) +} + +func (s *RedisStore) LoadSIPTrunk(ctx context.Context, id string) (*livekit.SIPTrunkInfo, error) { + tr, err := s.loadSIPLegacyTrunk(ctx, id) + if err == nil { + return tr, nil + } else if err != ErrSIPTrunkNotFound { + return nil, err + } + in, err := s.loadSIPInboundTrunk(ctx, id) + if err == nil { + return in.AsTrunkInfo(), nil + } else if err != ErrSIPTrunkNotFound { + return nil, err + } + out, err := s.loadSIPOutboundTrunk(ctx, id) + if err == nil { + return out.AsTrunkInfo(), nil + } else if err != ErrSIPTrunkNotFound { + return nil, err + } + return nil, ErrSIPTrunkNotFound +} + +func (s *RedisStore) LoadSIPInboundTrunk(ctx context.Context, id string) (*livekit.SIPInboundTrunkInfo, error) { + in, err := s.loadSIPInboundTrunk(ctx, id) + if err == nil { + return in, nil + } else if err != ErrSIPTrunkNotFound { + return nil, err + } + tr, err := s.loadSIPLegacyTrunk(ctx, id) + if err == nil { + return tr.AsInbound(), nil + } else if err != ErrSIPTrunkNotFound { + return nil, err + } + return nil, ErrSIPTrunkNotFound +} + +func (s *RedisStore) LoadSIPOutboundTrunk(ctx context.Context, id string) (*livekit.SIPOutboundTrunkInfo, error) { + in, err := s.loadSIPOutboundTrunk(ctx, id) + if err == nil { + return in, nil + } else if err != ErrSIPTrunkNotFound { + return nil, err + } + tr, err := s.loadSIPLegacyTrunk(ctx, id) + if err == nil { + return tr.AsOutbound(), nil + } else if err != ErrSIPTrunkNotFound { + return nil, err + } + return nil, ErrSIPTrunkNotFound +} + +func (s *RedisStore) deleteSIPTrunk(ctx context.Context, id string) error { + tx := s.rc.TxPipeline() + tx.HDel(s.ctx, SIPTrunkKey, id) + tx.HDel(s.ctx, SIPInboundTrunkKey, id) + tx.HDel(s.ctx, SIPOutboundTrunkKey, id) + _, err := tx.Exec(ctx) + return err +} + +func (s *RedisStore) DeleteSIPTrunk(ctx context.Context, info *livekit.SIPTrunkInfo) error { + return s.deleteSIPTrunk(ctx, info.SipTrunkId) +} + +func (s *RedisStore) listSIPLegacyTrunk(ctx context.Context) ([]*livekit.SIPTrunkInfo, error) { + return redisLoadMany[livekit.SIPTrunkInfo](ctx, s, SIPTrunkKey) +} + +func (s *RedisStore) listSIPInboundTrunk(ctx context.Context) ([]*livekit.SIPInboundTrunkInfo, error) { + return redisLoadMany[livekit.SIPInboundTrunkInfo](ctx, s, SIPInboundTrunkKey) +} + +func (s *RedisStore) listSIPOutboundTrunk(ctx context.Context) ([]*livekit.SIPOutboundTrunkInfo, error) { + return redisLoadMany[livekit.SIPOutboundTrunkInfo](ctx, s, SIPOutboundTrunkKey) +} + +func (s *RedisStore) ListSIPTrunk(ctx context.Context) ([]*livekit.SIPTrunkInfo, error) { + infos, err := s.listSIPLegacyTrunk(ctx) + if err != nil { + return nil, err + } + in, err := s.listSIPInboundTrunk(ctx) + if err != nil { + return infos, err + } + for _, t := range in { + infos = append(infos, t.AsTrunkInfo()) + } + out, err := s.listSIPOutboundTrunk(ctx) + if err != nil { + return infos, err + } + for _, t := range out { + infos = append(infos, t.AsTrunkInfo()) + } + return infos, nil +} + +func (s *RedisStore) ListSIPInboundTrunk(ctx context.Context) (infos []*livekit.SIPInboundTrunkInfo, err error) { + in, err := s.listSIPInboundTrunk(ctx) + if err != nil { + return in, err + } + old, err := s.listSIPLegacyTrunk(ctx) + if err != nil { + return nil, err + } + for _, t := range old { + in = append(in, t.AsInbound()) + } + return in, nil +} + +func (s *RedisStore) ListSIPOutboundTrunk(ctx context.Context) (infos []*livekit.SIPOutboundTrunkInfo, err error) { + out, err := s.listSIPOutboundTrunk(ctx) + if err != nil { + return out, err + } + old, err := s.listSIPLegacyTrunk(ctx) + if err != nil { + return nil, err + } + for _, t := range old { + out = append(out, t.AsOutbound()) + } + return out, nil +} + +func (s *RedisStore) StoreSIPDispatchRule(ctx context.Context, info *livekit.SIPDispatchRuleInfo) error { + return redisStoreOne(ctx, s, SIPDispatchRuleKey, info.SipDispatchRuleId, info) +} + +func (s *RedisStore) LoadSIPDispatchRule(ctx context.Context, sipDispatchRuleId string) (*livekit.SIPDispatchRuleInfo, error) { + return redisLoadOne[livekit.SIPDispatchRuleInfo](ctx, s, SIPDispatchRuleKey, sipDispatchRuleId, ErrSIPDispatchRuleNotFound) +} + +func (s *RedisStore) DeleteSIPDispatchRule(ctx context.Context, info *livekit.SIPDispatchRuleInfo) error { + return s.rc.HDel(s.ctx, SIPDispatchRuleKey, info.SipDispatchRuleId).Err() +} + +func (s *RedisStore) ListSIPDispatchRule(ctx context.Context) (infos []*livekit.SIPDispatchRuleInfo, err error) { + return redisLoadMany[livekit.SIPDispatchRuleInfo](ctx, s, SIPDispatchRuleKey) +} diff --git a/pkg/service/redisstore_sip_test.go b/pkg/service/redisstore_sip_test.go new file mode 100644 index 000000000..7e8e569f7 --- /dev/null +++ b/pkg/service/redisstore_sip_test.go @@ -0,0 +1,263 @@ +// Copyright 2024 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package service_test + +import ( + "context" + "slices" + "strings" + "testing" + + "github.com/livekit/protocol/livekit" + "github.com/livekit/protocol/utils" + "github.com/livekit/protocol/utils/guid" + "github.com/stretchr/testify/require" + "google.golang.org/protobuf/proto" + + "github.com/livekit/livekit-server/pkg/service" +) + +func TestSIPStoreDispatch(t *testing.T) { + ctx := context.Background() + rs := redisStore(t) + + id := guid.New(utils.SIPDispatchRulePrefix) + + // No dispatch rules initially. + list, err := rs.ListSIPDispatchRule(ctx) + require.NoError(t, err) + require.Empty(t, list) + + // Loading non-existent dispatch should return proper not found error. + got, err := rs.LoadSIPDispatchRule(ctx, id) + require.Equal(t, service.ErrSIPDispatchRuleNotFound, err) + require.Nil(t, got) + + // Creation without ID should fail. + rule := &livekit.SIPDispatchRuleInfo{ + TrunkIds: []string{"trunk"}, + Rule: &livekit.SIPDispatchRule{Rule: &livekit.SIPDispatchRule_DispatchRuleDirect{ + DispatchRuleDirect: &livekit.SIPDispatchRuleDirect{ + RoomName: "room", + Pin: "1234", + }, + }}, + } + err = rs.StoreSIPDispatchRule(ctx, rule) + require.Error(t, err) + + // Creation + rule.SipDispatchRuleId = id + err = rs.StoreSIPDispatchRule(ctx, rule) + require.NoError(t, err) + + // Loading + got, err = rs.LoadSIPDispatchRule(ctx, id) + require.NoError(t, err) + require.True(t, proto.Equal(rule, got)) + + // Listing + list, err = rs.ListSIPDispatchRule(ctx) + require.NoError(t, err) + require.Len(t, list, 1) + require.True(t, proto.Equal(rule, list[0])) + + // Deletion. Should not return error if not exists. + err = rs.DeleteSIPDispatchRule(ctx, &livekit.SIPDispatchRuleInfo{SipDispatchRuleId: id}) + require.NoError(t, err) + err = rs.DeleteSIPDispatchRule(ctx, &livekit.SIPDispatchRuleInfo{SipDispatchRuleId: id}) + require.NoError(t, err) + + // Check that it's deleted. + list, err = rs.ListSIPDispatchRule(ctx) + require.NoError(t, err) + require.Empty(t, list) + + got, err = rs.LoadSIPDispatchRule(ctx, id) + require.Equal(t, service.ErrSIPDispatchRuleNotFound, err) + require.Nil(t, got) +} + +func TestSIPStoreTrunk(t *testing.T) { + ctx := context.Background() + rs := redisStore(t) + + oldID := guid.New(utils.SIPTrunkPrefix) + inID := guid.New(utils.SIPTrunkPrefix) + outID := guid.New(utils.SIPTrunkPrefix) + + // No trunks initially. Check legacy, inbound, outbound. + // Loading non-existent trunk should return proper not found error. + oldList, err := rs.ListSIPTrunk(ctx) + require.NoError(t, err) + require.Empty(t, oldList) + + old, err := rs.LoadSIPTrunk(ctx, oldID) + require.Equal(t, service.ErrSIPTrunkNotFound, err) + require.Nil(t, old) + + inList, err := rs.ListSIPInboundTrunk(ctx) + require.NoError(t, err) + require.Empty(t, inList) + + in, err := rs.LoadSIPInboundTrunk(ctx, oldID) + require.Equal(t, service.ErrSIPTrunkNotFound, err) + require.Nil(t, in) + + outList, err := rs.ListSIPOutboundTrunk(ctx) + require.NoError(t, err) + require.Empty(t, outList) + + out, err := rs.LoadSIPOutboundTrunk(ctx, oldID) + require.Equal(t, service.ErrSIPTrunkNotFound, err) + require.Nil(t, out) + + // Creation without ID should fail. + oldT := &livekit.SIPTrunkInfo{ + Name: "Legacy", + } + err = rs.StoreSIPTrunk(ctx, oldT) + require.Error(t, err) + + inT := &livekit.SIPInboundTrunkInfo{ + Name: "Inbound", + } + err = rs.StoreSIPInboundTrunk(ctx, inT) + require.Error(t, err) + + outT := &livekit.SIPOutboundTrunkInfo{ + Name: "Outbound", + } + err = rs.StoreSIPOutboundTrunk(ctx, outT) + require.Error(t, err) + + // Creation + oldT.SipTrunkId = oldID + err = rs.StoreSIPTrunk(ctx, oldT) + require.NoError(t, err) + + inT.SipTrunkId = inID + err = rs.StoreSIPInboundTrunk(ctx, inT) + require.NoError(t, err) + + outT.SipTrunkId = outID + err = rs.StoreSIPOutboundTrunk(ctx, outT) + require.NoError(t, err) + + // Loading (with matching kind) + oldT2, err := rs.LoadSIPTrunk(ctx, oldID) + require.NoError(t, err) + require.True(t, proto.Equal(oldT, oldT2)) + + inT2, err := rs.LoadSIPInboundTrunk(ctx, inID) + require.NoError(t, err) + require.True(t, proto.Equal(inT, inT2)) + + outT2, err := rs.LoadSIPOutboundTrunk(ctx, outID) + require.NoError(t, err) + require.True(t, proto.Equal(outT, outT2)) + + // Loading (compat) + oldT2, err = rs.LoadSIPTrunk(ctx, inID) + require.NoError(t, err) + require.True(t, proto.Equal(inT.AsTrunkInfo(), oldT2)) + + oldT2, err = rs.LoadSIPTrunk(ctx, outID) + require.NoError(t, err) + require.True(t, proto.Equal(outT.AsTrunkInfo(), oldT2)) + + inT2, err = rs.LoadSIPInboundTrunk(ctx, oldID) + require.NoError(t, err) + require.True(t, proto.Equal(oldT.AsInbound(), inT2)) + + outT2, err = rs.LoadSIPOutboundTrunk(ctx, oldID) + require.NoError(t, err) + require.True(t, proto.Equal(oldT.AsOutbound(), outT2)) + + // Listing (always shows legacy + new) + listOld, err := rs.ListSIPTrunk(ctx) + require.NoError(t, err) + require.Len(t, listOld, 3) + slices.SortFunc(listOld, func(a, b *livekit.SIPTrunkInfo) int { + return strings.Compare(a.Name, b.Name) + }) + require.True(t, proto.Equal(inT.AsTrunkInfo(), listOld[0])) + require.True(t, proto.Equal(oldT, listOld[1])) + require.True(t, proto.Equal(outT.AsTrunkInfo(), listOld[2])) + + listIn, err := rs.ListSIPInboundTrunk(ctx) + require.NoError(t, err) + require.Len(t, listIn, 2) + slices.SortFunc(listIn, func(a, b *livekit.SIPInboundTrunkInfo) int { + return strings.Compare(a.Name, b.Name) + }) + require.True(t, proto.Equal(inT, listIn[0])) + require.True(t, proto.Equal(oldT.AsInbound(), listIn[1])) + + listOut, err := rs.ListSIPOutboundTrunk(ctx) + require.NoError(t, err) + require.Len(t, listOut, 2) + slices.SortFunc(listOut, func(a, b *livekit.SIPOutboundTrunkInfo) int { + return strings.Compare(a.Name, b.Name) + }) + require.True(t, proto.Equal(oldT.AsOutbound(), listOut[0])) + require.True(t, proto.Equal(outT, listOut[1])) + + // Deletion. Should not return error if not exists. + err = rs.DeleteSIPTrunk(ctx, &livekit.SIPTrunkInfo{SipTrunkId: oldID}) + require.NoError(t, err) + err = rs.DeleteSIPTrunk(ctx, &livekit.SIPTrunkInfo{SipTrunkId: oldID}) + require.NoError(t, err) + + // Other objects are still there. + inT2, err = rs.LoadSIPInboundTrunk(ctx, inID) + require.NoError(t, err) + require.True(t, proto.Equal(inT, inT2)) + + outT2, err = rs.LoadSIPOutboundTrunk(ctx, outID) + require.NoError(t, err) + require.True(t, proto.Equal(outT, outT2)) + + // Delete the rest + err = rs.DeleteSIPTrunk(ctx, &livekit.SIPTrunkInfo{SipTrunkId: inID}) + require.NoError(t, err) + err = rs.DeleteSIPTrunk(ctx, &livekit.SIPTrunkInfo{SipTrunkId: outID}) + require.NoError(t, err) + + // Check everything is deleted. + oldList, err = rs.ListSIPTrunk(ctx) + require.NoError(t, err) + require.Empty(t, oldList) + + inList, err = rs.ListSIPInboundTrunk(ctx) + require.NoError(t, err) + require.Empty(t, inList) + + outList, err = rs.ListSIPOutboundTrunk(ctx) + require.NoError(t, err) + require.Empty(t, outList) + + old, err = rs.LoadSIPTrunk(ctx, oldID) + require.Equal(t, service.ErrSIPTrunkNotFound, err) + require.Nil(t, old) + + in, err = rs.LoadSIPInboundTrunk(ctx, oldID) + require.Equal(t, service.ErrSIPTrunkNotFound, err) + require.Nil(t, in) + + out, err = rs.LoadSIPOutboundTrunk(ctx, oldID) + require.Equal(t, service.ErrSIPTrunkNotFound, err) + require.Nil(t, out) +} diff --git a/pkg/service/redisstore_test.go b/pkg/service/redisstore_test.go index 5f62a71f1..d62db1025 100644 --- a/pkg/service/redisstore_test.go +++ b/pkg/service/redisstore_test.go @@ -31,9 +31,13 @@ import ( "github.com/livekit/livekit-server/pkg/service" ) +func redisStore(t testing.TB) *service.RedisStore { + return service.NewRedisStore(redisClient(t)) +} + func TestRoomInternal(t *testing.T) { ctx := context.Background() - rs := service.NewRedisStore(redisClient()) + rs := redisStore(t) room := &livekit.Room{ Sid: "123", @@ -61,7 +65,7 @@ func TestRoomInternal(t *testing.T) { func TestParticipantPersistence(t *testing.T) { ctx := context.Background() - rs := service.NewRedisStore(redisClient()) + rs := redisStore(t) roomName := livekit.RoomName("room1") _ = rs.DeleteRoom(ctx, roomName) @@ -108,7 +112,7 @@ func TestParticipantPersistence(t *testing.T) { func TestRoomLock(t *testing.T) { ctx := context.Background() - rs := service.NewRedisStore(redisClient()) + rs := redisStore(t) lockInterval := 5 * time.Millisecond roomName := livekit.RoomName("myroom") @@ -158,8 +162,7 @@ func TestRoomLock(t *testing.T) { func TestEgressStore(t *testing.T) { ctx := context.Background() - rc := redisClient() - rs := service.NewRedisStore(rc) + rs := redisStore(t) roomName := "egress-test" @@ -229,7 +232,7 @@ func TestEgressStore(t *testing.T) { func TestIngressStore(t *testing.T) { ctx := context.Background() - rs := service.NewRedisStore(redisClient()) + rs := redisStore(t) info := &livekit.IngressInfo{ IngressId: "ingressId", diff --git a/pkg/service/servicefakes/fake_sipstore.go b/pkg/service/servicefakes/fake_sipstore.go index fe88a5359..d5ccde202 100644 --- a/pkg/service/servicefakes/fake_sipstore.go +++ b/pkg/service/servicefakes/fake_sipstore.go @@ -47,6 +47,32 @@ type FakeSIPStore struct { result1 []*livekit.SIPDispatchRuleInfo result2 error } + ListSIPInboundTrunkStub func(context.Context) ([]*livekit.SIPInboundTrunkInfo, error) + listSIPInboundTrunkMutex sync.RWMutex + listSIPInboundTrunkArgsForCall []struct { + arg1 context.Context + } + listSIPInboundTrunkReturns struct { + result1 []*livekit.SIPInboundTrunkInfo + result2 error + } + listSIPInboundTrunkReturnsOnCall map[int]struct { + result1 []*livekit.SIPInboundTrunkInfo + result2 error + } + ListSIPOutboundTrunkStub func(context.Context) ([]*livekit.SIPOutboundTrunkInfo, error) + listSIPOutboundTrunkMutex sync.RWMutex + listSIPOutboundTrunkArgsForCall []struct { + arg1 context.Context + } + listSIPOutboundTrunkReturns struct { + result1 []*livekit.SIPOutboundTrunkInfo + result2 error + } + listSIPOutboundTrunkReturnsOnCall map[int]struct { + result1 []*livekit.SIPOutboundTrunkInfo + result2 error + } ListSIPTrunkStub func(context.Context) ([]*livekit.SIPTrunkInfo, error) listSIPTrunkMutex sync.RWMutex listSIPTrunkArgsForCall []struct { @@ -74,6 +100,34 @@ type FakeSIPStore struct { result1 *livekit.SIPDispatchRuleInfo result2 error } + LoadSIPInboundTrunkStub func(context.Context, string) (*livekit.SIPInboundTrunkInfo, error) + loadSIPInboundTrunkMutex sync.RWMutex + loadSIPInboundTrunkArgsForCall []struct { + arg1 context.Context + arg2 string + } + loadSIPInboundTrunkReturns struct { + result1 *livekit.SIPInboundTrunkInfo + result2 error + } + loadSIPInboundTrunkReturnsOnCall map[int]struct { + result1 *livekit.SIPInboundTrunkInfo + result2 error + } + LoadSIPOutboundTrunkStub func(context.Context, string) (*livekit.SIPOutboundTrunkInfo, error) + loadSIPOutboundTrunkMutex sync.RWMutex + loadSIPOutboundTrunkArgsForCall []struct { + arg1 context.Context + arg2 string + } + loadSIPOutboundTrunkReturns struct { + result1 *livekit.SIPOutboundTrunkInfo + result2 error + } + loadSIPOutboundTrunkReturnsOnCall map[int]struct { + result1 *livekit.SIPOutboundTrunkInfo + result2 error + } LoadSIPTrunkStub func(context.Context, string) (*livekit.SIPTrunkInfo, error) loadSIPTrunkMutex sync.RWMutex loadSIPTrunkArgsForCall []struct { @@ -100,6 +154,30 @@ type FakeSIPStore struct { storeSIPDispatchRuleReturnsOnCall map[int]struct { result1 error } + StoreSIPInboundTrunkStub func(context.Context, *livekit.SIPInboundTrunkInfo) error + storeSIPInboundTrunkMutex sync.RWMutex + storeSIPInboundTrunkArgsForCall []struct { + arg1 context.Context + arg2 *livekit.SIPInboundTrunkInfo + } + storeSIPInboundTrunkReturns struct { + result1 error + } + storeSIPInboundTrunkReturnsOnCall map[int]struct { + result1 error + } + StoreSIPOutboundTrunkStub func(context.Context, *livekit.SIPOutboundTrunkInfo) error + storeSIPOutboundTrunkMutex sync.RWMutex + storeSIPOutboundTrunkArgsForCall []struct { + arg1 context.Context + arg2 *livekit.SIPOutboundTrunkInfo + } + storeSIPOutboundTrunkReturns struct { + result1 error + } + storeSIPOutboundTrunkReturnsOnCall map[int]struct { + result1 error + } StoreSIPTrunkStub func(context.Context, *livekit.SIPTrunkInfo) error storeSIPTrunkMutex sync.RWMutex storeSIPTrunkArgsForCall []struct { @@ -304,6 +382,134 @@ func (fake *FakeSIPStore) ListSIPDispatchRuleReturnsOnCall(i int, result1 []*liv }{result1, result2} } +func (fake *FakeSIPStore) ListSIPInboundTrunk(arg1 context.Context) ([]*livekit.SIPInboundTrunkInfo, error) { + fake.listSIPInboundTrunkMutex.Lock() + ret, specificReturn := fake.listSIPInboundTrunkReturnsOnCall[len(fake.listSIPInboundTrunkArgsForCall)] + fake.listSIPInboundTrunkArgsForCall = append(fake.listSIPInboundTrunkArgsForCall, struct { + arg1 context.Context + }{arg1}) + stub := fake.ListSIPInboundTrunkStub + fakeReturns := fake.listSIPInboundTrunkReturns + fake.recordInvocation("ListSIPInboundTrunk", []interface{}{arg1}) + fake.listSIPInboundTrunkMutex.Unlock() + if stub != nil { + return stub(arg1) + } + if specificReturn { + return ret.result1, ret.result2 + } + return fakeReturns.result1, fakeReturns.result2 +} + +func (fake *FakeSIPStore) ListSIPInboundTrunkCallCount() int { + fake.listSIPInboundTrunkMutex.RLock() + defer fake.listSIPInboundTrunkMutex.RUnlock() + return len(fake.listSIPInboundTrunkArgsForCall) +} + +func (fake *FakeSIPStore) ListSIPInboundTrunkCalls(stub func(context.Context) ([]*livekit.SIPInboundTrunkInfo, error)) { + fake.listSIPInboundTrunkMutex.Lock() + defer fake.listSIPInboundTrunkMutex.Unlock() + fake.ListSIPInboundTrunkStub = stub +} + +func (fake *FakeSIPStore) ListSIPInboundTrunkArgsForCall(i int) context.Context { + fake.listSIPInboundTrunkMutex.RLock() + defer fake.listSIPInboundTrunkMutex.RUnlock() + argsForCall := fake.listSIPInboundTrunkArgsForCall[i] + return argsForCall.arg1 +} + +func (fake *FakeSIPStore) ListSIPInboundTrunkReturns(result1 []*livekit.SIPInboundTrunkInfo, result2 error) { + fake.listSIPInboundTrunkMutex.Lock() + defer fake.listSIPInboundTrunkMutex.Unlock() + fake.ListSIPInboundTrunkStub = nil + fake.listSIPInboundTrunkReturns = struct { + result1 []*livekit.SIPInboundTrunkInfo + result2 error + }{result1, result2} +} + +func (fake *FakeSIPStore) ListSIPInboundTrunkReturnsOnCall(i int, result1 []*livekit.SIPInboundTrunkInfo, result2 error) { + fake.listSIPInboundTrunkMutex.Lock() + defer fake.listSIPInboundTrunkMutex.Unlock() + fake.ListSIPInboundTrunkStub = nil + if fake.listSIPInboundTrunkReturnsOnCall == nil { + fake.listSIPInboundTrunkReturnsOnCall = make(map[int]struct { + result1 []*livekit.SIPInboundTrunkInfo + result2 error + }) + } + fake.listSIPInboundTrunkReturnsOnCall[i] = struct { + result1 []*livekit.SIPInboundTrunkInfo + result2 error + }{result1, result2} +} + +func (fake *FakeSIPStore) ListSIPOutboundTrunk(arg1 context.Context) ([]*livekit.SIPOutboundTrunkInfo, error) { + fake.listSIPOutboundTrunkMutex.Lock() + ret, specificReturn := fake.listSIPOutboundTrunkReturnsOnCall[len(fake.listSIPOutboundTrunkArgsForCall)] + fake.listSIPOutboundTrunkArgsForCall = append(fake.listSIPOutboundTrunkArgsForCall, struct { + arg1 context.Context + }{arg1}) + stub := fake.ListSIPOutboundTrunkStub + fakeReturns := fake.listSIPOutboundTrunkReturns + fake.recordInvocation("ListSIPOutboundTrunk", []interface{}{arg1}) + fake.listSIPOutboundTrunkMutex.Unlock() + if stub != nil { + return stub(arg1) + } + if specificReturn { + return ret.result1, ret.result2 + } + return fakeReturns.result1, fakeReturns.result2 +} + +func (fake *FakeSIPStore) ListSIPOutboundTrunkCallCount() int { + fake.listSIPOutboundTrunkMutex.RLock() + defer fake.listSIPOutboundTrunkMutex.RUnlock() + return len(fake.listSIPOutboundTrunkArgsForCall) +} + +func (fake *FakeSIPStore) ListSIPOutboundTrunkCalls(stub func(context.Context) ([]*livekit.SIPOutboundTrunkInfo, error)) { + fake.listSIPOutboundTrunkMutex.Lock() + defer fake.listSIPOutboundTrunkMutex.Unlock() + fake.ListSIPOutboundTrunkStub = stub +} + +func (fake *FakeSIPStore) ListSIPOutboundTrunkArgsForCall(i int) context.Context { + fake.listSIPOutboundTrunkMutex.RLock() + defer fake.listSIPOutboundTrunkMutex.RUnlock() + argsForCall := fake.listSIPOutboundTrunkArgsForCall[i] + return argsForCall.arg1 +} + +func (fake *FakeSIPStore) ListSIPOutboundTrunkReturns(result1 []*livekit.SIPOutboundTrunkInfo, result2 error) { + fake.listSIPOutboundTrunkMutex.Lock() + defer fake.listSIPOutboundTrunkMutex.Unlock() + fake.ListSIPOutboundTrunkStub = nil + fake.listSIPOutboundTrunkReturns = struct { + result1 []*livekit.SIPOutboundTrunkInfo + result2 error + }{result1, result2} +} + +func (fake *FakeSIPStore) ListSIPOutboundTrunkReturnsOnCall(i int, result1 []*livekit.SIPOutboundTrunkInfo, result2 error) { + fake.listSIPOutboundTrunkMutex.Lock() + defer fake.listSIPOutboundTrunkMutex.Unlock() + fake.ListSIPOutboundTrunkStub = nil + if fake.listSIPOutboundTrunkReturnsOnCall == nil { + fake.listSIPOutboundTrunkReturnsOnCall = make(map[int]struct { + result1 []*livekit.SIPOutboundTrunkInfo + result2 error + }) + } + fake.listSIPOutboundTrunkReturnsOnCall[i] = struct { + result1 []*livekit.SIPOutboundTrunkInfo + result2 error + }{result1, result2} +} + func (fake *FakeSIPStore) ListSIPTrunk(arg1 context.Context) ([]*livekit.SIPTrunkInfo, error) { fake.listSIPTrunkMutex.Lock() ret, specificReturn := fake.listSIPTrunkReturnsOnCall[len(fake.listSIPTrunkArgsForCall)] @@ -433,6 +639,136 @@ func (fake *FakeSIPStore) LoadSIPDispatchRuleReturnsOnCall(i int, result1 *livek }{result1, result2} } +func (fake *FakeSIPStore) LoadSIPInboundTrunk(arg1 context.Context, arg2 string) (*livekit.SIPInboundTrunkInfo, error) { + fake.loadSIPInboundTrunkMutex.Lock() + ret, specificReturn := fake.loadSIPInboundTrunkReturnsOnCall[len(fake.loadSIPInboundTrunkArgsForCall)] + fake.loadSIPInboundTrunkArgsForCall = append(fake.loadSIPInboundTrunkArgsForCall, struct { + arg1 context.Context + arg2 string + }{arg1, arg2}) + stub := fake.LoadSIPInboundTrunkStub + fakeReturns := fake.loadSIPInboundTrunkReturns + fake.recordInvocation("LoadSIPInboundTrunk", []interface{}{arg1, arg2}) + fake.loadSIPInboundTrunkMutex.Unlock() + if stub != nil { + return stub(arg1, arg2) + } + if specificReturn { + return ret.result1, ret.result2 + } + return fakeReturns.result1, fakeReturns.result2 +} + +func (fake *FakeSIPStore) LoadSIPInboundTrunkCallCount() int { + fake.loadSIPInboundTrunkMutex.RLock() + defer fake.loadSIPInboundTrunkMutex.RUnlock() + return len(fake.loadSIPInboundTrunkArgsForCall) +} + +func (fake *FakeSIPStore) LoadSIPInboundTrunkCalls(stub func(context.Context, string) (*livekit.SIPInboundTrunkInfo, error)) { + fake.loadSIPInboundTrunkMutex.Lock() + defer fake.loadSIPInboundTrunkMutex.Unlock() + fake.LoadSIPInboundTrunkStub = stub +} + +func (fake *FakeSIPStore) LoadSIPInboundTrunkArgsForCall(i int) (context.Context, string) { + fake.loadSIPInboundTrunkMutex.RLock() + defer fake.loadSIPInboundTrunkMutex.RUnlock() + argsForCall := fake.loadSIPInboundTrunkArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2 +} + +func (fake *FakeSIPStore) LoadSIPInboundTrunkReturns(result1 *livekit.SIPInboundTrunkInfo, result2 error) { + fake.loadSIPInboundTrunkMutex.Lock() + defer fake.loadSIPInboundTrunkMutex.Unlock() + fake.LoadSIPInboundTrunkStub = nil + fake.loadSIPInboundTrunkReturns = struct { + result1 *livekit.SIPInboundTrunkInfo + result2 error + }{result1, result2} +} + +func (fake *FakeSIPStore) LoadSIPInboundTrunkReturnsOnCall(i int, result1 *livekit.SIPInboundTrunkInfo, result2 error) { + fake.loadSIPInboundTrunkMutex.Lock() + defer fake.loadSIPInboundTrunkMutex.Unlock() + fake.LoadSIPInboundTrunkStub = nil + if fake.loadSIPInboundTrunkReturnsOnCall == nil { + fake.loadSIPInboundTrunkReturnsOnCall = make(map[int]struct { + result1 *livekit.SIPInboundTrunkInfo + result2 error + }) + } + fake.loadSIPInboundTrunkReturnsOnCall[i] = struct { + result1 *livekit.SIPInboundTrunkInfo + result2 error + }{result1, result2} +} + +func (fake *FakeSIPStore) LoadSIPOutboundTrunk(arg1 context.Context, arg2 string) (*livekit.SIPOutboundTrunkInfo, error) { + fake.loadSIPOutboundTrunkMutex.Lock() + ret, specificReturn := fake.loadSIPOutboundTrunkReturnsOnCall[len(fake.loadSIPOutboundTrunkArgsForCall)] + fake.loadSIPOutboundTrunkArgsForCall = append(fake.loadSIPOutboundTrunkArgsForCall, struct { + arg1 context.Context + arg2 string + }{arg1, arg2}) + stub := fake.LoadSIPOutboundTrunkStub + fakeReturns := fake.loadSIPOutboundTrunkReturns + fake.recordInvocation("LoadSIPOutboundTrunk", []interface{}{arg1, arg2}) + fake.loadSIPOutboundTrunkMutex.Unlock() + if stub != nil { + return stub(arg1, arg2) + } + if specificReturn { + return ret.result1, ret.result2 + } + return fakeReturns.result1, fakeReturns.result2 +} + +func (fake *FakeSIPStore) LoadSIPOutboundTrunkCallCount() int { + fake.loadSIPOutboundTrunkMutex.RLock() + defer fake.loadSIPOutboundTrunkMutex.RUnlock() + return len(fake.loadSIPOutboundTrunkArgsForCall) +} + +func (fake *FakeSIPStore) LoadSIPOutboundTrunkCalls(stub func(context.Context, string) (*livekit.SIPOutboundTrunkInfo, error)) { + fake.loadSIPOutboundTrunkMutex.Lock() + defer fake.loadSIPOutboundTrunkMutex.Unlock() + fake.LoadSIPOutboundTrunkStub = stub +} + +func (fake *FakeSIPStore) LoadSIPOutboundTrunkArgsForCall(i int) (context.Context, string) { + fake.loadSIPOutboundTrunkMutex.RLock() + defer fake.loadSIPOutboundTrunkMutex.RUnlock() + argsForCall := fake.loadSIPOutboundTrunkArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2 +} + +func (fake *FakeSIPStore) LoadSIPOutboundTrunkReturns(result1 *livekit.SIPOutboundTrunkInfo, result2 error) { + fake.loadSIPOutboundTrunkMutex.Lock() + defer fake.loadSIPOutboundTrunkMutex.Unlock() + fake.LoadSIPOutboundTrunkStub = nil + fake.loadSIPOutboundTrunkReturns = struct { + result1 *livekit.SIPOutboundTrunkInfo + result2 error + }{result1, result2} +} + +func (fake *FakeSIPStore) LoadSIPOutboundTrunkReturnsOnCall(i int, result1 *livekit.SIPOutboundTrunkInfo, result2 error) { + fake.loadSIPOutboundTrunkMutex.Lock() + defer fake.loadSIPOutboundTrunkMutex.Unlock() + fake.LoadSIPOutboundTrunkStub = nil + if fake.loadSIPOutboundTrunkReturnsOnCall == nil { + fake.loadSIPOutboundTrunkReturnsOnCall = make(map[int]struct { + result1 *livekit.SIPOutboundTrunkInfo + result2 error + }) + } + fake.loadSIPOutboundTrunkReturnsOnCall[i] = struct { + result1 *livekit.SIPOutboundTrunkInfo + result2 error + }{result1, result2} +} + func (fake *FakeSIPStore) LoadSIPTrunk(arg1 context.Context, arg2 string) (*livekit.SIPTrunkInfo, error) { fake.loadSIPTrunkMutex.Lock() ret, specificReturn := fake.loadSIPTrunkReturnsOnCall[len(fake.loadSIPTrunkArgsForCall)] @@ -560,6 +896,130 @@ func (fake *FakeSIPStore) StoreSIPDispatchRuleReturnsOnCall(i int, result1 error }{result1} } +func (fake *FakeSIPStore) StoreSIPInboundTrunk(arg1 context.Context, arg2 *livekit.SIPInboundTrunkInfo) error { + fake.storeSIPInboundTrunkMutex.Lock() + ret, specificReturn := fake.storeSIPInboundTrunkReturnsOnCall[len(fake.storeSIPInboundTrunkArgsForCall)] + fake.storeSIPInboundTrunkArgsForCall = append(fake.storeSIPInboundTrunkArgsForCall, struct { + arg1 context.Context + arg2 *livekit.SIPInboundTrunkInfo + }{arg1, arg2}) + stub := fake.StoreSIPInboundTrunkStub + fakeReturns := fake.storeSIPInboundTrunkReturns + fake.recordInvocation("StoreSIPInboundTrunk", []interface{}{arg1, arg2}) + fake.storeSIPInboundTrunkMutex.Unlock() + if stub != nil { + return stub(arg1, arg2) + } + if specificReturn { + return ret.result1 + } + return fakeReturns.result1 +} + +func (fake *FakeSIPStore) StoreSIPInboundTrunkCallCount() int { + fake.storeSIPInboundTrunkMutex.RLock() + defer fake.storeSIPInboundTrunkMutex.RUnlock() + return len(fake.storeSIPInboundTrunkArgsForCall) +} + +func (fake *FakeSIPStore) StoreSIPInboundTrunkCalls(stub func(context.Context, *livekit.SIPInboundTrunkInfo) error) { + fake.storeSIPInboundTrunkMutex.Lock() + defer fake.storeSIPInboundTrunkMutex.Unlock() + fake.StoreSIPInboundTrunkStub = stub +} + +func (fake *FakeSIPStore) StoreSIPInboundTrunkArgsForCall(i int) (context.Context, *livekit.SIPInboundTrunkInfo) { + fake.storeSIPInboundTrunkMutex.RLock() + defer fake.storeSIPInboundTrunkMutex.RUnlock() + argsForCall := fake.storeSIPInboundTrunkArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2 +} + +func (fake *FakeSIPStore) StoreSIPInboundTrunkReturns(result1 error) { + fake.storeSIPInboundTrunkMutex.Lock() + defer fake.storeSIPInboundTrunkMutex.Unlock() + fake.StoreSIPInboundTrunkStub = nil + fake.storeSIPInboundTrunkReturns = struct { + result1 error + }{result1} +} + +func (fake *FakeSIPStore) StoreSIPInboundTrunkReturnsOnCall(i int, result1 error) { + fake.storeSIPInboundTrunkMutex.Lock() + defer fake.storeSIPInboundTrunkMutex.Unlock() + fake.StoreSIPInboundTrunkStub = nil + if fake.storeSIPInboundTrunkReturnsOnCall == nil { + fake.storeSIPInboundTrunkReturnsOnCall = make(map[int]struct { + result1 error + }) + } + fake.storeSIPInboundTrunkReturnsOnCall[i] = struct { + result1 error + }{result1} +} + +func (fake *FakeSIPStore) StoreSIPOutboundTrunk(arg1 context.Context, arg2 *livekit.SIPOutboundTrunkInfo) error { + fake.storeSIPOutboundTrunkMutex.Lock() + ret, specificReturn := fake.storeSIPOutboundTrunkReturnsOnCall[len(fake.storeSIPOutboundTrunkArgsForCall)] + fake.storeSIPOutboundTrunkArgsForCall = append(fake.storeSIPOutboundTrunkArgsForCall, struct { + arg1 context.Context + arg2 *livekit.SIPOutboundTrunkInfo + }{arg1, arg2}) + stub := fake.StoreSIPOutboundTrunkStub + fakeReturns := fake.storeSIPOutboundTrunkReturns + fake.recordInvocation("StoreSIPOutboundTrunk", []interface{}{arg1, arg2}) + fake.storeSIPOutboundTrunkMutex.Unlock() + if stub != nil { + return stub(arg1, arg2) + } + if specificReturn { + return ret.result1 + } + return fakeReturns.result1 +} + +func (fake *FakeSIPStore) StoreSIPOutboundTrunkCallCount() int { + fake.storeSIPOutboundTrunkMutex.RLock() + defer fake.storeSIPOutboundTrunkMutex.RUnlock() + return len(fake.storeSIPOutboundTrunkArgsForCall) +} + +func (fake *FakeSIPStore) StoreSIPOutboundTrunkCalls(stub func(context.Context, *livekit.SIPOutboundTrunkInfo) error) { + fake.storeSIPOutboundTrunkMutex.Lock() + defer fake.storeSIPOutboundTrunkMutex.Unlock() + fake.StoreSIPOutboundTrunkStub = stub +} + +func (fake *FakeSIPStore) StoreSIPOutboundTrunkArgsForCall(i int) (context.Context, *livekit.SIPOutboundTrunkInfo) { + fake.storeSIPOutboundTrunkMutex.RLock() + defer fake.storeSIPOutboundTrunkMutex.RUnlock() + argsForCall := fake.storeSIPOutboundTrunkArgsForCall[i] + return argsForCall.arg1, argsForCall.arg2 +} + +func (fake *FakeSIPStore) StoreSIPOutboundTrunkReturns(result1 error) { + fake.storeSIPOutboundTrunkMutex.Lock() + defer fake.storeSIPOutboundTrunkMutex.Unlock() + fake.StoreSIPOutboundTrunkStub = nil + fake.storeSIPOutboundTrunkReturns = struct { + result1 error + }{result1} +} + +func (fake *FakeSIPStore) StoreSIPOutboundTrunkReturnsOnCall(i int, result1 error) { + fake.storeSIPOutboundTrunkMutex.Lock() + defer fake.storeSIPOutboundTrunkMutex.Unlock() + fake.StoreSIPOutboundTrunkStub = nil + if fake.storeSIPOutboundTrunkReturnsOnCall == nil { + fake.storeSIPOutboundTrunkReturnsOnCall = make(map[int]struct { + result1 error + }) + } + fake.storeSIPOutboundTrunkReturnsOnCall[i] = struct { + result1 error + }{result1} +} + func (fake *FakeSIPStore) StoreSIPTrunk(arg1 context.Context, arg2 *livekit.SIPTrunkInfo) error { fake.storeSIPTrunkMutex.Lock() ret, specificReturn := fake.storeSIPTrunkReturnsOnCall[len(fake.storeSIPTrunkArgsForCall)] @@ -631,14 +1091,26 @@ func (fake *FakeSIPStore) Invocations() map[string][][]interface{} { defer fake.deleteSIPTrunkMutex.RUnlock() fake.listSIPDispatchRuleMutex.RLock() defer fake.listSIPDispatchRuleMutex.RUnlock() + fake.listSIPInboundTrunkMutex.RLock() + defer fake.listSIPInboundTrunkMutex.RUnlock() + fake.listSIPOutboundTrunkMutex.RLock() + defer fake.listSIPOutboundTrunkMutex.RUnlock() fake.listSIPTrunkMutex.RLock() defer fake.listSIPTrunkMutex.RUnlock() fake.loadSIPDispatchRuleMutex.RLock() defer fake.loadSIPDispatchRuleMutex.RUnlock() + fake.loadSIPInboundTrunkMutex.RLock() + defer fake.loadSIPInboundTrunkMutex.RUnlock() + fake.loadSIPOutboundTrunkMutex.RLock() + defer fake.loadSIPOutboundTrunkMutex.RUnlock() fake.loadSIPTrunkMutex.RLock() defer fake.loadSIPTrunkMutex.RUnlock() fake.storeSIPDispatchRuleMutex.RLock() defer fake.storeSIPDispatchRuleMutex.RUnlock() + fake.storeSIPInboundTrunkMutex.RLock() + defer fake.storeSIPInboundTrunkMutex.RUnlock() + fake.storeSIPOutboundTrunkMutex.RLock() + defer fake.storeSIPOutboundTrunkMutex.RUnlock() fake.storeSIPTrunkMutex.RLock() defer fake.storeSIPTrunkMutex.RUnlock() copiedInvocations := map[string][][]interface{}{} diff --git a/pkg/service/sip.go b/pkg/service/sip.go index 7421a5bd8..c70522560 100644 --- a/pkg/service/sip.go +++ b/pkg/service/sip.go @@ -16,6 +16,7 @@ package service import ( "context" + "errors" "fmt" "time" @@ -82,7 +83,38 @@ func (s *SIPService) CreateSIPTrunk(ctx context.Context, req *livekit.CreateSIPT } // Validate all trunks including the new one first. - list, err := s.store.ListSIPTrunk(ctx) + list, err := s.store.ListSIPInboundTrunk(ctx) + if err != nil { + return nil, err + } + list = append(list, info.AsInbound()) + if err = sip.ValidateTrunks(list); err != nil { + return nil, err + } + + // Now we can generate ID and store. + info.SipTrunkId = guid.New(utils.SIPTrunkPrefix) + if err := s.store.StoreSIPTrunk(ctx, info); err != nil { + return nil, err + } + return info, nil +} + +func (s *SIPService) CreateSIPInboundTrunk(ctx context.Context, req *livekit.CreateSIPInboundTrunkRequest) (*livekit.SIPInboundTrunkInfo, error) { + if s.store == nil { + return nil, ErrSIPNotConnected + } + info := req.Trunk + if info == nil { + return nil, errors.New("trunk info is required") + } else if info.SipTrunkId != "" { + return nil, errors.New("trunk ID must be empty") + } + + // Keep ID empty still, so that validation can print "" instead of a non-existent ID in the error. + + // Validate all trunks including the new one first. + list, err := s.store.ListSIPInboundTrunk(ctx) if err != nil { return nil, err } @@ -93,7 +125,26 @@ func (s *SIPService) CreateSIPTrunk(ctx context.Context, req *livekit.CreateSIPT // Now we can generate ID and store. info.SipTrunkId = guid.New(utils.SIPTrunkPrefix) - if err := s.store.StoreSIPTrunk(ctx, info); err != nil { + if err := s.store.StoreSIPInboundTrunk(ctx, info); err != nil { + return nil, err + } + return info, nil +} + +func (s *SIPService) CreateSIPOutboundTrunk(ctx context.Context, req *livekit.CreateSIPOutboundTrunkRequest) (*livekit.SIPOutboundTrunkInfo, error) { + if s.store == nil { + return nil, ErrSIPNotConnected + } + info := req.Trunk + if info == nil { + return nil, errors.New("trunk info is required") + } else if info.SipTrunkId != "" { + return nil, errors.New("trunk ID must be empty") + } + + // No additional validation needed for outbound. + info.SipTrunkId = guid.New(utils.SIPTrunkPrefix) + if err := s.store.StoreSIPOutboundTrunk(ctx, info); err != nil { return nil, err } return info, nil @@ -112,6 +163,32 @@ func (s *SIPService) ListSIPTrunk(ctx context.Context, req *livekit.ListSIPTrunk return &livekit.ListSIPTrunkResponse{Items: trunks}, nil } +func (s *SIPService) ListSIPInboundTrunk(ctx context.Context, req *livekit.ListSIPInboundTrunkRequest) (*livekit.ListSIPInboundTrunkResponse, error) { + if s.store == nil { + return nil, ErrSIPNotConnected + } + + trunks, err := s.store.ListSIPInboundTrunk(ctx) + if err != nil { + return nil, err + } + + return &livekit.ListSIPInboundTrunkResponse{Items: trunks}, nil +} + +func (s *SIPService) ListSIPOutboundTrunk(ctx context.Context, req *livekit.ListSIPOutboundTrunkRequest) (*livekit.ListSIPOutboundTrunkResponse, error) { + if s.store == nil { + return nil, ErrSIPNotConnected + } + + trunks, err := s.store.ListSIPOutboundTrunk(ctx) + if err != nil { + return nil, err + } + + return &livekit.ListSIPOutboundTrunkResponse{Items: trunks}, nil +} + func (s *SIPService) DeleteSIPTrunk(ctx context.Context, req *livekit.DeleteSIPTrunkRequest) (*livekit.SIPTrunkInfo, error) { if s.store == nil { return nil, ErrSIPNotConnected @@ -199,13 +276,17 @@ func (s *SIPService) CreateSIPParticipantWithToken(ctx context.Context, req *liv log := logger.GetLogger() log = log.WithValues("callId", callID, "roomName", req.RoomName, "sipTrunk", req.SipTrunkId, "toUser", req.SipCallTo) - trunk, err := s.store.LoadSIPTrunk(ctx, req.SipTrunkId) + trunk, err := s.store.LoadSIPOutboundTrunk(ctx, req.SipTrunkId) if err != nil { log.Errorw("cannot get trunk to update sip participant", err) return nil, err } - log = log.WithValues("fromUser", trunk.OutboundNumber, "toHost", trunk.OutboundAddress) - ireq := rpc.NewCreateSIPParticipantRequest(callID, wsUrl, token, req, trunk) + ireq, err := rpc.NewCreateSIPParticipantRequest(callID, wsUrl, token, req, trunk) + if err != nil { + log.Errorw("cannot create sip participant request", err) + return nil, err + } + log = log.WithValues("fromUser", ireq.Number, "toHost", trunk.Address) // CreateSIPParticipant will wait for LiveKit Participant to be created and that can take some time. // Thus, we must set a higher deadline for it, if it's not set already. diff --git a/pkg/service/utils_test.go b/pkg/service/utils_test.go index 99c19ac35..62d2975b8 100644 --- a/pkg/service/utils_test.go +++ b/pkg/service/utils_test.go @@ -15,7 +15,9 @@ package service_test import ( + "context" "testing" + "time" "github.com/redis/go-redis/v9" "github.com/stretchr/testify/require" @@ -23,10 +25,38 @@ import ( "github.com/livekit/livekit-server/pkg/service" ) -func redisClient() *redis.Client { - return redis.NewClient(&redis.Options{ +func redisClient(t testing.TB) *redis.Client { + cli := redis.NewClient(&redis.Options{ Addr: "localhost:6379", }) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + err := cli.Ping(ctx).Err() + if err == nil { + t.Cleanup(func() { + _ = cli.Close() + }) + return cli + } + _ = cli.Close() + t.Logf("local redis not available: %v", err) + + t.Logf("starting redis in docker") + addr := runRedis(t) + cli = redis.NewClient(&redis.Options{ + Addr: addr, + }) + ctx, cancel = context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + if err = cli.Ping(ctx).Err(); err != nil { + _ = cli.Close() + t.Fatal(err) + } + t.Cleanup(func() { + _ = cli.Close() + }) + return cli } func TestIsValidDomain(t *testing.T) {