Compare commits

..

14 Commits

Author SHA1 Message Date
Cédric Verstraeten
81cd95379b Merge pull request #315 from kerberos-io/fix/backchannel-reconnect
fix/backchannel-reconnect
2026-08-10 17:49:16 +02:00
Cédric Verstraeten
cd50f58138 Upgrade Go and RTSP dependencies
Upgrade to Go 1.25, gortsplib v5, and refreshed Pion dependencies. Adapt RTSP client APIs and normalize RTP timestamps for FPS tracking, with coverage for PTS conversion.
2026-08-10 15:46:02 +02:00
Cédric Verstraeten
b2f029117e Fix test cases by adding missing newlines and improving timeout handling 2026-08-10 12:52:37 +00:00
Cédric Verstraeten
db135acea9 Implement backchannel reconnection logic and enhance test coverage for write failures 2026-08-10 12:47:10 +00:00
Cédric Verstraeten
7b589b53f9 Harden RTSP backchannel streaming
Add paced, randomized RTP packetization with talkspurt markers and rollover-safe timestamps. Reconnect failed backchannel sessions with cancellable exponential backoff, initialize audio channels during bootstrap, and cover packetizer behavior with tests.
2026-08-10 14:45:18 +02:00
Cédric Verstraeten
c1740c752e Merge pull request #313 from kerberos-io/fix/moq-recovery-strategy
fix/moq-recovery-strategy
2026-08-07 16:28:27 +02:00
Cédric Verstraeten
5862786381 Deduplicate repeated H.264 keyframes
Remove exact duplicate IDR NALUs during normalization and drop repeated keyframes within a short timestamp window. Add normalization statistics, logging, reset handling, and coverage for deduplication behavior.
2026-08-07 15:46:50 +02:00
Cédric Verstraeten
e8dd64f54b Enhance MoQ streaming: implement quality tier broadcasting and subscriber management 2026-08-07 13:32:20 +00:00
Cédric Verstraeten
ba96b63002 Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-08-07 13:28:38 +02:00
Cédric Verstraeten
c7c6bcbdf2 Add live stream recovery gating
Drop stale H.264 packets until a recent keyframe arrives, with lifecycle logging and slow MoQ write diagnostics. Add focused FrameGate tests and configure the UI package registry.
2026-08-07 13:23:44 +02:00
Cédric Verstraeten
faa3b4eabb Merge pull request #312 from kerberos-io/fix/moq-double-pts-insertion
fix/moq-double-pts-insertion
2026-08-06 22:46:30 +02:00
Cédric Verstraeten
f0a6eb7d98 Merge pull request #311 from kerberos-io/fix/preserve-recording-fps-precision
preserve-recording-fps-precision
2026-08-06 21:19:18 +02:00
Cédric Verstraeten
33a58cddf7 Fix live stream presentation timestamps
Use capture presentation time directly for MoQ timestamps instead of adding composition time, which is already reflected in PTS.
2026-08-06 21:05:13 +02:00
Kilian Boute
fea6d81246 prevent fps rounding 2026-08-06 17:46:38 +02:00
23 changed files with 1073 additions and 222 deletions

View File

@@ -1,4 +1,4 @@
FROM mcr.microsoft.com/devcontainers/go:1.24-trixie
FROM mcr.microsoft.com/devcontainers/go:1.25-trixie
# Install node environment
RUN apt-get update && \

View File

@@ -1,5 +1,5 @@
ARG GO_IMAGE=golang:1.24-trixie
ARG GO_IMAGE=golang:1.25-trixie
ARG RUNTIME_IMAGE=debian:trixie-slim
ARG VERSION=0.0.0
FROM ${GO_IMAGE} AS build-machinery

View File

@@ -429,10 +429,20 @@ uses Debian Trixie. The publisher is disabled unless explicitly enabled at runti
-e AGENT_LIVE_MOQ_URL=https://relay.uug.ai/anon \
kerberos/agent
`AGENT_LIVE_MOQ_BROADCAST_PREFIX` defaults to `devices`, producing the broadcast
`devices/<agent-key>/live.hang`. `AGENT_LIVE_MOQ_QUALITY` accepts `auto` (the
default), `high`, or `low` and selects the main or sub camera stream when the
Agent starts. The initial implementation publishes H.264 video only.
`AGENT_LIVE_MOQ_BROADCAST_PREFIX` defaults to `devices`. MoQ viewers subscribe to
a relay and never negotiate with the Agent, so every quality tier is published as
its own broadcast and switching quality is simply a resubscribe:
| Tier | Broadcast | Source |
| ------ | ------------------------------------- | ------------------------------------------ |
| `high` | `devices/<agent-key>/live.hang` | highest-resolution camera stream |
| `low` | `devices/<agent-key>/live-low.hang` | sub stream (main stream when none is set) |
Each tier only uploads while it has at least one subscriber, so the tier nobody
watches costs virtually no bandwidth. `AGENT_LIVE_MOQ_QUALITY` accepts `high` or
`low` to pin the Agent to a single tier; viewers requesting the other tier then
find no broadcast. Any other value (including the default) publishes both. The
initial implementation publishes H.264 video only.
The `/anon` relay route is intended for interoperability testing. Production
deployments must set `AGENT_LIVE_MOQ_URL` to a short-lived, device-scoped

View File

@@ -1,6 +1,6 @@
module github.com/kerberos-io/agent/machinery
go 1.24.2
go 1.25.0
replace google.golang.org/genproto => google.golang.org/genproto v0.0.0-20250519155744-55703ea1f237
@@ -8,7 +8,7 @@ require (
github.com/Eyevinn/mp4ff v0.48.0
github.com/InVisionApp/conjungo v1.1.0
github.com/appleboy/gin-jwt/v2 v2.10.3
github.com/bluenviron/gortsplib/v4 v4.14.1
github.com/bluenviron/gortsplib/v5 v5.6.3
github.com/bluenviron/mediacommon v1.14.0
github.com/cedricve/go-onvif v0.0.0-20200222191200-567e8ce298f6
github.com/dromara/carbon/v2 v2.6.8
@@ -29,11 +29,11 @@ require (
github.com/moq-dev/moq-go v0.5.7
github.com/nfnt/resize v0.0.0-20180221191011-83c6a9932646
github.com/op/go-logging v0.0.0-20160315200505-970db520ece7
github.com/pion/interceptor v0.1.40
github.com/pion/rtp v1.8.19
github.com/pion/webrtc/v4 v4.1.2
github.com/pion/interceptor v0.1.47
github.com/pion/rtp v1.10.5
github.com/pion/webrtc/v4 v4.2.18
github.com/sirupsen/logrus v1.9.3
github.com/stretchr/testify v1.10.0
github.com/stretchr/testify v1.11.1
github.com/swaggo/files v1.0.1
github.com/swaggo/gin-swagger v1.6.0
github.com/swaggo/swag v1.16.4
@@ -53,7 +53,7 @@ require (
github.com/PuerkitoBio/purell v1.1.1 // indirect
github.com/PuerkitoBio/urlesc v0.0.0-20170810143723-de5bf2ad4578 // indirect
github.com/beevik/etree v1.2.0 // indirect
github.com/bluenviron/mediacommon/v2 v2.2.0 // indirect
github.com/bluenviron/mediacommon/v2 v2.9.2 // indirect
github.com/bytedance/sonic v1.13.2 // indirect
github.com/bytedance/sonic/loader v0.2.4 // indirect
github.com/cenkalti/backoff/v5 v5.0.2 // indirect
@@ -99,19 +99,19 @@ require (
github.com/moq-dev/moq-go-ffi v0.3.7 // indirect
github.com/nxadm/tail v1.4.11 // indirect
github.com/pelletier/go-toml/v2 v2.2.3 // indirect
github.com/pion/datachannel v1.5.10 // indirect
github.com/pion/dtls/v3 v3.0.6 // indirect
github.com/pion/ice/v4 v4.0.10 // indirect
github.com/pion/logging v0.2.3 // indirect
github.com/pion/mdns/v2 v2.0.7 // indirect
github.com/pion/datachannel v1.6.2 // indirect
github.com/pion/dtls/v3 v3.1.5 // indirect
github.com/pion/ice/v4 v4.4.0 // indirect
github.com/pion/logging v0.2.4 // indirect
github.com/pion/mdns/v2 v2.1.0 // indirect
github.com/pion/randutil v0.1.0 // indirect
github.com/pion/rtcp v1.2.15 // indirect
github.com/pion/sctp v1.8.39 // indirect
github.com/pion/sdp/v3 v3.0.13 // indirect
github.com/pion/srtp/v3 v3.0.5 // indirect
github.com/pion/stun/v3 v3.0.0 // indirect
github.com/pion/transport/v3 v3.0.7 // indirect
github.com/pion/turn/v4 v4.0.0 // indirect
github.com/pion/rtcp v1.2.17 // indirect
github.com/pion/sctp v1.11.1 // indirect
github.com/pion/sdp/v3 v3.0.19 // indirect
github.com/pion/srtp/v3 v3.0.12 // indirect
github.com/pion/stun/v3 v3.1.6 // indirect
github.com/pion/transport/v4 v4.0.2 // indirect
github.com/pion/turn/v5 v5.0.12 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/prometheus/procfs v0.15.1 // indirect
github.com/twitchyliquid64/golang-asm v0.15.1 // indirect
@@ -126,13 +126,14 @@ require (
go.opentelemetry.io/otel/metric v1.36.0 // indirect
go.opentelemetry.io/proto/otlp v1.6.0 // indirect
golang.org/x/arch v0.16.0 // indirect
golang.org/x/crypto v0.38.0 // indirect
golang.org/x/net v0.40.0 // indirect
golang.org/x/crypto v0.54.0 // indirect
golang.org/x/net v0.57.0 // indirect
golang.org/x/oauth2 v0.30.0 // indirect
golang.org/x/sync v0.14.0 // indirect
golang.org/x/sys v0.33.0 // indirect
golang.org/x/text v0.25.0 // indirect
golang.org/x/tools v0.30.0 // indirect
golang.org/x/sync v0.22.0 // indirect
golang.org/x/sys v0.47.0 // indirect
golang.org/x/text v0.40.0 // indirect
golang.org/x/time v0.14.0 // indirect
golang.org/x/tools v0.47.0 // indirect
google.golang.org/genproto/googleapis/api v0.0.0-20250519155744-55703ea1f237 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20250519155744-55703ea1f237 // indirect
google.golang.org/grpc v1.72.1 // indirect

View File

@@ -393,12 +393,12 @@ github.com/appleboy/gofight/v2 v2.1.2/go.mod h1:frW+U1QZEdDgixycTj4CygQ48yLTUhpl
github.com/bazelbuild/rules_go v0.49.0/go.mod h1:Dhcz716Kqg1RHNWos+N6MlXNkjNP2EwZQ0LukRKJfMs=
github.com/beevik/etree v1.2.0 h1:l7WETslUG/T+xOPs47dtd6jov2Ii/8/OjCldk5fYfQw=
github.com/beevik/etree v1.2.0/go.mod h1:aiPf89g/1k3AShMVAzriilpcE4R/Vuor90y83zVZWFc=
github.com/bluenviron/gortsplib/v4 v4.14.1 h1:v99NmXeeJFfbrO+ipPzPxYGibQaR5ZOUESOA9UQZhsI=
github.com/bluenviron/gortsplib/v4 v4.14.1/go.mod h1:3LaEcg0d47+kfXju5KSlsSxCiZ3IKBI/sqIrBPcsS64=
github.com/bluenviron/gortsplib/v5 v5.6.3 h1:OXvHthQZ9fZbLh6r3Go2wuF4XQ4/QW4WTIM2f4bv/W4=
github.com/bluenviron/gortsplib/v5 v5.6.3/go.mod h1:kzHgUtvl8NWNsQ5Vsez6Vuugk6ItFT4ByCPm8J/kbSQ=
github.com/bluenviron/mediacommon v1.14.0 h1:lWCwOBKNKgqmspRpwpvvg3CidYm+XOc2+z/Jw7LM5dQ=
github.com/bluenviron/mediacommon v1.14.0/go.mod h1:z5LP9Tm1ZNfQV5Co54PyOzaIhGMusDfRKmh42nQSnyo=
github.com/bluenviron/mediacommon/v2 v2.2.0 h1:fGXEX0OEvv5VhGHOv3Q2ABzOtSkIpl9UbwOHrnKWNTk=
github.com/bluenviron/mediacommon/v2 v2.2.0/go.mod h1:a6MbPmXtYda9mKibKVMZlW20GYLLrX2R7ZkUE+1pwV0=
github.com/bluenviron/mediacommon/v2 v2.9.2 h1:jvYeBjvhHKFOBRMTMm4hvrSjyOlCelkOkx6708DidQM=
github.com/bluenviron/mediacommon/v2 v2.9.2/go.mod h1:jMf/OJDaJl02xRgkLM2zbidUHnDYqLnO1dMMveCmyyU=
github.com/boombuler/barcode v1.0.0/go.mod h1:paBWMcWSl3LHKBqUq+rly7CNSldXjb2rDl3JlRe0mD8=
github.com/boombuler/barcode v1.0.1/go.mod h1:paBWMcWSl3LHKBqUq+rly7CNSldXjb2rDl3JlRe0mD8=
github.com/bytedance/sonic v1.13.2 h1:8/H1FempDZqC4VqjptGo14QQlJx8VdZJegxs6wwfqpQ=
@@ -866,38 +866,40 @@ github.com/phpdave11/gofpdf v1.4.2/go.mod h1:zpO6xFn9yxo3YLyMvW8HcKWVdbNqgIfOOp2
github.com/phpdave11/gofpdi v1.0.12/go.mod h1:vBmVV0Do6hSBHC8uKUQ71JGW+ZGQq74llk/7bXwjDoI=
github.com/phpdave11/gofpdi v1.0.13/go.mod h1:vBmVV0Do6hSBHC8uKUQ71JGW+ZGQq74llk/7bXwjDoI=
github.com/pierrec/lz4/v4 v4.1.18/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4=
github.com/pion/datachannel v1.5.10 h1:ly0Q26K1i6ZkGf42W7D4hQYR90pZwzFOjTq5AuCKk4o=
github.com/pion/datachannel v1.5.10/go.mod h1:p/jJfC9arb29W7WrxyKbepTU20CFgyx5oLo8Rs4Py/M=
github.com/pion/dtls/v3 v3.0.6 h1:7Hkd8WhAJNbRgq9RgdNh1aaWlZlGpYTzdqjy9x9sK2E=
github.com/pion/dtls/v3 v3.0.6/go.mod h1:iJxNQ3Uhn1NZWOMWlLxEEHAN5yX7GyPvvKw04v9bzYU=
github.com/pion/ice/v4 v4.0.10 h1:P59w1iauC/wPk9PdY8Vjl4fOFL5B+USq1+xbDcN6gT4=
github.com/pion/ice/v4 v4.0.10/go.mod h1:y3M18aPhIxLlcO/4dn9X8LzLLSma84cx6emMSu14FGw=
github.com/pion/interceptor v0.1.40 h1:e0BjnPcGpr2CFQgKhrQisBU7V3GXK6wrfYrGYaU6Jq4=
github.com/pion/interceptor v0.1.40/go.mod h1:Z6kqH7M/FYirg3frjGJ21VLSRJGBXB/KqaTIrdqnOic=
github.com/pion/logging v0.2.3 h1:gHuf0zpoh1GW67Nr6Gj4cv5Z9ZscU7g/EaoC/Ke/igI=
github.com/pion/logging v0.2.3/go.mod h1:z8YfknkquMe1csOrxK5kc+5/ZPAzMxbKLX5aXpbpC90=
github.com/pion/mdns/v2 v2.0.7 h1:c9kM8ewCgjslaAmicYMFQIde2H9/lrZpjBkN8VwoVtM=
github.com/pion/mdns/v2 v2.0.7/go.mod h1:vAdSYNAT0Jy3Ru0zl2YiW3Rm/fJCwIeM0nToenfOJKA=
github.com/pion/datachannel v1.6.2 h1:7EXQ8TH3vTouBUdRWYbcX2edSx9Yj6k5zl5P+qyxEPc=
github.com/pion/datachannel v1.6.2/go.mod h1:pzbdAZvyGtXbcHM1hBbsFaOTf40lZizU/dNlvVOak6E=
github.com/pion/dtls/v3 v3.1.5 h1:9xJtVsHwMYeSjPp5Hh1FTis4DchnQWtnOa5o+6ygqfc=
github.com/pion/dtls/v3 v3.1.5/go.mod h1:gz1K4jg6c+fq86oQMH4pilpCEOEPwmEr2jY+VcF/mkU=
github.com/pion/ice/v4 v4.4.0 h1:wvHDDqimaC38Y7MVpD46Y63p246ChvXd87VKoLYS5b4=
github.com/pion/ice/v4 v4.4.0/go.mod h1:obAyD+J+Hzs7QA7Y8YXHp5uIn6gb7z87pKedXZkrcFU=
github.com/pion/interceptor v0.1.47 h1:yw8t5pJ2f8t78NgU+8EmxhaqYLXS7uFCC/tAGOaSDBo=
github.com/pion/interceptor v0.1.47/go.mod h1:7yoRBzaIDETPC6cIN8Zj9EyGqHv1ImOpcTFPha6MuOM=
github.com/pion/logging v0.2.4 h1:tTew+7cmQ+Mc1pTBLKH2puKsOvhm32dROumOZ655zB8=
github.com/pion/logging v0.2.4/go.mod h1:DffhXTKYdNZU+KtJ5pyQDjvOAh/GsNSyv1lbkFbe3so=
github.com/pion/mdns/v2 v2.1.0 h1:3IJ9+Xio6tWYjhN6WwuY142P/1jA0D5ERaIqawg/fOY=
github.com/pion/mdns/v2 v2.1.0/go.mod h1:pcez23GdynwcfRU1977qKU0mDxSeucttSHbCSfFOd9A=
github.com/pion/randutil v0.1.0 h1:CFG1UdESneORglEsnimhUjf33Rwjubwj6xfiOXBa3mA=
github.com/pion/randutil v0.1.0/go.mod h1:XcJrSMMbbMRhASFVOlj/5hQial/Y8oH/HVo7TBZq+j8=
github.com/pion/rtcp v1.2.15 h1:LZQi2JbdipLOj4eBjK4wlVoQWfrZbh3Q6eHtWtJBZBo=
github.com/pion/rtcp v1.2.15/go.mod h1:jlGuAjHMEXwMUHK78RgX0UmEJFV4zUKOFHR7OP+D3D0=
github.com/pion/rtp v1.8.19 h1:jhdO/3XhL/aKm/wARFVmvTfq0lC/CvN1xwYKmduly3c=
github.com/pion/rtp v1.8.19/go.mod h1:bAu2UFKScgzyFqvUKmbvzSdPr+NGbZtv6UB2hesqXBk=
github.com/pion/sctp v1.8.39 h1:PJma40vRHa3UTO3C4MyeJDQ+KIobVYRZQZ0Nt7SjQnE=
github.com/pion/sctp v1.8.39/go.mod h1:cNiLdchXra8fHQwmIoqw0MbLLMs+f7uQ+dGMG2gWebE=
github.com/pion/sdp/v3 v3.0.13 h1:uN3SS2b+QDZnWXgdr69SM8KB4EbcnPnPf2Laxhty/l4=
github.com/pion/sdp/v3 v3.0.13/go.mod h1:88GMahN5xnScv1hIMTqLdu/cOcUkj6a9ytbncwMCq2E=
github.com/pion/srtp/v3 v3.0.5 h1:8XLB6Dt3QXkMkRFpoqC3314BemkpMQK2mZeJc4pUKqo=
github.com/pion/srtp/v3 v3.0.5/go.mod h1:r1G7y5r1scZRLe2QJI/is+/O83W2d+JoEsuIexpw+uM=
github.com/pion/stun/v3 v3.0.0 h1:4h1gwhWLWuZWOJIJR9s2ferRO+W3zA/b6ijOI6mKzUw=
github.com/pion/stun/v3 v3.0.0/go.mod h1:HvCN8txt8mwi4FBvS3EmDghW6aQJ24T+y+1TKjB5jyU=
github.com/pion/transport/v3 v3.0.7 h1:iRbMH05BzSNwhILHoBoAPxoB9xQgOaJk+591KC9P1o0=
github.com/pion/transport/v3 v3.0.7/go.mod h1:YleKiTZ4vqNxVwh77Z0zytYi7rXHl7j6uPLGhhz9rwo=
github.com/pion/turn/v4 v4.0.0 h1:qxplo3Rxa9Yg1xXDxxH8xaqcyGUtbHYw4QSCvmFWvhM=
github.com/pion/turn/v4 v4.0.0/go.mod h1:MuPDkm15nYSklKpN8vWJ9W2M0PlyQZqYt1McGuxG7mA=
github.com/pion/webrtc/v4 v4.1.2 h1:mpuUo/EJ1zMNKGE79fAdYNFZBX790KE7kQQpLMjjR54=
github.com/pion/webrtc/v4 v4.1.2/go.mod h1:xsCXiNAmMEjIdFxAYU0MbB3RwRieJsegSB2JZsGN+8U=
github.com/pion/rtcp v1.2.17 h1:PxiT6L79yPZKtXIsXdG1eakBl6dtBj4x+4oVEL0DlSw=
github.com/pion/rtcp v1.2.17/go.mod h1:7kBpuBJaWwax4hzc/pgexY8vkOpvh8atgYDbaKZq0iU=
github.com/pion/rtp v1.10.5 h1:ip0HhO/wYZqQ4bKS+R99KnZh/GRCmIT0jDXikub7vlE=
github.com/pion/rtp v1.10.5/go.mod h1:Au8fc6cEByy8RLTwKTQTEeQqDB/SJDxwL4mZuxYA5Pk=
github.com/pion/sctp v1.11.1 h1:O4dIFyURw1KTST7w+gtD4gLeYXkhPa0xXLHMMoe/OSA=
github.com/pion/sctp v1.11.1/go.mod h1:7KFmTwLcoYgJs/Z+99nJvsWL0qDpuyloSI0RbAqlrz0=
github.com/pion/sdp/v3 v3.0.19 h1:1VMKs3gIkTQV5M3hNKfTAPrDXSNrYtOlmOD8+mSZUGQ=
github.com/pion/sdp/v3 v3.0.19/go.mod h1:dE5WOSlzXrtiE/iuZqe9n+AcEbOjtAd3k5m5NtlV/qU=
github.com/pion/srtp/v3 v3.0.12 h1:U7V17bckl7sI4mb3sepiojByDuBY0wNCqQE+6IlQBbc=
github.com/pion/srtp/v3 v3.0.12/go.mod h1:EeZOi/sd6glM1EXapg051gdNWO9yWT1YSsgQ4SlJkns=
github.com/pion/stun/v3 v3.1.6 h1:WnhsD0eHCiwCfKNkVx0VJJwr2Y3eV4Ueih3KJ+dfZy8=
github.com/pion/stun/v3 v3.1.6/go.mod h1:zRUghXSQU32Lx5orJsz3uYMkIihweXb3mu5gIns02fs=
github.com/pion/transport/v3 v3.1.1 h1:Tr684+fnnKlhPceU+ICdrw6KKkTms+5qHMgw6bIkYOM=
github.com/pion/transport/v3 v3.1.1/go.mod h1:+c2eewC5WJQHiAA46fkMMzoYZSuGzA/7E2FPrOYHctQ=
github.com/pion/transport/v4 v4.0.2 h1:ifYlPqNwsy6aKQ9y8yzxXlHae5431ZrH2avkD/Rn6Tk=
github.com/pion/transport/v4 v4.0.2/go.mod h1:06hFI+jCFcok2X2MekVufNZ/uzNZXivGBPfviSVcjgM=
github.com/pion/turn/v5 v5.0.12 h1:6+b69ivQQXSlyfkp2AKripqD2k3W32qXK8QzCzpJWPI=
github.com/pion/turn/v5 v5.0.12/go.mod h1:CQACsRDJtjQ+6RSrGHrS2PCIerLwbW3uqXRqOvtjAFg=
github.com/pion/webrtc/v4 v4.2.18 h1:smA/3g6Gy4RohM0VIZ5KKY/12TQbxv3XFgpUMyb2EUI=
github.com/pion/webrtc/v4 v4.2.18/go.mod h1:vmzi6s+rvhoIuT94DPqivB+0xJXs9rG4QRD+4MgBtlY=
github.com/pkg/diff v0.0.0-20210226163009-20ebb0f2a09e/go.mod h1:pJLUxLENpZxwdsKMEsNbx1VGcRFpLqf3715MtcvvzbA=
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
@@ -957,8 +959,9 @@ github.com/stretchr/testify v1.8.2/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o
github.com/stretchr/testify v1.8.3/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo=
github.com/stretchr/testify v1.8.4/go.mod h1:sz/lmYIOXD/1dqDmKjjqLyZ2RngseejIcXlSw2iwfAo=
github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA=
github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U=
github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U=
github.com/substrait-io/substrait-go v0.4.2/go.mod h1:qhpnLmrcvAnlZsUyPXZRqldiHapPTXC3t7xFgDi3aQg=
github.com/swaggo/files v1.0.1 h1:J1bVJ4XHZNq0I46UU90611i9/YzdrF7x92oX1ig5IdE=
github.com/swaggo/files v1.0.1/go.mod h1:0qXmMNH6sXNf+73t65aKeB+ApmgxdnkQzVTAj2uaMUg=
@@ -1166,8 +1169,8 @@ golang.org/x/crypto v0.33.0/go.mod h1:bVdXmD7IV/4GdElGPozy6U7lWdRXA4qyRVGJV57uQ5
golang.org/x/crypto v0.35.0/go.mod h1:dy7dXNW32cAb/6/PRuTNsix8T+vJAqvuIy5Bli/x0YQ=
golang.org/x/crypto v0.36.0/go.mod h1:Y4J0ReaxCR1IMaabaSMugxJES1EpwhBHhv2bDHklZvc=
golang.org/x/crypto v0.37.0/go.mod h1:vg+k43peMZ0pUMhYmVAWysMK35e6ioLh3wB8ZCAfbVc=
golang.org/x/crypto v0.38.0 h1:jt+WWG8IZlBnVbomuhg2Mdq0+BBQaHbtqHEFEigjUV8=
golang.org/x/crypto v0.38.0/go.mod h1:MvrbAqul58NNYPKnOra203SB9vpuZW0e+RRZV+Ggqjw=
golang.org/x/crypto v0.54.0 h1:YLIA59K4fiNzHzjnZt2tUJQjQtUWfWbeHBqKtk3eScw=
golang.org/x/crypto v0.54.0/go.mod h1:KWL8ny2AZdGR2cWmzeHrp2azQPGogOv+HeQaVEXC2dk=
golang.org/x/exp v0.0.0-20180321215751-8460e604b9de/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
golang.org/x/exp v0.0.0-20180807140117-3d87b88a115f/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
@@ -1236,8 +1239,9 @@ golang.org/x/mod v0.15.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
golang.org/x/mod v0.17.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
golang.org/x/mod v0.18.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
golang.org/x/mod v0.19.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c=
golang.org/x/mod v0.23.0 h1:Zb7khfcRGKk+kqfxFaP5tZqCnDZMjC5VtUBs87Hr6QM=
golang.org/x/mod v0.23.0/go.mod h1:6SkKJ3Xj0I0BrPOZoBy3bdMptDDU9oJrpohJ3eWZ1fY=
golang.org/x/mod v0.37.0 h1:vF1DjpVEshcIqoEaauuHebaLk1O1forxjxBaVn884JQ=
golang.org/x/mod v0.37.0/go.mod h1:m8S8VeM9r4dzDwjrKO0a1sZP3YjeMamRRlD+fmR2Q/0=
golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
golang.org/x/net v0.0.0-20180826012351-8a410e7b638d/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
golang.org/x/net v0.0.0-20190108225652-1e06a53dbb7e/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4=
@@ -1325,8 +1329,8 @@ golang.org/x/net v0.34.0/go.mod h1:di0qlW3YNM5oh6GqDGQr92MyTozJPmybPK4Ev/Gm31k=
golang.org/x/net v0.35.0/go.mod h1:EglIi67kWsHKlRzzVMUD93VMSWGFOMSZgxFjparz1Qk=
golang.org/x/net v0.37.0/go.mod h1:ivrbrMbzFq5J41QOQh0siUuly180yBYtLp+CKbEaFx8=
golang.org/x/net v0.39.0/go.mod h1:X7NRbYVEA+ewNkCNyJ513WmMdQ3BineSwVtN2zD/d+E=
golang.org/x/net v0.40.0 h1:79Xs7wF06Gbdcg4kdCCIQArK11Z1hr5POQ6+fIYHNuY=
golang.org/x/net v0.40.0/go.mod h1:y0hY0exeL2Pku80/zKK7tpntoX23cqL3Oa6njdgRtds=
golang.org/x/net v0.57.0 h1:K5+3DljvIuDG9/Jv9rvyMywYNFCQ9RSUY6OOTTkT+tE=
golang.org/x/net v0.57.0/go.mod h1:KpXc8iv+r3XplLAG/f7Jsf9RPszJzdR0f58q9vGOuEU=
golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U=
golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw=
golang.org/x/oauth2 v0.0.0-20190604053449-0f29369cfe45/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw=
@@ -1405,8 +1409,9 @@ golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
golang.org/x/sync v0.11.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk=
golang.org/x/sync v0.12.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA=
golang.org/x/sync v0.13.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA=
golang.org/x/sync v0.14.0 h1:woo0S4Yywslg6hp4eUFjTVOyKt0RookbpAHG4c1HmhQ=
golang.org/x/sync v0.14.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA=
golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek=
golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190312061237-fead79001313/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
@@ -1515,8 +1520,8 @@ golang.org/x/sys v0.29.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.30.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.31.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
golang.org/x/sys v0.32.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
golang.org/x/sys v0.33.0 h1:q3i8TbbEz+JRD9ywIRlyRAQbM0qF7hu24q3teo2hbuw=
golang.org/x/sys v0.33.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k=
golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs=
golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE=
golang.org/x/telemetry v0.0.0-20240521205824-bda55230c457/go.mod h1:pRgIJT+bRLFKnoM1ldnzKoxTIn14Yxz928LQRYYgIN0=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
@@ -1583,8 +1588,8 @@ golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ=
golang.org/x/text v0.22.0/go.mod h1:YRoo4H8PVmsu+E3Ou7cqLVH8oXWIHVoX0jqUWALQhfY=
golang.org/x/text v0.23.0/go.mod h1:/BLNzu4aZCJ1+kcD0DNRotWKage4q2rGVAg4o22unh4=
golang.org/x/text v0.24.0/go.mod h1:L8rBsPeo2pSS+xqN0d5u2ikmjtmoJbDBT1b7nHvFCdU=
golang.org/x/text v0.25.0 h1:qVyWApTSYLk/drJRO5mDlNYskwQznZmkpV2c8q9zls4=
golang.org/x/text v0.25.0/go.mod h1:WEdwpYrmk1qmdHvhkSTNPm3app7v4rsT8F2UD6+VHIA=
golang.org/x/text v0.40.0 h1:Ub2Z6/xjgF1WrYQz2nuITOEegKFtiIy+rieRJ5lHZKs=
golang.org/x/text v0.40.0/go.mod h1:hpnzDAfGV753zIKo+wk3u1bVKCGPbrnF7+7LBF/UHVY=
golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
golang.org/x/time v0.0.0-20190308202827-9d24e82272b4/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
golang.org/x/time v0.0.0-20191024005414-555d28b269f0/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
@@ -1596,6 +1601,8 @@ golang.org/x/time v0.8.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM=
golang.org/x/time v0.9.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM=
golang.org/x/time v0.10.0/go.mod h1:3BpzKBy/shNhVucY/MWOyx10tF3SFh9QdLuxbVysPQM=
golang.org/x/time v0.11.0/go.mod h1:CDIdPxbZBQxdj6cxyCIdrNogrJKMJ7pr37NYpMcMDSg=
golang.org/x/time v0.14.0 h1:MRx4UaLrDotUKUdCIqzPC48t1Y9hANFKIRpNx+Te8PI=
golang.org/x/time v0.14.0/go.mod h1:eL/Oa2bBBK0TkX57Fyni+NgnyQQN4LitPmob2Hjnqw4=
golang.org/x/tools v0.0.0-20180525024113-a5b4c53f6e8b/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20190114222345-bf090417da8b/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
@@ -1669,8 +1676,9 @@ golang.org/x/tools v0.15.0/go.mod h1:hpksKq4dtpQWS1uQ61JkdqWM3LscIS6Slf+VVkm+wQk
golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk=
golang.org/x/tools v0.22.0/go.mod h1:aCwcsjqvq7Yqt6TNyX7QMU2enbQ/Gt0bo6krSeEri+c=
golang.org/x/tools v0.23.0/go.mod h1:pnu6ufv6vQkll6szChhK3C3L/ruaIv5eBeztNG8wtsI=
golang.org/x/tools v0.30.0 h1:BgcpHewrV5AUp2G9MebG4XPFI1E2W41zU1SaqVA9vJY=
golang.org/x/tools v0.30.0/go.mod h1:c347cR/OJfw5TI+GfX7RUPNMdDRRbjvYTS0jPyvsVtY=
golang.org/x/tools v0.47.0 h1:7Kn5x/d1svx/PzryTsqeoZN4TZwqeH5pGWjefhLi/1Q=
golang.org/x/tools v0.47.0/go.mod h1:dFHnyTvFWY212G+h7ZY4Vsp/K3U4/7W9TyVaAul8uCA=
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=

View File

@@ -17,15 +17,15 @@ import (
"time"
"unsafe"
"github.com/bluenviron/gortsplib/v4"
"github.com/bluenviron/gortsplib/v4/pkg/base"
"github.com/bluenviron/gortsplib/v4/pkg/description"
"github.com/bluenviron/gortsplib/v4/pkg/format"
"github.com/bluenviron/gortsplib/v4/pkg/format/rtph264"
"github.com/bluenviron/gortsplib/v4/pkg/format/rtph265"
"github.com/bluenviron/gortsplib/v4/pkg/format/rtplpcm"
"github.com/bluenviron/gortsplib/v4/pkg/format/rtpmpeg4audio"
"github.com/bluenviron/gortsplib/v4/pkg/format/rtpsimpleaudio"
"github.com/bluenviron/gortsplib/v5"
"github.com/bluenviron/gortsplib/v5/pkg/base"
"github.com/bluenviron/gortsplib/v5/pkg/description"
"github.com/bluenviron/gortsplib/v5/pkg/format"
"github.com/bluenviron/gortsplib/v5/pkg/format/rtph264"
"github.com/bluenviron/gortsplib/v5/pkg/format/rtph265"
"github.com/bluenviron/gortsplib/v5/pkg/format/rtplpcm"
"github.com/bluenviron/gortsplib/v5/pkg/format/rtpmpeg4audio"
"github.com/bluenviron/gortsplib/v5/pkg/format/rtpsimpleaudio"
"github.com/bluenviron/mediacommon/pkg/codecs/h264"
"github.com/bluenviron/mediacommon/pkg/codecs/h265"
"github.com/bluenviron/mediacommon/pkg/codecs/mpeg4audio"
@@ -317,11 +317,11 @@ func (g *Golibrtsp) Connect(ctx context.Context, ctxOtel context.Context) (err e
_, span := tracer.Start(ctxOtel, "Connect")
defer span.End()
transport := gortsplib.TransportTCP
protocol := gortsplib.ProtocolTCP
g.health = newStreamHealth()
g.Client = gortsplib.Client{
RequestBackChannels: false,
Transport: &transport,
Protocol: &protocol,
// Route gortsplib's packet-loss / decode-error reporting through our
// structured logger with stream context (replaces its plain stdout
// logging). These hooks are what let us tell whether the camera is
@@ -342,7 +342,9 @@ func (g *Golibrtsp) Connect(ctx context.Context, ctxOtel context.Context) (err e
}
// connect to the server
err = g.Client.Start(u.Scheme, u.Host)
g.Client.Scheme = u.Scheme
g.Client.Host = u.Host
err = g.Client.Start()
if err != nil {
log.Log.Debug("capture.golibrtsp.Connect(Start): " + err.Error())
}
@@ -593,10 +595,10 @@ func (g *Golibrtsp) ConnectBackChannel(ctx context.Context, ctxRunAgent context.
defer span.End()
// Transport TCP
transport := gortsplib.TransportTCP
protocol := gortsplib.ProtocolTCP
g.Client = gortsplib.Client{
RequestBackChannels: true,
Transport: &transport,
Protocol: &protocol,
}
// parse URL
u, err := base.ParseURL(g.Url)
@@ -606,7 +608,9 @@ func (g *Golibrtsp) ConnectBackChannel(ctx context.Context, ctxRunAgent context.
}
// connect to the server
err = g.Client.Start(u.Scheme, u.Host)
g.Client.Scheme = u.Scheme
g.Client.Host = u.Host
err = g.Client.Start()
if err != nil {
log.Log.Error("capture.golibrtsp.ConnectBackChannel(): " + err.Error())
}
@@ -680,6 +684,12 @@ func compositionOffsetMs(ext dtsExtractor, au [][]byte, pts int64, clockRate int
return offset * 1000 / int64(clockRate)
}
func ptsToDuration(pts int64, clockRate int) time.Duration {
rate := int64(clockRate)
return time.Duration(pts/rate)*time.Second +
time.Duration(pts%rate)*time.Second/time.Duration(rate)
}
// Start the RTSP client, and start reading packets.
func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets.Queue, configuration *models.Configuration, communication *models.Communication) (err error) {
log.Log.Debug("capture.golibrtsp.Start(): started")
@@ -693,13 +703,13 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets
// called when a MULAW audio RTP packet arrives
if g.AudioG711Media != nil && g.AudioG711Forma != nil {
g.Client.OnPacketRTP(g.AudioG711Media, g.AudioG711Forma, func(rtppkt *rtp.Packet) {
pts, ok := g.Client.PacketPTS(g.AudioG711Media, rtppkt)
// decode timestamp
pts2, ok := g.Client.PacketPTS2(g.AudioG711Media, rtppkt)
pts2, ok := g.Client.PacketPTS(g.AudioG711Media, rtppkt)
if !ok {
log.Log.Debug("capture.golibrtsp.Start(): " + "unable to get PTS")
return
}
pts := ptsToDuration(pts2, g.AudioG711Forma.ClockRate())
// extract LPCM samples from RTP packets
op, err := g.AudioG711Decoder.Decode(rtppkt)
@@ -729,12 +739,12 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets
if g.AudioMPEG4Media != nil && g.AudioMPEG4Forma != nil {
g.Client.OnPacketRTP(g.AudioMPEG4Media, g.AudioMPEG4Forma, func(rtppkt *rtp.Packet) {
// decode timestamp
pts, ok := g.Client.PacketPTS(g.AudioMPEG4Media, rtppkt)
pts2, ok := g.Client.PacketPTS2(g.AudioMPEG4Media, rtppkt)
pts2, ok := g.Client.PacketPTS(g.AudioMPEG4Media, rtppkt)
if !ok {
log.Log.Error("capture.golibrtsp.Start(): " + "unable to get PTS")
return
}
pts := ptsToDuration(pts2, g.AudioMPEG4Forma.ClockRate())
// Encode the AAC samples from RTP packets
// extract access units from RTP packets
@@ -788,12 +798,12 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets
if len(rtppkt.Payload) > 0 {
// decode timestamps — validate each call separately
pts, okPTS := g.Client.PacketPTS(g.VideoH264Media, rtppkt)
pts2, okPTS2 := g.Client.PacketPTS2(g.VideoH264Media, rtppkt)
pts2, okPTS2 := g.Client.PacketPTS(g.VideoH264Media, rtppkt)
if !okPTS2 {
log.Log.Debug("capture.golibrtsp.Start(): unable to get PTS2 from PacketPTS2")
log.Log.Debug("capture.golibrtsp.Start(): unable to get PTS")
return
}
pts := ptsToDuration(pts2, g.VideoH264Forma.ClockRate())
// Extract access units from RTP packets.
// We need a complete access unit to determine whether
@@ -807,15 +817,13 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets
}
// Frame is complete — update per-stream FPS from PTS.
if okPTS {
ft := g.fpsTrackers[g.VideoH264Index]
if ft == nil {
ft = newFPSTracker(30)
g.fpsTrackers[g.VideoH264Index] = ft
}
if ptsFPS := ft.update(pts); ptsFPS > 0 && ptsFPS <= 120 {
g.Streams[g.VideoH264Index].FPS = ptsFPS
}
ft := g.fpsTrackers[g.VideoH264Index]
if ft == nil {
ft = newFPSTracker(30)
g.fpsTrackers[g.VideoH264Index] = ft
}
if ptsFPS := ft.update(pts); ptsFPS > 0 && ptsFPS <= 120 {
g.Streams[g.VideoH264Index].FPS = ptsFPS
}
// We'll need to read out a few things.
@@ -1035,12 +1043,12 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets
if len(rtppkt.Payload) > 0 {
// decode timestamps — validate each call separately
pts, okPTS := g.Client.PacketPTS(g.VideoH265Media, rtppkt)
pts2, okPTS2 := g.Client.PacketPTS2(g.VideoH265Media, rtppkt)
pts2, okPTS2 := g.Client.PacketPTS(g.VideoH265Media, rtppkt)
if !okPTS2 {
log.Log.Debug("capture.golibrtsp.Start(): unable to get PTS")
return
}
pts := ptsToDuration(pts2, g.VideoH265Forma.ClockRate())
// Extract access units from RTP packets.
// We need a complete access unit to determine whether
@@ -1054,16 +1062,14 @@ func (g *Golibrtsp) Start(ctx context.Context, streamType string, queue *packets
}
// Frame is complete — update per-stream FPS from PTS.
if okPTS {
ft := g.fpsTrackers[g.VideoH265Index]
if ft == nil {
ft = newFPSTracker(30)
g.fpsTrackers[g.VideoH265Index] = ft
}
if ptsFPS := ft.update(pts); ptsFPS > 0 && ptsFPS <= 120 {
g.Streams[g.VideoH265Index].FPS = ptsFPS
g.persistStreamFPS(configuration, streamType, ptsFPS)
}
ft := g.fpsTrackers[g.VideoH265Index]
if ft == nil {
ft = newFPSTracker(30)
g.fpsTrackers[g.VideoH265Index] = ft
}
if ptsFPS := ft.update(pts); ptsFPS > 0 && ptsFPS <= 120 {
g.Streams[g.VideoH265Index].FPS = ptsFPS
g.persistStreamFPS(configuration, streamType, ptsFPS)
}
// Preserve the decoded access unit (in decode order) for DTS

View File

@@ -62,7 +62,7 @@ func recordingUploadMetadata(name, deviceKey string, timestamp int64, mp4Video *
}
value := mp4Video.AverageFPS()
if value > 0 && value <= 240 && !math.IsInf(value, 0) && !math.IsNaN(value) {
metadata.FPS = int(math.Floor(value))
metadata.FPS = value
}
return metadata
}

View File

@@ -6,11 +6,35 @@ import (
"os"
"path/filepath"
"testing"
"time"
"github.com/kerberos-io/agent/machinery/src/models"
"github.com/kerberos-io/agent/machinery/src/video"
)
func TestPTSToDuration(t *testing.T) {
tests := []struct {
name string
pts int64
clockRate int
want time.Duration
}{
{name: "one video second", pts: 90_000, clockRate: 90_000, want: time.Second},
{name: "one audio frame", pts: 1_024, clockRate: 8_000, want: 128 * time.Millisecond},
{name: "fractional millisecond", pts: 45_045, clockRate: 90_000, want: 500*time.Millisecond + 500*time.Microsecond},
{name: "negative timestamp", pts: -45_045, clockRate: 90_000, want: -500*time.Millisecond - 500*time.Microsecond},
{name: "large timestamp", pts: 90_000 * 60 * 60 * 24, clockRate: 90_000, want: 24 * time.Hour},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
if got := ptsToDuration(test.pts, test.clockRate); got != test.want {
t.Fatalf("ptsToDuration(%d, %d) = %s, want %s", test.pts, test.clockRate, got, test.want)
}
})
}
}
func TestQueueRecordingForUploadStoresFinalizedMetadata(t *testing.T) {
configDirectory := t.TempDir()
if err := os.MkdirAll(filepath.Join(configDirectory, "data", "cloud"), 0o755); err != nil {
@@ -29,9 +53,13 @@ func TestQueueRecordingForUploadStoresFinalizedMetadata(t *testing.T) {
if err := json.Unmarshal(got, &stored); err != nil {
t.Fatalf("decode upload marker: %v", err)
}
if stored.FileName != "recording.mp4" || stored.DeviceKey != "device-key" || stored.Timestamp != 1785934709414 || stored.Duration != 20452 || stored.FPS != 29 {
expectedFPS := mp4Video.AverageFPS()
if stored.FileName != "recording.mp4" || stored.DeviceKey != "device-key" || stored.Timestamp != 1785934709414 || stored.Duration != 20452 || math.Abs(stored.FPS-expectedFPS) > 1e-9 {
t.Fatalf("upload marker = %+v", stored)
}
if stored.FPS == math.Floor(stored.FPS) {
t.Fatalf("upload marker FPS = %v, want fractional precision", stored.FPS)
}
}
func TestQueueRecordingForUploadKeepsUnknownFPSCompatible(t *testing.T) {
@@ -44,7 +72,7 @@ func TestQueueRecordingForUploadKeepsUnknownFPSCompatible(t *testing.T) {
metadata := models.RecordingUploadMetadata{FileName: "recording.mp4"}
if fps >= 1 && fps <= 240 && !math.IsNaN(fps) && !math.IsInf(fps, 0) {
metadata.FPS = int(math.Floor(fps))
metadata.FPS = fps
}
queueRecordingForUpload(configDirectory, metadata)

View File

@@ -5,10 +5,16 @@ import (
"strings"
"github.com/bluenviron/mediacommon/pkg/codecs/h264"
"github.com/kerberos-io/agent/machinery/src/models"
)
var annexBStartCode = []byte{0x00, 0x00, 0x00, 0x01}
// H264NormalizationStats reports malformed duplication removed from an access unit.
type H264NormalizationStats struct {
DuplicateIDRNALUs int
}
// EnsureAnnexB restores the start code stripped by the Agent capture queue.
func EnsureAnnexB(payload []byte) []byte {
if hasAnnexBStartCode(payload) {
@@ -20,20 +26,28 @@ func EnsureAnnexB(payload []byte) []byte {
return append(framed, payload...)
}
// NormalizeH264AccessUnit removes delimiters and duplicate parameter sets that
// can make older MoQ splitters emit a parameter-only frame before the IDR.
// NormalizeH264AccessUnit removes delimiters and exact duplicate parameter-set
// or IDR NALUs that can confuse older MoQ splitters and decoders.
func NormalizeH264AccessUnit(payload []byte) ([]byte, error) {
normalized, _, err := NormalizeH264AccessUnitWithStats(payload)
return normalized, err
}
// NormalizeH264AccessUnitWithStats also reports exact duplicate IDR NALUs.
func NormalizeH264AccessUnitWithStats(payload []byte) ([]byte, H264NormalizationStats, error) {
nalus, err := h264.AnnexBUnmarshal(EnsureAnnexB(payload))
if err != nil {
return nil, err
return nil, H264NormalizationStats{}, err
}
stats := H264NormalizationStats{}
normalized := make([][]byte, 0, len(nalus))
for _, nalu := range nalus {
if len(nalu) == 0 || nalu[0]&0x1f == 9 {
continue
}
if nalu[0]&0x1f == 7 || nalu[0]&0x1f == 8 {
naluType := nalu[0] & 0x1f
if naluType == 7 || naluType == 8 || naluType == 5 {
duplicate := false
for _, existing := range normalized {
if bytes.Equal(existing, nalu) {
@@ -42,21 +56,45 @@ func NormalizeH264AccessUnit(payload []byte) ([]byte, error) {
}
}
if duplicate {
if naluType == 5 {
stats.DuplicateIDRNALUs++
}
continue
}
}
normalized = append(normalized, nalu)
}
return h264.AnnexBMarshal(normalized)
result, err := h264.AnnexBMarshal(normalized)
return result, stats, err
}
func BroadcastPath(prefix string, deviceKey string) string {
// BroadcastPath returns the relay path a quality tier is published on. Every
// tier gets its own broadcast so a viewer switches between the camera's main and
// sub stream by resubscribing to another path, without any control channel back
// to the Agent. The high tier keeps the historical ".../live.hang" path so
// existing viewers keep working; the low tier lives next to it on
// ".../live-low.hang".
func BroadcastPath(prefix string, deviceKey string, quality string) string {
prefix = strings.Trim(prefix, "/")
if prefix == "" {
prefix = "devices"
}
return prefix + "/" + strings.Trim(deviceKey, "/") + "/live.hang"
name := "live.hang"
if quality == models.StreamQualityLow {
name = "live-low.hang"
}
return prefix + "/" + strings.Trim(deviceKey, "/") + "/" + name
}
// TimestampUs converts the capture presentation timestamp from milliseconds.
// CompositionTime must not be added: it is already represented in the PTS and
// is only used by muxers to derive DTS for streams containing B-frames.
func TimestampUs(presentationTimeMs int64) uint64 {
if presentationTimeMs < 0 {
return 0
}
return uint64(presentationTimeMs) * 1000
}
func hasAnnexBStartCode(payload []byte) bool {

View File

@@ -3,6 +3,8 @@ package livemoq
import (
"bytes"
"testing"
"github.com/kerberos-io/agent/machinery/src/models"
)
func TestEnsureAnnexB(t *testing.T) {
@@ -64,11 +66,51 @@ func TestNormalizeH264AccessUnit(t *testing.T) {
}
}
func TestBroadcastPath(t *testing.T) {
if got := BroadcastPath("/devices/", "/camera-1/"); got != "devices/camera-1/live.hang" {
t.Fatalf("BroadcastPath() = %q", got)
func TestNormalizeH264AccessUnitRemovesOnlyExactDuplicateIDRSlices(t *testing.T) {
startCode := []byte{0x00, 0x00, 0x00, 0x01}
idrSlice1 := []byte{0x65, 0x88, 0x84}
idrSlice2 := []byte{0x65, 0x44, 0x22}
payload := make([]byte, 0)
for _, nalu := range [][]byte{idrSlice1, idrSlice1, idrSlice2} {
payload = append(payload, startCode...)
payload = append(payload, nalu...)
}
if got := BroadcastPath("", "camera-1"); got != "devices/camera-1/live.hang" {
t.Fatalf("BroadcastPath() default = %q", got)
got, stats, err := NormalizeH264AccessUnitWithStats(payload)
if err != nil {
t.Fatal(err)
}
want := make([]byte, 0)
for _, nalu := range [][]byte{idrSlice1, idrSlice2} {
want = append(want, startCode...)
want = append(want, nalu...)
}
if !bytes.Equal(got, want) {
t.Fatalf("NormalizeH264AccessUnitWithStats() = %x, want %x", got, want)
}
if stats.DuplicateIDRNALUs != 1 {
t.Fatalf("DuplicateIDRNALUs = %d, want 1", stats.DuplicateIDRNALUs)
}
}
func TestBroadcastPath(t *testing.T) {
if got := BroadcastPath("/devices/", "/camera-1/", models.StreamQualityHigh); got != "devices/camera-1/live.hang" {
t.Fatalf("BroadcastPath() high = %q", got)
}
if got := BroadcastPath("", "camera-1", models.StreamQualityHigh); got != "devices/camera-1/live.hang" {
t.Fatalf("BroadcastPath() default = %q", got)
}
if got := BroadcastPath("", "camera-1", models.StreamQualityLow); got != "devices/camera-1/live-low.hang" {
t.Fatalf("BroadcastPath() low = %q", got)
}
}
func TestTimestampUs(t *testing.T) {
if got := TimestampUs(1234); got != 1_234_000 {
t.Fatalf("TimestampUs() = %d, want 1234000", got)
}
if got := TimestampUs(-1); got != 0 {
t.Fatalf("TimestampUs() negative = %d, want 0", got)
}
}

View File

@@ -0,0 +1,41 @@
package livemoq
import (
"crypto/sha256"
"time"
)
type KeyframeDeduplicator struct {
hasPrevious bool
timestampMs int64
capturedAtMs int64
observedAt time.Time
digest [sha256.Size]byte
}
func (d *KeyframeDeduplicator) Reset() {
*d = KeyframeDeduplicator{}
}
// IsDuplicate reports exact repeated keyframe access units observed close
// together. Distinct IDR slices within one access unit remain untouched.
func (d *KeyframeDeduplicator) IsDuplicate(timestampMs int64, capturedAtMs int64, payload []byte, observedAt time.Time, window time.Duration) bool {
digest := sha256.Sum256(payload)
duplicate := d.hasPrevious && d.timestampMs == timestampMs && d.digest == digest
if duplicate {
if capturedAtMs > 0 && d.capturedAtMs > 0 {
gap := time.Duration(capturedAtMs-d.capturedAtMs) * time.Millisecond
duplicate = gap >= 0 && gap <= window
} else {
gap := observedAt.Sub(d.observedAt)
duplicate = gap >= 0 && gap <= window
}
}
d.hasPrevious = true
d.timestampMs = timestampMs
d.capturedAtMs = capturedAtMs
d.observedAt = observedAt
d.digest = digest
return duplicate
}

View File

@@ -0,0 +1,60 @@
package livemoq
import (
"testing"
"time"
)
func TestKeyframeDeduplicator(t *testing.T) {
now := time.UnixMilli(10_000)
window := 500 * time.Millisecond
payload := []byte{0x00, 0x00, 0x00, 0x01, 0x65, 0x88}
deduplicator := KeyframeDeduplicator{}
if deduplicator.IsDuplicate(1_000, 10_000, payload, now, window) {
t.Fatal("first keyframe reported as duplicate")
}
if !deduplicator.IsDuplicate(1_000, 10_020, payload, now.Add(20*time.Millisecond), window) {
t.Fatal("exact repeated keyframe was not reported as duplicate")
}
if deduplicator.IsDuplicate(2_000, 11_000, payload, now.Add(time.Second), window) {
t.Fatal("same payload with a new timestamp reported as duplicate")
}
if deduplicator.IsDuplicate(2_000, 11_020, append(payload, 0x01), now.Add(1020*time.Millisecond), window) {
t.Fatal("different payload with the same timestamp reported as duplicate")
}
}
func TestKeyframeDeduplicatorAllowsTimestampReuseOutsideWindow(t *testing.T) {
now := time.UnixMilli(10_000)
payload := []byte{0x00, 0x00, 0x00, 0x01, 0x65, 0x88}
deduplicator := KeyframeDeduplicator{}
deduplicator.IsDuplicate(1_000, 10_000, payload, now, 500*time.Millisecond)
if deduplicator.IsDuplicate(1_000, 20_000, payload, now.Add(10*time.Second), 500*time.Millisecond) {
t.Fatal("later keyframe after timestamp reset reported as duplicate")
}
}
func TestKeyframeDeduplicatorFallsBackToObservationTime(t *testing.T) {
now := time.UnixMilli(10_000)
payload := []byte{0x65, 0x88}
deduplicator := KeyframeDeduplicator{}
deduplicator.IsDuplicate(1_000, 0, payload, now, 500*time.Millisecond)
if !deduplicator.IsDuplicate(1_000, 0, payload, now.Add(20*time.Millisecond), 500*time.Millisecond) {
t.Fatal("duplicate without capture time was not reported")
}
}
func TestKeyframeDeduplicatorReset(t *testing.T) {
now := time.UnixMilli(10_000)
payload := []byte{0x65, 0x88}
deduplicator := KeyframeDeduplicator{}
deduplicator.IsDuplicate(1_000, 10_000, payload, now, 500*time.Millisecond)
deduplicator.Reset()
if deduplicator.IsDuplicate(1_000, 10_020, payload, now.Add(20*time.Millisecond), 500*time.Millisecond) {
t.Fatal("first keyframe after reset reported as duplicate")
}
}

View File

@@ -0,0 +1,55 @@
package livemoq
import "time"
type FrameGateEvent uint8
const (
FrameGateEventNone FrameGateEvent = iota
FrameGateEventStarted
FrameGateEventLagging
FrameGateEventRecovered
)
// FrameGate keeps publication on a decodable, recent GOP.
type FrameGate struct {
started bool
recovering bool
}
// Reset closes the gate so publication resumes on the next keyframe. It is used
// when the publisher stopped writing for a reason unrelated to the stream health
// (no subscribers), so the next viewer never receives a partial GOP.
func (g *FrameGate) Reset() {
g.started = false
g.recovering = false
}
// Allow rejects stale frames and waits for a fresh keyframe before reopening.
func (g *FrameGate) Allow(isKeyFrame bool, capturedAtMs int64, now time.Time, maxAge time.Duration) (bool, FrameGateEvent) {
if capturedAtMs > 0 && now.Sub(time.UnixMilli(capturedAtMs)) > maxAge {
event := FrameGateEventNone
if g.started {
if !g.recovering {
event = FrameGateEventLagging
}
g.started = false
g.recovering = true
}
return false, event
}
if !g.started {
if !isKeyFrame {
return false, FrameGateEventNone
}
g.started = true
if g.recovering {
g.recovering = false
return true, FrameGateEventRecovered
}
return true, FrameGateEventStarted
}
return true, FrameGateEventNone
}

View File

@@ -0,0 +1,49 @@
package livemoq
import (
"testing"
"time"
)
func TestFrameGateRecoversAtFreshKeyframe(t *testing.T) {
now := time.UnixMilli(10_000)
maxAge := 1500 * time.Millisecond
gate := FrameGate{}
tests := []struct {
name string
isKeyFrame bool
capturedAtMs int64
wantAllowed bool
wantEvent FrameGateEvent
}{
{name: "waits for initial keyframe", capturedAtMs: 10_000},
{name: "starts at initial keyframe", isKeyFrame: true, capturedAtMs: 10_000, wantAllowed: true, wantEvent: FrameGateEventStarted},
{name: "publishes fresh delta", capturedAtMs: 10_020, wantAllowed: true},
{name: "detects stale packet", capturedAtMs: 8_000, wantEvent: FrameGateEventLagging},
{name: "rejects fresh delta while recovering", capturedAtMs: 10_040},
{name: "rejects stale keyframe without duplicate event", isKeyFrame: true, capturedAtMs: 8_000},
{name: "recovers at fresh keyframe", isKeyFrame: true, capturedAtMs: 10_060, wantAllowed: true, wantEvent: FrameGateEventRecovered},
{name: "publishes delta after recovery", capturedAtMs: 10_080, wantAllowed: true},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
allowed, event := gate.Allow(test.isKeyFrame, test.capturedAtMs, now, maxAge)
if allowed != test.wantAllowed {
t.Fatalf("Allow() allowed = %t, want %t", allowed, test.wantAllowed)
}
if event != test.wantEvent {
t.Fatalf("Allow() event = %d, want %d", event, test.wantEvent)
}
})
}
}
func TestFrameGateAllowsMissingCaptureTime(t *testing.T) {
gate := FrameGate{}
allowed, event := gate.Allow(true, 0, time.Now(), time.Second)
if !allowed || event != FrameGateEventStarted {
t.Fatalf("Allow() = (%t, %d), want (true, %d)", allowed, event, FrameGateEventStarted)
}
}

View File

@@ -7,6 +7,7 @@ import (
"fmt"
"os"
"strings"
"sync/atomic"
"time"
"github.com/kerberos-io/agent/machinery/src/cloud/livemoq"
@@ -17,9 +18,13 @@ import (
)
const (
defaultMoQRelayURL = "https://relay.uug.ai/anon"
minMoQRetryDelay = time.Second
maxMoQRetryDelay = 30 * time.Second
defaultMoQRelayURL = "https://relay.uug.ai/anon"
minMoQRetryDelay = time.Second
maxMoQRetryDelay = 30 * time.Second
maxMoQLivePacketAge = 1500 * time.Millisecond
slowMoQWriteThreshold = 100 * time.Millisecond
moQWriteWarningInterval = 10 * time.Second
duplicateKeyframeWindow = 500 * time.Millisecond
)
type liveMoQConfig struct {
@@ -30,8 +35,23 @@ type liveMoQConfig struct {
queue *packets.Queue
}
// label identifies the tier in log lines, since one Agent runs a publisher per
// quality tier.
func (c liveMoQConfig) label() string {
return c.quality + " (" + c.sourceLabel + " stream)"
}
// StartLiveStreamMoQ starts the publisher only in the dedicated MoQ build and
// only when explicitly enabled by the deployment.
//
// Unlike WebRTC and HLS — where a viewer negotiates a session with the Agent and
// can therefore ask for another quality on the fly — MoQ viewers subscribe to a
// relay and never talk to the Agent. The quality selector is honoured by
// publishing each tier as its OWN broadcast (see livemoq.BroadcastPath): the
// high tier from the camera's highest-resolution stream and the low tier from
// its sub stream, so switching quality in the frontend is a resubscribe to the
// other path. Each tier only uploads while it actually has subscribers, so the
// second broadcast is close to free when nobody watches it.
func StartLiveStreamMoQ(configuration *models.Configuration, communication *models.Communication, subStreamEnabled bool) {
if os.Getenv("AGENT_LIVE_MOQ_ENABLED") != "true" {
return
@@ -47,39 +67,47 @@ func StartLiveStreamMoQ(configuration *models.Configuration, communication *mode
return
}
quality := os.Getenv("AGENT_LIVE_MOQ_QUALITY")
if quality == "" {
quality = models.StreamQualityAuto
}
useSub := models.SelectSubStreamForQuality(config, quality, subStreamEnabled)
queue := communication.Queue
sourceLabel := "main"
if useSub && communication.SubQueue != nil {
queue = communication.SubQueue
sourceLabel = "sub"
}
if queue == nil {
log.Log.Warning("cloud.StartLiveStreamMoQ(): selected packet queue is unavailable")
return
// Both tiers are published by default. AGENT_LIVE_MOQ_QUALITY pins the Agent
// to a single tier for deployments that must never publish the other one
// (viewers asking for the pinned-away tier then find no broadcast).
qualities := []string{models.StreamQualityHigh, models.StreamQualityLow}
switch strings.ToLower(strings.TrimSpace(os.Getenv("AGENT_LIVE_MOQ_QUALITY"))) {
case models.StreamQualityHigh:
qualities = []string{models.StreamQualityHigh}
case models.StreamQualityLow:
qualities = []string{models.StreamQualityLow}
}
relayURL := os.Getenv("AGENT_LIVE_MOQ_URL")
if relayURL == "" {
relayURL = defaultMoQRelayURL
}
publisherConfig := liveMoQConfig{
relayURL: relayURL,
broadcast: livemoq.BroadcastPath(os.Getenv("AGENT_LIVE_MOQ_BROADCAST_PREFIX"), config.Key),
quality: quality,
sourceLabel: sourceLabel,
queue: queue,
}
broadcastPrefix := os.Getenv("AGENT_LIVE_MOQ_BROADCAST_PREFIX")
ctx := context.Background()
if communication.Context != nil {
ctx = *communication.Context
}
go runLiveStreamMoQ(ctx, publisherConfig)
for _, quality := range qualities {
queue := communication.Queue
sourceLabel := "main"
if models.SelectSubStreamForQuality(config, quality, subStreamEnabled) && communication.SubQueue != nil {
queue = communication.SubQueue
sourceLabel = "sub"
}
if queue == nil {
log.Log.Warning("cloud.StartLiveStreamMoQ(): packet queue for the " + quality + " tier is unavailable")
continue
}
go runLiveStreamMoQ(ctx, liveMoQConfig{
relayURL: relayURL,
broadcast: livemoq.BroadcastPath(broadcastPrefix, config.Key, quality),
quality: quality,
sourceLabel: sourceLabel,
queue: queue,
})
}
}
func runLiveStreamMoQ(ctx context.Context, config liveMoQConfig) {
@@ -135,37 +163,123 @@ func publishLiveStreamMoQ(ctx context.Context, config liveMoQConfig) error {
}
defer stream.Finish()
// Only upload while this tier is actually being watched. `publishing` starts
// true so the track becomes discoverable on the relay even before the first
// subscriber ever arrives; from the moment a viewer has attached once, the
// subscriber watcher takes over and idles the tier again when everybody left.
watchCtx, cancelWatch := context.WithCancel(ctx)
defer cancelWatch()
publishing := &atomic.Bool{}
publishing.Store(true)
go watchLiveStreamMoQSubscribers(watchCtx, stream, publishing, config)
cursor := config.queue.Latest()
writing := false
gate := livemoq.FrameGate{}
deduplicator := livemoq.KeyframeDeduplicator{}
var lastSlowWriteWarning time.Time
var lastDuplicateKeyframeWarning time.Time
idle := false
for {
packet, err := cursor.ReadPacket()
if err != nil {
return fmt.Errorf("read packet: %w", err)
}
if !publishing.Load() {
// Keep draining the cursor so we stay at the live edge, but publish
// nothing. The gate is closed so the next viewer resumes on a keyframe.
if !idle {
gate.Reset()
deduplicator.Reset()
idle = true
}
continue
}
idle = false
if !packet.IsVideo || len(packet.Data) == 0 || !strings.EqualFold(packet.Codec, "H264") {
continue
}
if !writing {
if !packet.IsKeyFrame {
continue
}
writing = true
log.Log.Info("cloud.publishLiveStreamMoQ(): first H.264 keyframe received; broadcast is live")
allowed, event := gate.Allow(packet.IsKeyFrame, packet.CurrentTime, time.Now(), maxMoQLivePacketAge)
switch event {
case livemoq.FrameGateEventStarted:
log.Log.Info("cloud.publishLiveStreamMoQ(): first H.264 keyframe received; " + config.label() + " broadcast is live")
case livemoq.FrameGateEventLagging:
log.Log.Warning("cloud.publishLiveStreamMoQ(): " + config.label() + " stream is lagging; dropping packets until a recent keyframe")
case livemoq.FrameGateEventRecovered:
log.Log.Info("cloud.publishLiveStreamMoQ(): caught up with the " + config.label() + " live stream at a recent keyframe")
}
presentationTimeMs := packet.Time + packet.CompositionTime
if presentationTimeMs < 0 {
presentationTimeMs = 0
if !allowed {
continue
}
payload, err := livemoq.NormalizeH264AccessUnit(packet.Data)
payload, normalizationStats, err := livemoq.NormalizeH264AccessUnitWithStats(packet.Data)
if err != nil {
return fmt.Errorf("normalize H.264 access unit: %w", err)
}
if normalizationStats.DuplicateIDRNALUs > 0 && time.Since(lastDuplicateKeyframeWarning) >= moQWriteWarningInterval {
log.Log.Warning(fmt.Sprintf(
"cloud.publishLiveStreamMoQ(): %s removed %d duplicate IDR NALU(s) from H.264 keyframe (timestamp_ms=%d)",
config.label(), normalizationStats.DuplicateIDRNALUs, packet.Time,
))
lastDuplicateKeyframeWarning = time.Now()
}
if packet.IsKeyFrame && deduplicator.IsDuplicate(packet.Time, packet.CurrentTime, payload, time.Now(), duplicateKeyframeWindow) {
if time.Since(lastDuplicateKeyframeWarning) >= moQWriteWarningInterval {
log.Log.Warning(fmt.Sprintf(
"cloud.publishLiveStreamMoQ(): %s dropping duplicate H.264 keyframe (timestamp_ms=%d)",
config.label(), packet.Time,
))
lastDuplicateKeyframeWarning = time.Now()
}
continue
}
frame := moq.Frame{
Payload: payload,
TimestampUs: uint64(presentationTimeMs) * 1000,
TimestampUs: livemoq.TimestampUs(packet.Time),
}
writeStartedAt := time.Now()
if err := stream.WriteFrame(frame); err != nil {
return fmt.Errorf("write H.264 access unit: %w", err)
}
writeDuration := time.Since(writeStartedAt)
if writeDuration >= slowMoQWriteThreshold && time.Since(lastSlowWriteWarning) >= moQWriteWarningInterval {
packetAge := time.Duration(0)
if packet.CurrentTime > 0 {
packetAge = time.Since(time.UnixMilli(packet.CurrentTime))
if packetAge < 0 {
packetAge = 0
}
}
log.Log.Warning(fmt.Sprintf(
"cloud.publishLiveStreamMoQ(): %s WriteFrame blocked for %s (packet_age=%s keyframe=%t)",
config.label(), writeDuration.Round(time.Millisecond), packetAge.Round(time.Millisecond), packet.IsKeyFrame,
))
lastSlowWriteWarning = time.Now()
}
}
}
// watchLiveStreamMoQSubscribers flips the publisher between uploading and idling
// as viewers subscribe to and leave this tier's broadcast. Used and Unused both
// block, so they are followed from their own goroutine.
//
// It deliberately never turns publishing off before the first subscriber has
// been observed: the relay catalog is only complete once media has flowed, so
// going idle up front could keep the tier undiscoverable. On any error it fails
// open (keeps publishing) — a stalled watcher must never take the live view down.
func watchLiveStreamMoQSubscribers(ctx context.Context, stream *moq.MediaProducer, publishing *atomic.Bool, config liveMoQConfig) {
for ctx.Err() == nil {
if err := stream.Used(ctx); err != nil {
publishing.Store(true)
return
}
if publishing.CompareAndSwap(false, true) {
log.Log.Info("cloud.watchLiveStreamMoQSubscribers(): viewer subscribed, resuming the " + config.label() + " broadcast")
}
if err := stream.Unused(ctx); err != nil {
publishing.Store(true)
return
}
publishing.Store(false)
log.Log.Info("cloud.watchLiveStreamMoQSubscribers(): no viewers left, idling the " + config.label() + " broadcast")
}
}

View File

@@ -28,10 +28,10 @@ func queuedRecordingFPS(fileName string) string {
marker := strings.TrimSpace(string(value))
if strings.HasPrefix(marker, "{") {
metadata, ok := decodeRecordingUploadMetadata(value)
if !ok || metadata.FPS <= 0 || metadata.FPS > 240 {
if !ok || metadata.FPS <= 0 || metadata.FPS > 240 || math.IsInf(metadata.FPS, 0) || math.IsNaN(metadata.FPS) {
return ""
}
return strconv.Itoa(metadata.FPS)
return strconv.FormatFloat(metadata.FPS, 'f', -1, 64)
}
// Compatibility with markers created before upload metadata used JSON.

View File

@@ -272,7 +272,7 @@ func TestUploadVaultResumable_HappyPath(t *testing.T) {
fileName := "1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4"
payload := bytes.Repeat([]byte("x"), 4096)
withRecording(t, fileName, payload)
withQueuedRecordingFPS(t, fileName, `{"filename":"recording.mp4","device_key":"device-key","timestamp":1785934709414,"duration":20452,"fps":29}`)
withQueuedRecordingFPS(t, fileName, `{"filename":"recording.mp4","device_key":"device-key","timestamp":1785934709414,"duration":20452,"fps":29.97}`)
uploaded, responded, supported, _, err := uploadVaultResumable(testVault(ts.URL), "pk", "dev", fileName, "test", "primary")
if err != nil {
@@ -289,8 +289,8 @@ func TestUploadVaultResumable_HappyPath(t *testing.T) {
}
posts := srv.requestsForMethod(http.MethodPost)
metadata := decodeTusMetadata(posts[0].header.Get("Upload-Metadata"))
if got := metadata["fps"]; got != "29" {
t.Fatalf("POST metadata fps = %q, want %q", got, "29")
if got := metadata["fps"]; got != "29.97" {
t.Fatalf("POST metadata fps = %q, want %q", got, "29.97")
}
if got := metadata["duration"]; got != "20452" {
t.Fatalf("POST metadata duration = %q, want %q", got, "20452")
@@ -307,6 +307,7 @@ func TestQueuedRecordingFPSValidation(t *testing.T) {
want string
}{
{name: "json", fps: `{"fps":29}`, want: "29"},
{name: "json fractional", fps: `{"fps":17.35}`, want: "17.35"},
{name: "json with future field", fps: `{"fps":29,"codec":"h264"}`, want: "29"},
{name: "json without fps", fps: `{}`},
{name: "json invalid fps", fps: `{"fps":241}`},

View File

@@ -2,7 +2,9 @@ package components
import (
"bufio"
"context"
"fmt"
"math/rand"
"os"
"time"
@@ -15,6 +17,62 @@ import (
"github.com/zaf/g711"
)
const (
backchannelSampleRate = 8000
backchannelTalkspurtGap = 500 * time.Millisecond
backchannelReconnectInitial = time.Second
backchannelReconnectMax = 30 * time.Second
)
type backchannelClient interface {
ConnectBackChannel(ctx context.Context, otelContext context.Context) error
StartBackChannel(ctx context.Context, otelContext context.Context) error
WritePacket(pkt packets.Packet) error
Close(otelContext context.Context) error
}
type backchannelPacketizer struct {
sequenceNumber uint16
timestamp uint32
ssrc uint32
lastPacketAt time.Time
}
func newBackchannelPacketizer() backchannelPacketizer {
return backchannelPacketizer{
sequenceNumber: uint16(rand.Uint32()),
timestamp: rand.Uint32(),
ssrc: rand.Uint32(),
}
}
func (p *backchannelPacketizer) packet(audio models.AudioDataPartial, now time.Time) packets.Packet {
bufferUlaw := make([]byte, len(audio.Data))
for index, sample := range audio.Data {
bufferUlaw[index] = g711.EncodeUlawFrame(sample)
}
pkt := packets.Packet{
Packet: &rtp.Packet{
Header: rtp.Header{
Version: 2,
Marker: p.lastPacketAt.IsZero() || now.Sub(p.lastPacketAt) >= backchannelTalkspurtGap,
PayloadType: 0,
SequenceNumber: p.sequenceNumber,
Timestamp: p.timestamp,
SSRC: p.ssrc,
},
Payload: bufferUlaw,
},
}
p.timestamp += uint32(len(bufferUlaw))
p.sequenceNumber++
p.lastPacketAt = now
return pkt
}
func GetBackChannelAudioCodec(streams []av.CodecData, communication *models.Communication) av.AudioCodecData {
for _, stream := range streams {
if stream.Type().IsAudio() {
@@ -31,41 +89,115 @@ func GetBackChannelAudioCodec(streams []av.CodecData, communication *models.Comm
}
func WriteAudioToBackchannel(communication *models.Communication, rtspClient capture.RTSPClient) {
log.Log.Info("Audio.WriteAudioToBackchannel(): writing to backchannel audio codec")
length := uint32(0)
sequenceNumber := uint16(0)
for audio := range communication.HandleAudio {
// Encode PCM to MULAW
var bufferUlaw []byte
for _, v := range audio.Data {
b := g711.EncodeUlawFrame(v)
bufferUlaw = append(bufferUlaw, b)
}
pkt := packets.Packet{
Packet: &rtp.Packet{
Header: rtp.Header{
Version: 2,
Marker: true, // should be true
PayloadType: 0, //packet.PayloadType, // will be owerwriten
SequenceNumber: sequenceNumber,
Timestamp: uint32(length),
SSRC: 1293847657,
},
Payload: bufferUlaw,
},
}
err := rtspClient.WritePacket(pkt)
if err != nil {
log.Log.Error("Audio.WriteAudioToBackchannel(): error writing packet to backchannel")
}
length = (length + uint32(len(bufferUlaw))) % 65536
sequenceNumber = (sequenceNumber + 1) % 65535
time.Sleep(128 * time.Millisecond)
ctx := context.Background()
if communication.Context != nil {
ctx = *communication.Context
}
log.Log.Info("Audio.WriteAudioToBackchannel(): finished")
writeAudioToBackchannel(ctx, ctx, communication.HandleAudio, rtspClient)
}
func writeAudioToBackchannel(ctx context.Context, otelContext context.Context, audioChannel <-chan models.AudioDataPartial, rtspClient backchannelClient) {
log.Log.Info("Audio.WriteAudioToBackchannel(): writing to backchannel audio codec")
if err := rtspClient.StartBackChannel(ctx, otelContext); err != nil {
log.Log.Error("Audio.WriteAudioToBackchannel(): error starting backchannel: " + err.Error())
if !reconnectBackchannel(ctx, otelContext, rtspClient) {
log.Log.Info("Audio.WriteAudioToBackchannel(): stopped while reconnecting")
return
}
}
packetizer := newBackchannelPacketizer()
for {
select {
case <-ctx.Done():
log.Log.Info("Audio.WriteAudioToBackchannel(): stopped")
return
case audio, ok := <-audioChannel:
if !ok {
log.Log.Info("Audio.WriteAudioToBackchannel(): finished")
return
}
audio = latestBackchannelAudio(audio, audioChannel)
if len(audio.Data) == 0 {
continue
}
pkt := packetizer.packet(audio, time.Now())
if err := rtspClient.WritePacket(pkt); err != nil {
log.Log.Error("Audio.WriteAudioToBackchannel(): error writing packet to backchannel: " + err.Error())
if !reconnectBackchannel(ctx, otelContext, rtspClient) {
log.Log.Info("Audio.WriteAudioToBackchannel(): stopped while reconnecting")
return
}
packetizer = newBackchannelPacketizer()
continue
}
if !waitForBackchannel(ctx, time.Duration(len(audio.Data))*time.Second/backchannelSampleRate) {
log.Log.Info("Audio.WriteAudioToBackchannel(): stopped")
return
}
}
}
}
func latestBackchannelAudio(audio models.AudioDataPartial, audioChannel <-chan models.AudioDataPartial) models.AudioDataPartial {
for {
select {
case next, ok := <-audioChannel:
if !ok {
return audio
}
audio = next
default:
return audio
}
}
}
func reconnectBackchannel(ctx context.Context, otelContext context.Context, rtspClient backchannelClient) bool {
backoff := backchannelReconnectInitial
for {
if err := rtspClient.Close(otelContext); err != nil {
log.Log.Error("Audio.WriteAudioToBackchannel(): error closing failed backchannel: " + err.Error())
}
if ctx.Err() != nil {
return false
}
err := rtspClient.ConnectBackChannel(ctx, otelContext)
if err == nil {
err = rtspClient.StartBackChannel(ctx, otelContext)
}
if err == nil {
log.Log.Info("Audio.WriteAudioToBackchannel(): reconnected backchannel")
return true
}
log.Log.Error("Audio.WriteAudioToBackchannel(): error reconnecting backchannel: " + err.Error())
if !waitForBackchannel(ctx, backoff) {
return false
}
backoff *= 2
if backoff > backchannelReconnectMax {
backoff = backchannelReconnectMax
}
}
}
func waitForBackchannel(ctx context.Context, duration time.Duration) bool {
timer := time.NewTimer(duration)
defer timer.Stop()
select {
case <-ctx.Done():
return false
case <-timer.C:
return true
}
}
func WriteFileToBackChannel(infile av.DemuxCloser) {

View File

@@ -0,0 +1,209 @@
package components
import (
"context"
"errors"
"sync"
"testing"
"time"
"github.com/kerberos-io/agent/machinery/src/models"
"github.com/kerberos-io/agent/machinery/src/packets"
)
type fakeBackchannelClient struct {
mutex sync.Mutex
startErrors []error
connectError error
writeErrors []error
startCalls int
connectCalls int
closeCalls int
writeCalls int
connectAttempt chan struct{}
successfulWrite chan packets.Packet
}
func (f *fakeBackchannelClient) ConnectBackChannel(context.Context, context.Context) error {
f.mutex.Lock()
f.connectCalls++
err := f.connectError
f.mutex.Unlock()
select {
case f.connectAttempt <- struct{}{}:
default:
}
return err
}
func (f *fakeBackchannelClient) StartBackChannel(context.Context, context.Context) error {
f.mutex.Lock()
defer f.mutex.Unlock()
f.startCalls++
if len(f.startErrors) == 0 {
return nil
}
err := f.startErrors[0]
f.startErrors = f.startErrors[1:]
return err
}
func (f *fakeBackchannelClient) WritePacket(pkt packets.Packet) error {
f.mutex.Lock()
f.writeCalls++
var err error
if len(f.writeErrors) != 0 {
err = f.writeErrors[0]
f.writeErrors = f.writeErrors[1:]
}
f.mutex.Unlock()
if err == nil {
select {
case f.successfulWrite <- pkt:
default:
}
}
return err
}
func (f *fakeBackchannelClient) Close(context.Context) error {
f.mutex.Lock()
f.closeCalls++
f.mutex.Unlock()
return nil
}
func (f *fakeBackchannelClient) callCounts() (start, connect, close, write int) {
f.mutex.Lock()
defer f.mutex.Unlock()
return f.startCalls, f.connectCalls, f.closeCalls, f.writeCalls
}
func TestBackchannelPacketizerUsesFullRTPClock(t *testing.T) {
packetizer := backchannelPacketizer{ssrc: 1}
audio := models.AudioDataPartial{Data: make([]int16, 1024)}
startedAt := time.Unix(1, 0)
var timestamp uint32
for index := 0; index <= 64; index++ {
pkt := packetizer.packet(audio, startedAt.Add(time.Duration(index)*128*time.Millisecond))
timestamp = pkt.Packet.Timestamp
}
if timestamp != 65536 {
t.Fatalf("timestamp after 64 frames = %d, want 65536", timestamp)
}
}
func TestBackchannelPacketizerUsesNaturalSequenceRollover(t *testing.T) {
packetizer := backchannelPacketizer{sequenceNumber: ^uint16(0), ssrc: 1}
audio := models.AudioDataPartial{Data: []int16{0}}
startedAt := time.Unix(1, 0)
last := packetizer.packet(audio, startedAt)
firstAfterRollover := packetizer.packet(audio, startedAt.Add(time.Millisecond))
if last.Packet.SequenceNumber != ^uint16(0) {
t.Fatalf("last sequence number = %d, want %d", last.Packet.SequenceNumber, ^uint16(0))
}
if firstAfterRollover.Packet.SequenceNumber != 0 {
t.Fatalf("first sequence number after rollover = %d, want 0", firstAfterRollover.Packet.SequenceNumber)
}
}
func TestBackchannelPacketizerMarksTalkspurtStart(t *testing.T) {
packetizer := backchannelPacketizer{ssrc: 1}
audio := models.AudioDataPartial{Data: []int16{0}}
startedAt := time.Unix(1, 0)
first := packetizer.packet(audio, startedAt)
continuous := packetizer.packet(audio, startedAt.Add(128*time.Millisecond))
afterGap := packetizer.packet(audio, startedAt.Add(backchannelTalkspurtGap+128*time.Millisecond))
if !first.Packet.Marker {
t.Fatal("first packet must mark the start of a talkspurt")
}
if continuous.Packet.Marker {
t.Fatal("continuous packet must not carry the marker bit")
}
if !afterGap.Packet.Marker {
t.Fatal("packet after an audio gap must mark a new talkspurt")
}
}
func TestWriteAudioToBackchannelReconnectsAfterWriteFailure(t *testing.T) {
writeFailure := errors.New("EOF")
client := &fakeBackchannelClient{
writeErrors: []error{writeFailure, nil},
connectAttempt: make(chan struct{}, 1),
successfulWrite: make(chan packets.Packet, 1),
}
audioChannel := make(chan models.AudioDataPartial, 2)
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
writeAudioToBackchannel(ctx, ctx, audioChannel, client)
close(done)
}()
audioChannel <- models.AudioDataPartial{Data: make([]int16, 1024)}
select {
case <-client.connectAttempt:
case <-time.After(time.Second):
t.Fatal("backchannel was not reconnected after the write failure")
}
audioChannel <- models.AudioDataPartial{Data: make([]int16, 1024)}
select {
case pkt := <-client.successfulWrite:
if !pkt.Packet.Marker {
t.Fatal("first packet after reconnect must mark a new talkspurt")
}
case <-time.After(time.Second):
t.Fatal("fresh audio was not written after reconnect")
}
cancel()
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("backchannel writer did not stop after cancellation")
}
startCalls, connectCalls, closeCalls, writeCalls := client.callCounts()
if startCalls != 2 || connectCalls != 1 || closeCalls != 1 || writeCalls != 2 {
t.Fatalf("calls (start, connect, close, write) = (%d, %d, %d, %d), want (2, 1, 1, 2)", startCalls, connectCalls, closeCalls, writeCalls)
}
}
func TestWriteAudioToBackchannelCancellationStopsReconnect(t *testing.T) {
client := &fakeBackchannelClient{
connectError: errors.New("camera unavailable"),
writeErrors: []error{errors.New("EOF")},
connectAttempt: make(chan struct{}, 1),
successfulWrite: make(chan packets.Packet, 1),
}
audioChannel := make(chan models.AudioDataPartial, 1)
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
writeAudioToBackchannel(ctx, ctx, audioChannel, client)
close(done)
}()
audioChannel <- models.AudioDataPartial{Data: make([]int16, 1024)}
select {
case <-client.connectAttempt:
case <-time.After(time.Second):
t.Fatal("expected a reconnect attempt")
}
cancel()
select {
case <-done:
case <-time.After(250 * time.Millisecond):
t.Fatal("cancellation did not interrupt reconnect backoff")
}
}

View File

@@ -74,6 +74,7 @@ func Bootstrap(ctx context.Context, configDirectory string, configuration *model
communication.HandleLiveHDKeepalive = make(chan string, 1)
communication.HandleLiveHDPeers = make(chan string, 1)
communication.HandleLiveHLS = make(chan string, 1)
communication.HandleAudio = make(chan models.AudioDataPartial, 10)
communication.IsConfiguring = abool.New()
communication.IsRecordingManual = abool.New()
communication.RecordingManualHeartbeat = &atomic.Int64{}
@@ -269,11 +270,11 @@ func RunAgent(configDirectory string, configuration *models.Configuration, commu
communication.MainStreamConnected = true
// Try to create backchannel
communication.HasBackChannel = false
rtspBackChannelClient := captureDevice.SetBackChannelClient(rtspUrl)
err = rtspBackChannelClient.ConnectBackChannel(ctx, ctxRunAgent)
if err == nil {
log.Log.Info("components.Kerberos.RunAgent(): opened RTSP backchannel stream: " + rtspUrl)
go rtspBackChannelClient.StartBackChannel(ctx, ctxRunAgent)
}
rtspSubClient := captureDevice.RTSPSubClient
@@ -350,7 +351,6 @@ func RunAgent(configDirectory string, configuration *models.Configuration, commu
// is a no-op if ONVIFMotion is not enabled.
go onvif.HandleONVIFEventStream(*communication.Context, configuration, communication)
communication.HandleAudio = make(chan models.AudioDataPartial, 10)
if rtspBackChannelClient.HasBackChannel {
communication.HasBackChannel = true
go WriteAudioToBackchannel(communication, rtspBackChannelClient)
@@ -441,9 +441,6 @@ func RunAgent(configDirectory string, configuration *models.Configuration, commu
close(communication.HandleMotion)
communication.HandleMotion = nil
close(communication.HandleAudio)
communication.HandleAudio = nil
close(communication.HandleONVIF)
communication.HandleONVIF = nil

View File

@@ -11,11 +11,11 @@ const RecordingUploadMetadataExtension = ".metadata"
// with a recording. New optional fields can be added without changing the queue
// mechanism or breaking older agents.
type RecordingUploadMetadata struct {
FileName string `json:"filename"`
DeviceKey string `json:"device_key"`
Timestamp int64 `json:"timestamp"` // Unix milliseconds.
Duration uint64 `json:"duration"` // Milliseconds.
FPS int `json:"fps,omitempty"`
FileName string `json:"filename"`
DeviceKey string `json:"device_key"`
Timestamp int64 `json:"timestamp"` // Unix milliseconds.
Duration uint64 `json:"duration"` // Milliseconds.
FPS float64 `json:"fps,omitempty"`
}
// RecordingUploadMetadataFileName returns the queue marker name associated

View File

@@ -332,7 +332,7 @@ func MQTTListenerHandler(mqttClient mqtt.Client, hubKey string, configDirectory
case "record":
go HandleRecording(mqttClient, hubKey, payload, configuration, communication)
case "get-audio-backchannel":
go HandleAudio(mqttClient, hubKey, payload, configuration, communication)
HandleAudio(mqttClient, hubKey, payload, configuration, communication)
case "get-ptz-position":
go HandleGetPTZPosition(mqttClient, hubKey, payload, configuration, communication)
case "update-ptz-position":
@@ -442,10 +442,31 @@ func HandleAudio(mqttClient mqtt.Client, hubKey string, payload models.Payload,
Timestamp: audioPayload.Timestamp,
Data: audioPayload.Data,
}
communication.HandleAudio <- audioDataPartial
if enqueueLatestAudio(communication.HandleAudio, audioDataPartial) {
log.Log.Debug("routers.mqtt.main.HandleAudio(): dropped stale audio because the backchannel queue was full")
}
}
}
func enqueueLatestAudio(audioChannel chan models.AudioDataPartial, audio models.AudioDataPartial) bool {
select {
case audioChannel <- audio:
return false
default:
}
select {
case <-audioChannel:
default:
}
select {
case audioChannel <- audio:
default:
}
return true
}
func HandleGetPTZPosition(mqttClient mqtt.Client, hubKey string, payload models.Payload, configuration *models.Configuration, communication *models.Communication) {
value := payload.Value

View File

@@ -0,0 +1,39 @@
package mqtt
import (
"testing"
"time"
"github.com/kerberos-io/agent/machinery/src/models"
)
func TestEnqueueLatestAudioReplacesOldestFrameWhenFull(t *testing.T) {
audioChannel := make(chan models.AudioDataPartial, 2)
audioChannel <- models.AudioDataPartial{Timestamp: 1}
audioChannel <- models.AudioDataPartial{Timestamp: 2}
dropped := enqueueLatestAudio(audioChannel, models.AudioDataPartial{Timestamp: 3})
if !dropped {
t.Fatal("enqueueLatestAudio() dropped = false, want true")
}
first := <-audioChannel
second := <-audioChannel
if first.Timestamp != 2 || second.Timestamp != 3 {
t.Fatalf("queued timestamps = (%d, %d), want (2, 3)", first.Timestamp, second.Timestamp)
}
}
func TestEnqueueLatestAudioDoesNotBlockNilChannel(t *testing.T) {
done := make(chan struct{})
go func() {
enqueueLatestAudio(nil, models.AudioDataPartial{Timestamp: 1})
close(done)
}()
select {
case <-done:
case <-time.After(100 * time.Millisecond):
t.Fatal("enqueueLatestAudio() blocked on a nil channel")
}
}