mirror of
https://github.com/kerberos-io/agent.git
synced 2026-08-23 15:08:32 +00:00
Compare commits
23 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
203d7b5518 | ||
|
|
d0a7efff85 | ||
|
|
95ea92b9ce | ||
|
|
6890d1889c | ||
|
|
6c71ff5039 | ||
|
|
efc90c76c9 | ||
|
|
7459eb02ee | ||
|
|
fc46da1398 | ||
|
|
51a11edb71 | ||
|
|
29e7f26c0e | ||
|
|
01d270fcfa | ||
|
|
92d3311192 | ||
|
|
97a8c1fcaf | ||
|
|
de5b0666bd | ||
|
|
ffdb8b6f22 | ||
|
|
250e3b0b20 | ||
|
|
2bb8144e79 | ||
|
|
72d4fca63c | ||
|
|
81cd95379b | ||
|
|
cd50f58138 | ||
|
|
b2f029117e | ||
|
|
db135acea9 | ||
|
|
7b589b53f9 |
@@ -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 && \
|
||||
|
||||
2
.vscode/launch.json
vendored
2
.vscode/launch.json
vendored
@@ -17,7 +17,7 @@
|
||||
"8080"
|
||||
],
|
||||
"envFile": "${workspaceFolder}/machinery/.env.local",
|
||||
"buildFlags": "--tags dynamic",
|
||||
"buildFlags": "--tags dynamic,moq",
|
||||
"env": {
|
||||
"GOWORK": "off"
|
||||
},
|
||||
|
||||
@@ -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
|
||||
@@ -103,7 +103,8 @@ RUN apt-get update && apt-get install -y --no-install-recommends \
|
||||
ca-certificates curl ffmpeg libatomic1 libcap2-bin libstdc++6 && \
|
||||
rm -rf /var/lib/apt/lists/* && \
|
||||
groupadd --system kerberosio && \
|
||||
useradd --system --gid kerberosio --groups video --create-home agent
|
||||
useradd --system --gid kerberosio --groups video --create-home agent && \
|
||||
chmod 0755 /home/agent
|
||||
|
||||
#################################
|
||||
# Copy files from previous images
|
||||
|
||||
645
README-RTSPS-TLS.md
Normal file
645
README-RTSPS-TLS.md
Normal file
@@ -0,0 +1,645 @@
|
||||
# RTSPS and TLS certificates
|
||||
|
||||
This guide explains how Kerberos Agent connects to an IP camera over RTSPS,
|
||||
how to issue a camera certificate with a private CA, and how to validate the
|
||||
complete trust path. It also explains why some apparently corrupted trust
|
||||
bundles can still allow a connection.
|
||||
|
||||
The camera-specific steps were verified with a Bosch FLEXIDOME micro 3100i.
|
||||
Other Bosch firmware versions may use different labels or ports.
|
||||
|
||||
The commands were tested with Smallstep CLI `0.30.6` and OpenSSL `3.5.6` on
|
||||
Debian. Check `step certificate sign --help` when using an older Smallstep CLI.
|
||||
The OpenSSL isolation flags `-no-CApath` and `-no-CAstore` require a version that
|
||||
lists them in `openssl s_client -help`.
|
||||
|
||||
## Tested configuration
|
||||
|
||||
| Setting | Value |
|
||||
| --- | --- |
|
||||
| Camera | Bosch FLEXIDOME micro 3100i |
|
||||
| Example camera address | `10.0.30.11` |
|
||||
| RTSPS port | `9554` |
|
||||
| Main stream | `rtsps://<user>:<password>@10.0.30.11:9554/?inst=1` |
|
||||
| Sub stream | `rtsps://<user>:<password>@10.0.30.11:9554/?inst=2` |
|
||||
| Certificate SAN | `IP Address:10.0.30.11` |
|
||||
| Bosch certificate usage | `HTTPS` |
|
||||
| Agent trust input | Issuing intermediate plus root CA |
|
||||
|
||||
Replace the example address and certificate names throughout this guide. Keep
|
||||
camera credentials out of source control and percent-encode reserved URL
|
||||
characters in usernames and passwords.
|
||||
|
||||
## Mental model
|
||||
|
||||
### RTSPS, SRTSP, TLS, and SRTP
|
||||
|
||||
- The standard URL scheme is `rtsps://`. Do not use `srtsp://`.
|
||||
- Bosch interfaces and documentation may use SRTSP or Secure RTSP as product
|
||||
terminology.
|
||||
- RTSPS carries the RTSP control connection over TLS. With gortsplib, media is
|
||||
normally interleaved over the same TCP/TLS connection for this camera.
|
||||
- SRTP is a separate media protection mechanism and is negotiated only when the
|
||||
camera advertises an appropriate secure RTP profile.
|
||||
|
||||
Encryption alone does not prove which camera the Agent reached. Verified TLS
|
||||
also checks that:
|
||||
|
||||
1. The camera certificate is signed by a trusted authority.
|
||||
2. The certificate is valid at the current time.
|
||||
3. The URL host matches a certificate Subject Alternative Name (SAN).
|
||||
|
||||
Modern Go verification uses SANs for identity. A Common Name alone is not
|
||||
sufficient. Connecting to `10.0.30.11` requires an IP SAN with that exact value,
|
||||
not `DNS:10.0.30.11` and not only a device-name DNS SAN.
|
||||
|
||||
### Agent behavior
|
||||
|
||||
Kerberos Agent uses gortsplib for RTSP and RTSPS. With the normal configuration,
|
||||
gortsplib receives a nil custom TLS configuration and Go performs standard
|
||||
certificate and hostname verification with the process trust pool.
|
||||
|
||||
`AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=true` is an explicit escape hatch that
|
||||
sets `InsecureSkipVerify` for camera clients. It should be false in a verified
|
||||
deployment.
|
||||
|
||||
## Communication and certificate flow
|
||||
|
||||
The certificate is used during the TLS handshake, before the first RTSP command
|
||||
is exchanged. It is not attached to `DESCRIBE`, `SETUP`, or `PLAY`, and the CA
|
||||
trust bundle is never sent to the camera.
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant Agent as Kerberos Agent
|
||||
participant Trust as Go trust pool
|
||||
participant Camera as Camera RTSPS :9554
|
||||
|
||||
Agent->>Trust: Load system roots and append AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE
|
||||
Agent->>Camera: Open TCP connection
|
||||
Agent->>Camera: Send TLS ClientHello
|
||||
Camera-->>Agent: Send TLS ServerHello and camera certificate
|
||||
Agent->>Trust: Verify chain, validity, serverAuth, and URL host against SAN
|
||||
Trust-->>Agent: Accept or reject the camera identity
|
||||
Agent->>Camera: Complete TLS handshake
|
||||
Note over Agent,Camera: All following traffic is encrypted by TLS
|
||||
Agent->>Camera: DESCRIBE with RTSP authentication
|
||||
Camera-->>Agent: Return SDP and available media tracks
|
||||
Agent->>Camera: SETUP selected video and audio tracks over TCP
|
||||
Agent->>Camera: PLAY
|
||||
Camera-->>Agent: Send interleaved RTP and RTCP media over TLS
|
||||
```
|
||||
|
||||
The files and keys have distinct roles:
|
||||
|
||||
| Material | Location | Purpose | Sent over the connection |
|
||||
| --- | --- | --- | --- |
|
||||
| Camera leaf certificate | Camera | Identifies the camera and binds its public key to its SAN | Yes, by the camera during the TLS handshake |
|
||||
| Camera private key | Camera | Proves that the camera owns the presented certificate | No |
|
||||
| Intermediate and root CA PEM bundle | Agent | Lets Go build and trust the camera certificate chain | No |
|
||||
| RTSP username and password | Agent configuration or URL | Authenticates the Agent to the RTSP service after TLS succeeds | An authentication response is sent inside TLS; its form depends on the RTSP authentication method |
|
||||
|
||||
For an `rtsps://` URL, the Agent parses the URL and gives gortsplib the host and
|
||||
TLS settings. gortsplib opens the TCP connection and starts TLS. Go compares the
|
||||
certificate presented by the camera with the local trust pool, checks its
|
||||
validity period and server usage, and matches the URL hostname or IP address to
|
||||
the certificate SAN. Only a successful handshake creates the encrypted channel
|
||||
needed for the RTSP exchange.
|
||||
|
||||
The Agent then sends `DESCRIBE`, selects the advertised video and audio tracks,
|
||||
sends `SETUP`, and starts delivery with `PLAY`. For the tested camera, gortsplib
|
||||
uses interleaved TCP, so the RTSP control messages and RTP/RTCP media remain
|
||||
inside the same encrypted TLS connection. Main stream, sub stream, and enabled
|
||||
audio backchannel clients each establish and verify their own connection.
|
||||
|
||||
If certificate verification fails, the TLS handshake does not complete and no
|
||||
usable RTSP session is established. Setting
|
||||
`AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=true` keeps traffic encrypted but skips
|
||||
certificate-chain and hostname verification, so an attacker could impersonate
|
||||
the camera. It is not equivalent to trusting the camera certificate.
|
||||
|
||||
## Decide the certificate identity first
|
||||
|
||||
Choose the stable name used in every Agent URL before creating the certificate:
|
||||
|
||||
- For an IP URL, add that address as an IP SAN.
|
||||
- For a DNS URL, add the exact hostname as a DNS SAN.
|
||||
- Add both when clients legitimately use both forms.
|
||||
|
||||
A certificate stops matching if the camera IP changes. Use a static address,
|
||||
DHCP reservation, or stable DNS name.
|
||||
|
||||
## Configure RTSPS in the Bosch UI
|
||||
|
||||
1. Sign in to the camera as an administrator.
|
||||
2. Open **Configuration**.
|
||||
3. Open **Network > Network Services**.
|
||||
4. Enable **RTSPS**.
|
||||
5. Confirm port `9554`, or record the configured alternative.
|
||||
6. Click **Set**.
|
||||
|
||||
RTSP on port `554` and RTSPS on port `9554` are separate services. Enabling
|
||||
RTSPS does not make an `rtsp://` URL secure.
|
||||
|
||||
## Generate the private key and CSR on the camera
|
||||
|
||||
Keeping the TLS private key on the camera avoids exporting it to an operator
|
||||
workstation or deployment system.
|
||||
|
||||
1. Open **Service > Certificates**.
|
||||
2. Click **Add**.
|
||||
3. Select **Generate signing request**.
|
||||
4. Select `RSA 2048bit` or the stronger option supported by all clients.
|
||||
5. Enter a unique file name, such as `agent-rtsps`.
|
||||
6. Enter a descriptive Common Name and any required organization fields.
|
||||
7. Click **Generate**.
|
||||
8. Download the resulting CSR from the certificate table.
|
||||
|
||||
On the tested firmware, this form contains no SAN field. The downloaded CSR
|
||||
therefore has no IP SAN. The CA must add the SAN while signing.
|
||||
|
||||
Inspect the CSR before signing:
|
||||
|
||||
```bash
|
||||
openssl req -in camera.csr.pem -noout -verify -subject
|
||||
openssl req -in camera.csr.pem -noout -text
|
||||
```
|
||||
|
||||
The first command must report `Certificate request self-signature verify OK`.
|
||||
An absent `Subject Alternative Name` section is expected for this firmware.
|
||||
|
||||
## Prepare Smallstep
|
||||
|
||||
Use an existing organizational CA when one is available. Creating a new CA
|
||||
creates a new long-lived trust domain that must be distributed, protected,
|
||||
backed up, and eventually rotated.
|
||||
|
||||
### Install the CLI on Debian amd64
|
||||
|
||||
```bash
|
||||
curl -fsSL \
|
||||
https://dl.smallstep.com/cli/docs-ca-install/latest/step-cli_amd64.deb \
|
||||
-o /tmp/step-cli_amd64.deb
|
||||
sudo dpkg -i /tmp/step-cli_amd64.deb
|
||||
rm /tmp/step-cli_amd64.deb
|
||||
step version
|
||||
```
|
||||
|
||||
Use the official package matching the host architecture on other systems.
|
||||
|
||||
### Create a dedicated offline CA
|
||||
|
||||
Skip this section when using an existing CA.
|
||||
|
||||
```bash
|
||||
umask 077
|
||||
mkdir -p "$HOME/.step/secrets" "$HOME/.step/camera"
|
||||
|
||||
openssl rand -base64 48 > "$HOME/.step/secrets/camera_ca_password"
|
||||
chmod 600 "$HOME/.step/secrets/camera_ca_password"
|
||||
|
||||
step ca init \
|
||||
--pki \
|
||||
--name "UUG Camera CA" \
|
||||
--password-file "$HOME/.step/secrets/camera_ca_password"
|
||||
```
|
||||
|
||||
This produces:
|
||||
|
||||
```text
|
||||
$HOME/.step/certs/root_ca.crt
|
||||
$HOME/.step/certs/intermediate_ca.crt
|
||||
$HOME/.step/secrets/root_ca_key
|
||||
$HOME/.step/secrets/intermediate_ca_key
|
||||
$HOME/.step/secrets/camera_ca_password
|
||||
```
|
||||
|
||||
The files under `secrets/` are sensitive. Keep them mode `600`, never commit
|
||||
them, and back them up to encrypted persistent storage. A devcontainer can be
|
||||
rebuilt or deleted; it is not sufficient as the only CA backup.
|
||||
|
||||
## Add the SAN while signing
|
||||
|
||||
Copy the camera CSR into a protected working directory:
|
||||
|
||||
```bash
|
||||
cp camera.csr.pem "$HOME/.step/camera/camera.csr.pem"
|
||||
```
|
||||
|
||||
Create `$HOME/.step/camera/bosch-rtsps.tpl`:
|
||||
|
||||
```json
|
||||
{
|
||||
"subject": {
|
||||
"commonName": {{ toJson .Insecure.CR.Subject.CommonName }}
|
||||
},
|
||||
"ipAddresses": ["10.0.30.11"],
|
||||
"keyUsage": ["keyEncipherment", "digitalSignature"],
|
||||
"extKeyUsage": ["serverAuth", "clientAuth"]
|
||||
}
|
||||
```
|
||||
|
||||
The template preserves the camera CSR public key, sets the IP identity, and
|
||||
creates a TLS leaf rather than a CA certificate.
|
||||
|
||||
Sign it with a validity period that ends before the intermediate CA expires.
|
||||
A one-year leaf is preferable to a ten-year leaf when automated renewal is
|
||||
available:
|
||||
|
||||
```bash
|
||||
step certificate sign \
|
||||
--template "$HOME/.step/camera/bosch-rtsps.tpl" \
|
||||
--bundle \
|
||||
--not-after 8760h \
|
||||
--password-file "$HOME/.step/secrets/camera_ca_password" \
|
||||
"$HOME/.step/camera/camera.csr.pem" \
|
||||
"$HOME/.step/certs/intermediate_ca.crt" \
|
||||
"$HOME/.step/secrets/intermediate_ca_key" \
|
||||
> "$HOME/.step/camera/bosch-rtsps-chain.pem"
|
||||
```
|
||||
|
||||
For an online `step-ca`, do not assume `step ca sign` accepts a `--san` flag. It
|
||||
does not. Authorize SANs in the one-time token or configure a provisioner
|
||||
template that produces the required SANs.
|
||||
|
||||
## Validate before upload
|
||||
|
||||
Inspect the leaf certificate, which is the first PEM block in the chain file:
|
||||
|
||||
```bash
|
||||
openssl x509 \
|
||||
-in "$HOME/.step/camera/bosch-rtsps-chain.pem" \
|
||||
-noout -subject -issuer -dates -ext subjectAltName -ext extendedKeyUsage
|
||||
```
|
||||
|
||||
Confirm the SAN separately because some OpenSSL versions display only the last
|
||||
requested extension:
|
||||
|
||||
```bash
|
||||
openssl x509 \
|
||||
-in "$HOME/.step/camera/bosch-rtsps-chain.pem" \
|
||||
-noout -ext subjectAltName
|
||||
```
|
||||
|
||||
Verify the path and IP identity:
|
||||
|
||||
```bash
|
||||
openssl verify \
|
||||
-CAfile "$HOME/.step/certs/root_ca.crt" \
|
||||
-untrusted "$HOME/.step/certs/intermediate_ca.crt" \
|
||||
-verify_ip 10.0.30.11 \
|
||||
"$HOME/.step/camera/bosch-rtsps-chain.pem"
|
||||
```
|
||||
|
||||
Confirm that the signed leaf uses the exact public key from the camera CSR:
|
||||
|
||||
```bash
|
||||
csr_key=$(
|
||||
openssl req -in "$HOME/.step/camera/camera.csr.pem" -pubkey -noout |
|
||||
openssl pkey -pubin -outform DER 2>/dev/null |
|
||||
sha256sum | cut -d' ' -f1
|
||||
)
|
||||
|
||||
cert_key=$(
|
||||
openssl x509 -in "$HOME/.step/camera/bosch-rtsps-chain.pem" -pubkey -noout |
|
||||
openssl pkey -pubin -outform DER 2>/dev/null |
|
||||
sha256sum | cut -d' ' -f1
|
||||
)
|
||||
|
||||
test "$csr_key" = "$cert_key"
|
||||
```
|
||||
|
||||
Do not upload a certificate when any of these checks fail.
|
||||
|
||||
## Upload and assign the certificate
|
||||
|
||||
1. Return to **Service > Certificates**.
|
||||
2. Click **Add > Upload certificate**.
|
||||
3. Select the leaf-plus-intermediate PEM chain.
|
||||
4. Click **Upload** and wait for `100%`.
|
||||
5. Confirm that the former CSR row is now a `Certificate`.
|
||||
6. Confirm that the key icon is present. It proves that the camera associated
|
||||
the certificate with its retained private key.
|
||||
7. Open the new certificate's **Usage** selector.
|
||||
8. Select only **HTTPS**.
|
||||
9. Leave **CBS client** assigned to the original Bosch `DeviceCertificate`.
|
||||
10. Click **Set** and wait for the table to reload.
|
||||
|
||||
On the tested firmware, there is no separate SRTSP usage. RTSPS presents the
|
||||
certificate assigned to HTTPS. Reassigning HTTPS therefore changes both the
|
||||
web interface and RTSPS certificate.
|
||||
|
||||
After saving, the expected split is:
|
||||
|
||||
| Certificate | Usage |
|
||||
| --- | --- |
|
||||
| Private-CA camera certificate | `HTTPS` |
|
||||
| Bosch `DeviceCertificate` | `CBS client` |
|
||||
|
||||
The browser may warn about the new HTTPS certificate until the private root CA
|
||||
is trusted by the workstation.
|
||||
|
||||
## Account for the Bosch chain behavior
|
||||
|
||||
The tested firmware served only the leaf certificate on ports `443` and `9554`,
|
||||
even when the uploaded file contained the leaf and intermediate. Uploading the
|
||||
intermediate separately as a trusted camera certificate did not change the
|
||||
served chain.
|
||||
|
||||
Confirm the behavior:
|
||||
|
||||
```bash
|
||||
openssl s_client \
|
||||
-connect 10.0.30.11:9554 \
|
||||
-showcerts </dev/null 2>/dev/null |
|
||||
grep -c '^-----BEGIN CERTIFICATE-----$'
|
||||
```
|
||||
|
||||
A result of `1` means the client must already have the issuing intermediate.
|
||||
Create a portable trust bundle containing the intermediate and root:
|
||||
|
||||
```bash
|
||||
step certificate bundle \
|
||||
"$HOME/.step/certs/intermediate_ca.crt" \
|
||||
"$HOME/.step/certs/root_ca.crt" \
|
||||
"$HOME/.step/camera/uug-camera-trust-bundle.pem"
|
||||
|
||||
chmod 644 "$HOME/.step/camera/uug-camera-trust-bundle.pem"
|
||||
```
|
||||
|
||||
The trust bundle is public material. The CA private keys and password are not.
|
||||
|
||||
## Configure Kerberos Agent
|
||||
|
||||
For a process running directly in the same environment:
|
||||
|
||||
```dotenv
|
||||
AGENT_CAPTURE_IPCAMERA_RTSP="rtsps://<user>:<password>@10.0.30.11:9554/?inst=1"
|
||||
AGENT_CAPTURE_IPCAMERA_SUB_RTSP="rtsps://<user>:<password>@10.0.30.11:9554/?inst=2"
|
||||
AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=false
|
||||
AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE=/home/agent/data/config/uug-camera-trust-bundle.pem
|
||||
```
|
||||
|
||||
For a container, mount the public bundle read-only at the exact path visible
|
||||
inside the container. The Agent image creates `/home/agent/data/config` and
|
||||
includes Debian's `ca-certificates` package:
|
||||
|
||||
```bash
|
||||
docker run \
|
||||
-v /secure/config/uug-camera-trust-bundle.pem:/home/agent/data/config/uug-camera-trust-bundle.pem:ro \
|
||||
-e AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE=/home/agent/data/config/uug-camera-trust-bundle.pem \
|
||||
-e AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=false \
|
||||
-e 'AGENT_CAPTURE_IPCAMERA_RTSP=rtsps://<user>:<password>@10.0.30.11:9554/?inst=1' \
|
||||
-e 'AGENT_CAPTURE_IPCAMERA_SUB_RTSP=rtsps://<user>:<password>@10.0.30.11:9554/?inst=2' \
|
||||
kerberos/agent:latest
|
||||
```
|
||||
|
||||
`AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE` starts with the operating system's roots
|
||||
and appends the camera bundle only to the gortsplib TLS configuration. Other
|
||||
clients, including MoQ, Hub, and Vault, retain the normal public CA chain.
|
||||
Restart the Agent after changing trust files.
|
||||
|
||||
Do not set `SSL_CERT_FILE` or `SSL_CERT_DIR` in production solely for camera
|
||||
trust. They are process-wide and can prevent other clients from validating
|
||||
public services. For a deliberate process-wide isolation test, mount an empty
|
||||
directory and set `SSL_CERT_DIR` to its path:
|
||||
|
||||
```bash
|
||||
-v /secure/config/empty-ca-dir:/home/agent/data/config/empty-ca-dir:ro \
|
||||
-e SSL_CERT_DIR=/home/agent/data/config/empty-ca-dir
|
||||
```
|
||||
|
||||
Do not use `SSL_CERT_DIR=`. Go treats an empty value as unset and scans its
|
||||
default certificate directories.
|
||||
|
||||
Only use that mode when the Agent does not need public roots for other TLS
|
||||
connections.
|
||||
|
||||
## Validate the live endpoints
|
||||
|
||||
### Strict TLS and identity check
|
||||
|
||||
Use only the specified bundle, without OpenSSL's default CA locations:
|
||||
|
||||
```bash
|
||||
openssl s_client \
|
||||
-brief \
|
||||
-connect 10.0.30.11:9554 \
|
||||
-verify_ip 10.0.30.11 \
|
||||
-verify_return_error \
|
||||
-CAfile "$HOME/.step/camera/uug-camera-trust-bundle.pem" \
|
||||
-no-CApath \
|
||||
-no-CAstore \
|
||||
</dev/null
|
||||
```
|
||||
|
||||
Repeat with port `443`. Both must report `Verification: OK`.
|
||||
|
||||
Confirm that identity checking is active by repeating the command with a wrong
|
||||
address, such as `-verify_ip 10.0.30.12`. It must fail with an IP address
|
||||
mismatch.
|
||||
|
||||
### Confirm the live leaf is the generated leaf
|
||||
|
||||
```bash
|
||||
live_fingerprint=$(
|
||||
openssl s_client -connect 10.0.30.11:9554 -showcerts </dev/null 2>/dev/null |
|
||||
openssl x509 -noout -fingerprint -sha256 |
|
||||
cut -d= -f2
|
||||
)
|
||||
|
||||
local_fingerprint=$(
|
||||
openssl x509 \
|
||||
-in "$HOME/.step/camera/bosch-rtsps-chain.pem" \
|
||||
-noout -fingerprint -sha256 |
|
||||
cut -d= -f2
|
||||
)
|
||||
|
||||
test "$live_fingerprint" = "$local_fingerprint"
|
||||
```
|
||||
|
||||
### Validate the media path
|
||||
|
||||
A successful TLS handshake does not prove that RTSP authentication, DESCRIBE,
|
||||
SETUP, PLAY, and RTP delivery work. Start a fresh Agent with verified TLS and
|
||||
confirm that it connects without an x509 error and receives frames. During the
|
||||
verified setup described here, a gortsplib probe completed all RTSP operations
|
||||
and received an RTP packet over TCP.
|
||||
|
||||
## Why a tampered bundle may still connect
|
||||
|
||||
Editing PEM text is not always a useful negative TLS test.
|
||||
|
||||
### A certificate can still parse after a byte change
|
||||
|
||||
Base64 can remain syntactically valid when one character changes. OpenSSL may
|
||||
still list the certificate subject and issuer even though a signature is now
|
||||
invalid. Parsing and signature verification are different operations.
|
||||
|
||||
### Trust anchors are not validated through a parent
|
||||
|
||||
Every certificate loaded into Go's root pool is a trust anchor, including a
|
||||
non-self-signed intermediate CA. Verification can terminate at that certificate.
|
||||
|
||||
If tampering changes only the intermediate's signature from its parent root,
|
||||
but does not change its public key, that intermediate can still validate the
|
||||
camera leaf when it is trusted directly. Its now-invalid parent signature is
|
||||
not consulted at the trust boundary.
|
||||
|
||||
This is equivalent to OpenSSL's partial-chain behavior:
|
||||
|
||||
```bash
|
||||
openssl s_client \
|
||||
-connect 10.0.30.11:9554 \
|
||||
-verify_ip 10.0.30.11 \
|
||||
-verify_return_error \
|
||||
-partial_chain \
|
||||
-CAfile tampered-bundle.pem \
|
||||
-no-CApath \
|
||||
-no-CAstore \
|
||||
</dev/null
|
||||
```
|
||||
|
||||
### `SSL_CERT_FILE` does not isolate Go from CA directories
|
||||
|
||||
On Unix, Go uses `SSL_CERT_FILE` instead of its default aggregate CA file, but it
|
||||
still scans default certificate directories such as `/etc/ssl/certs`. Setting
|
||||
`SSL_CERT_FILE` alone therefore does not remove CA certificates installed with
|
||||
`update-ca-certificates`. `AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE` is appended
|
||||
after this system pool is loaded; it does not replace the system roots.
|
||||
|
||||
Blank values do not select empty trust sources. Both `SSL_CERT_FILE=` and
|
||||
`SSL_CERT_DIR=` are treated as unset, so Go falls back to its default aggregate
|
||||
CA file and certificate directories. To test with no trusted certificates on
|
||||
Linux, use a non-empty file path that contains no certificates and a non-empty
|
||||
directory path that contains no certificates:
|
||||
|
||||
```bash
|
||||
mkdir -p /tmp/empty-ca-dir
|
||||
SSL_CERT_FILE=/dev/null \
|
||||
SSL_CERT_DIR=/tmp/empty-ca-dir \
|
||||
AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE= \
|
||||
AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=false \
|
||||
GOWORK=off \
|
||||
go run -tags moq . -action run -port 8080
|
||||
```
|
||||
|
||||
That fresh process must fail with `x509: certificate signed by unknown
|
||||
authority`.
|
||||
|
||||
Use exactly one camera trust-distribution approach when possible:
|
||||
|
||||
1. Mount a private trust bundle and set
|
||||
`AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE`; or
|
||||
2. Install the CA certificates into the operating-system trust store.
|
||||
|
||||
Using both is valid, but makes isolation tests less obvious.
|
||||
|
||||
### Running processes can retain old roots
|
||||
|
||||
A long-running Go process may already have loaded and cached the trust pool.
|
||||
Always start a new process after changing trust configuration during a negative
|
||||
test.
|
||||
|
||||
## Perform a meaningful negative test
|
||||
|
||||
Do not corrupt only the root or intermediate signature. Instead, give a fresh
|
||||
Agent process a completely unrelated CA and hide the default CA directories.
|
||||
|
||||
```bash
|
||||
mkdir -p /tmp/empty-ca-dir
|
||||
|
||||
openssl req \
|
||||
-x509 -newkey rsa:2048 -nodes -days 1 \
|
||||
-subj '/CN=Unrelated Test Root' \
|
||||
-keyout /tmp/unrelated-test-root.key \
|
||||
-out /tmp/unrelated-test-root.crt
|
||||
|
||||
SSL_CERT_FILE=/dev/null \
|
||||
SSL_CERT_DIR=/tmp/empty-ca-dir \
|
||||
AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE=/tmp/unrelated-test-root.crt \
|
||||
AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=false \
|
||||
GOWORK=off \
|
||||
go run -tags moq . -action run -port 8080
|
||||
```
|
||||
|
||||
The connection must fail with an unknown-authority or chain-building error.
|
||||
Delete the temporary test key and certificate afterward.
|
||||
|
||||
To test bundle integrity rather than client distrust, validate the intermediate
|
||||
against the root explicitly:
|
||||
|
||||
```bash
|
||||
openssl verify \
|
||||
-CAfile "$HOME/.step/certs/root_ca.crt" \
|
||||
-no-CApath \
|
||||
-no-CAstore \
|
||||
"$HOME/.step/certs/intermediate_ca.crt"
|
||||
```
|
||||
|
||||
Store and compare approved SHA-256 fingerprints when detecting unauthorized
|
||||
certificate-file changes is a requirement.
|
||||
|
||||
Restore an accidentally edited bundle from the protected CA certificates, then
|
||||
restart the Agent:
|
||||
|
||||
```bash
|
||||
step certificate bundle -f \
|
||||
"$HOME/.step/certs/intermediate_ca.crt" \
|
||||
"$HOME/.step/certs/root_ca.crt" \
|
||||
"$HOME/.step/camera/uug-camera-trust-bundle.pem"
|
||||
|
||||
openssl verify \
|
||||
-CAfile "$HOME/.step/certs/root_ca.crt" \
|
||||
-no-CApath \
|
||||
-no-CAstore \
|
||||
"$HOME/.step/certs/intermediate_ca.crt"
|
||||
```
|
||||
|
||||
## Optional system trust installation
|
||||
|
||||
On Debian, install both public CA certificates when every process in the system
|
||||
should trust this camera PKI:
|
||||
|
||||
```bash
|
||||
sudo install -m 0644 \
|
||||
"$HOME/.step/certs/root_ca.crt" \
|
||||
/usr/local/share/ca-certificates/uug-camera-ca.crt
|
||||
|
||||
sudo install -m 0644 \
|
||||
"$HOME/.step/certs/intermediate_ca.crt" \
|
||||
/usr/local/share/ca-certificates/uug-camera-intermediate-ca.crt
|
||||
|
||||
sudo update-ca-certificates
|
||||
```
|
||||
|
||||
This creates links below `/etc/ssl/certs`. Remove those files and rerun
|
||||
`update-ca-certificates` before attempting an isolated trust-bundle test.
|
||||
|
||||
## Renewal and recovery
|
||||
|
||||
- Renew before the leaf or intermediate expires.
|
||||
- Generate a new camera CSR if the firmware cannot renew the existing key.
|
||||
- Sign the new CSR with all required SANs.
|
||||
- Upload and validate the new certificate before deleting the old one.
|
||||
- Preserve an alternate administrative access path while changing HTTPS usage.
|
||||
- Back up the CA certificates, encrypted CA keys, and password separately.
|
||||
- If the CA private keys are lost, create a new CA and redistribute its trust
|
||||
before replacing camera certificates.
|
||||
|
||||
## Production checklist
|
||||
|
||||
- [ ] The Agent URL uses `rtsps://`, not `rtsp://` or `srtsp://`.
|
||||
- [ ] RTSPS is enabled on the camera and the configured port is reachable.
|
||||
- [ ] The certificate SAN exactly matches the Agent URL host.
|
||||
- [ ] The leaf public key matches the camera-generated CSR.
|
||||
- [ ] The certificate has `serverAuth` extended key usage.
|
||||
- [ ] The certificate expires before its issuer.
|
||||
- [ ] HTTPS is assigned to the private-CA certificate.
|
||||
- [ ] CBS client remains assigned to the Bosch device certificate.
|
||||
- [ ] The Agent has the intermediate and root CA certificates it needs.
|
||||
- [ ] `AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=false`.
|
||||
- [ ] The Agent was restarted after trust changes.
|
||||
- [ ] A strict TLS check reports `Verification: OK`.
|
||||
- [ ] A real Agent connection receives RTP packets.
|
||||
- [ ] CA private keys and passwords are backed up outside the devcontainer.
|
||||
21
README.md
21
README.md
@@ -189,6 +189,21 @@ Next to attaching the configuration file, it is also possible to override the co
|
||||
-e AGENT_CAPTURE_CONTINUOUS=true \
|
||||
-d --restart=always kerberos/agent:latest
|
||||
|
||||
### Secure camera streams (RTSPS)
|
||||
|
||||
The Agent accepts `rtsps://` camera URLs. Do not use `srtsp://`; RTSPS is RTSP over TLS. For a Bosch FLEXIDOME micro 3100i, enable **Secure RTSP** under **Network > Network Services** and use port `9554`:
|
||||
|
||||
```bash
|
||||
AGENT_CAPTURE_IPCAMERA_RTSP='rtsps://username:password@camera.example:9554/?inst=1'
|
||||
AGENT_CAPTURE_IPCAMERA_SUB_RTSP='rtsps://username:password@camera.example:9554/?inst=2'
|
||||
```
|
||||
|
||||
Certificate verification is enabled by default. The URL hostname or IP address must match the camera certificate SAN. On this Bosch firmware, RTSPS presents the certificate assigned to **HTTPS**; there is no separate SRTSP certificate usage. Leave **CBS client** assigned to the Bosch device certificate.
|
||||
|
||||
For a private CA, mount a PEM trust bundle containing every CA certificate needed to build the camera certificate chain and set `AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE` to its path inside the Agent. The bundle is appended to the system roots for camera RTSPS connections only. This Bosch firmware presents only its leaf certificate, so include both the issuing intermediate and root certificates in the bundle. As a temporary fallback, `AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE=true` disables certificate verification for camera streams only.
|
||||
|
||||
See [RTSPS and TLS certificates](README-RTSPS-TLS.md) for the complete Bosch UI, private-CA, deployment, validation, and troubleshooting procedure.
|
||||
|
||||
| Name | Description | Default Value |
|
||||
| --------------------------------------- | ----------------------------------------------------------------------------------------------- | ------------------------------ |
|
||||
| `LOG_LEVEL` | Level for logging, could be "info", "warning", "debug", "error" or "fatal". | "info" |
|
||||
@@ -208,8 +223,10 @@ Next to attaching the configuration file, it is also possible to override the co
|
||||
| `AGENT_TIME` | Enable the timetable for Kerberos Agent | "false" |
|
||||
| `AGENT_TIMETABLE` | A (weekly) time table to specify when to make recordings "start1,end1,start2,end2;start1.. | "" |
|
||||
| `AGENT_REGION_POLYGON` | A single polygon set for motion detection: "x1,y1;x2,y2;x3,y3;... | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_RTSP` | Full-HD RTSP endpoint to the camera you're targetting. | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_SUB_RTSP` | Sub-stream RTSP endpoint used for livestreaming (WebRTC). | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_RTSP` | Full-HD RTSP or RTSPS endpoint for the target camera. | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_SUB_RTSP` | RTSP or RTSPS sub-stream endpoint used for livestreaming (WebRTC). | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE` | PEM CA bundle appended to the system roots for RTSPS camera certificate verification. | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE` | Disable RTSPS camera certificate verification; use only when a trusted CA cannot be installed. | "false" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_BASE_WIDTH` | Force a specific width resolution for live view processing. | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_BASE_HEIGHT` | Force a specific height resolution for live view processing. | "" |
|
||||
| `AGENT_CAPTURE_IPCAMERA_ONVIF` | Mark as a compliant ONVIF device. | "" |
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
{"upload_url":"https://vault.kerberos.io/api/storage/tus/19e42fbc666a38064904caf8c46d182a","vault_uri":"https://vault.kerberos.io/api/storage/tus/","size":1591581}
|
||||
@@ -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
|
||||
|
||||
100
machinery/go.sum
100
machinery/go.sum
@@ -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=
|
||||
|
||||
@@ -8,24 +8,27 @@ import "C"
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"errors"
|
||||
"fmt"
|
||||
"image"
|
||||
"os"
|
||||
"reflect"
|
||||
"strconv"
|
||||
"sync"
|
||||
"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"
|
||||
@@ -38,6 +41,36 @@ import (
|
||||
|
||||
var tracer = otel.Tracer("github.com/kerberos-io/agent/machinery/src/capture")
|
||||
|
||||
const (
|
||||
rtspsCAFileEnv = "AGENT_CAPTURE_IPCAMERA_RTSPS_CA_FILE"
|
||||
rtspsInsecureEnv = "AGENT_CAPTURE_IPCAMERA_RTSPS_INSECURE"
|
||||
)
|
||||
|
||||
func rtspsTLSConfig() (*tls.Config, error) {
|
||||
if os.Getenv(rtspsInsecureEnv) == "true" {
|
||||
return &tls.Config{InsecureSkipVerify: true}, nil // #nosec G402 -- explicit opt-in for cameras with self-signed certificates
|
||||
}
|
||||
|
||||
caFile := os.Getenv(rtspsCAFileEnv)
|
||||
if caFile == "" {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
rootCAs, err := x509.SystemCertPool()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("load system CA pool: %w", err)
|
||||
}
|
||||
caPEM, err := os.ReadFile(caFile)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("read RTSPS CA file %q: %w", caFile, err)
|
||||
}
|
||||
if !rootCAs.AppendCertsFromPEM(caPEM) {
|
||||
return nil, fmt.Errorf("RTSPS CA file %q contains no certificates", caFile)
|
||||
}
|
||||
|
||||
return &tls.Config{RootCAs: rootCAs}, nil
|
||||
}
|
||||
|
||||
// Implements the RTSPClient interface.
|
||||
type Golibrtsp struct {
|
||||
RTSPClient
|
||||
@@ -317,11 +350,17 @@ func (g *Golibrtsp) Connect(ctx context.Context, ctxOtel context.Context) (err e
|
||||
_, span := tracer.Start(ctxOtel, "Connect")
|
||||
defer span.End()
|
||||
|
||||
transport := gortsplib.TransportTCP
|
||||
tlsConfig, err := rtspsTLSConfig()
|
||||
if err != nil {
|
||||
return fmt.Errorf("configure RTSPS TLS: %w", err)
|
||||
}
|
||||
|
||||
protocol := gortsplib.ProtocolTCP
|
||||
g.health = newStreamHealth()
|
||||
g.Client = gortsplib.Client{
|
||||
RequestBackChannels: false,
|
||||
Transport: &transport,
|
||||
Protocol: &protocol,
|
||||
TLSConfig: tlsConfig,
|
||||
// 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 +381,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())
|
||||
}
|
||||
@@ -592,11 +633,17 @@ func (g *Golibrtsp) ConnectBackChannel(ctx context.Context, ctxRunAgent context.
|
||||
_, span := tracer.Start(ctxRunAgent, "ConnectBackChannel")
|
||||
defer span.End()
|
||||
|
||||
tlsConfig, err := rtspsTLSConfig()
|
||||
if err != nil {
|
||||
return fmt.Errorf("configure RTSPS TLS: %w", err)
|
||||
}
|
||||
|
||||
// Transport TCP
|
||||
transport := gortsplib.TransportTCP
|
||||
protocol := gortsplib.ProtocolTCP
|
||||
g.Client = gortsplib.Client{
|
||||
RequestBackChannels: true,
|
||||
Transport: &transport,
|
||||
Protocol: &protocol,
|
||||
TLSConfig: tlsConfig,
|
||||
}
|
||||
// parse URL
|
||||
u, err := base.ParseURL(g.Url)
|
||||
@@ -606,7 +653,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 +729,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 +748,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 +784,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 +843,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 +862,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 +1088,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 +1107,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
|
||||
|
||||
79
machinery/src/capture/gortsplib_test.go
Normal file
79
machinery/src/capture/gortsplib_test.go
Normal file
@@ -0,0 +1,79 @@
|
||||
package capture
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/pem"
|
||||
"net/http/httptest"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestRTSPSTLSConfig(t *testing.T) {
|
||||
t.Run("verifies certificates by default", func(t *testing.T) {
|
||||
t.Setenv(rtspsCAFileEnv, "")
|
||||
t.Setenv(rtspsInsecureEnv, "")
|
||||
|
||||
got, err := rtspsTLSConfig()
|
||||
if err != nil {
|
||||
t.Fatalf("rtspsTLSConfig() error = %v", err)
|
||||
}
|
||||
if got != nil {
|
||||
t.Fatalf("rtspsTLSConfig() = %#v, want nil", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("allows explicit insecure mode", func(t *testing.T) {
|
||||
t.Setenv(rtspsCAFileEnv, "/missing/ignored-in-insecure-mode.pem")
|
||||
t.Setenv(rtspsInsecureEnv, "true")
|
||||
|
||||
got, err := rtspsTLSConfig()
|
||||
if err != nil {
|
||||
t.Fatalf("rtspsTLSConfig() error = %v", err)
|
||||
}
|
||||
if got == nil || !got.InsecureSkipVerify {
|
||||
t.Fatalf("rtspsTLSConfig() = %#v, want InsecureSkipVerify enabled", got)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("adds a camera CA to system roots", func(t *testing.T) {
|
||||
t.Setenv(rtspsInsecureEnv, "")
|
||||
server := httptest.NewTLSServer(nil)
|
||||
defer server.Close()
|
||||
|
||||
certificate := server.Certificate()
|
||||
caFile := filepath.Join(t.TempDir(), "camera-ca.pem")
|
||||
caPEM := pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: certificate.Raw})
|
||||
if err := os.WriteFile(caFile, caPEM, 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Setenv(rtspsCAFileEnv, caFile)
|
||||
|
||||
got, err := rtspsTLSConfig()
|
||||
if err != nil {
|
||||
t.Fatalf("rtspsTLSConfig() error = %v", err)
|
||||
}
|
||||
if got == nil || got.RootCAs == nil {
|
||||
t.Fatalf("rtspsTLSConfig() = %#v, want custom RootCAs", got)
|
||||
}
|
||||
for _, subject := range got.RootCAs.Subjects() {
|
||||
if bytes.Equal(subject, certificate.RawSubject) {
|
||||
return
|
||||
}
|
||||
}
|
||||
t.Fatal("camera CA was not added to RootCAs")
|
||||
})
|
||||
|
||||
t.Run("rejects an invalid camera CA file", func(t *testing.T) {
|
||||
t.Setenv(rtspsInsecureEnv, "")
|
||||
caFile := filepath.Join(t.TempDir(), "camera-ca.pem")
|
||||
if err := os.WriteFile(caFile, []byte("not a certificate"), 0o600); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
t.Setenv(rtspsCAFileEnv, caFile)
|
||||
|
||||
if _, err := rtspsTLSConfig(); err == nil {
|
||||
t.Fatal("rtspsTLSConfig() error = nil, want invalid CA error")
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -229,6 +229,54 @@ func rawJSONOrEmptyArray(b []byte) json.RawMessage {
|
||||
return json.RawMessage(b)
|
||||
}
|
||||
|
||||
const heartbeatResponseBodyLogLimit = 4 * 1024
|
||||
|
||||
func readHeartbeatResponseBody(response *http.Response) (string, bool, error) {
|
||||
if response == nil || response.Body == nil {
|
||||
return "", false, nil
|
||||
}
|
||||
defer response.Body.Close()
|
||||
|
||||
body, err := io.ReadAll(io.LimitReader(response.Body, heartbeatResponseBodyLogLimit+1))
|
||||
if err != nil {
|
||||
return "", false, err
|
||||
}
|
||||
|
||||
truncated := len(body) > heartbeatResponseBodyLogLimit
|
||||
if truncated {
|
||||
body = body[:heartbeatResponseBodyLogLimit]
|
||||
}
|
||||
|
||||
return strings.TrimSpace(string(body)), truncated, nil
|
||||
}
|
||||
|
||||
func formatHeartbeatFailureLog(response *http.Response, requestErr error, responseBody string, responseBodyTruncated bool, responseBodyErr error, elapsed time.Duration) string {
|
||||
details := make([]string, 0, 6)
|
||||
if response != nil {
|
||||
details = append(details, "status_code="+strconv.Itoa(response.StatusCode))
|
||||
if response.Status != "" {
|
||||
details = append(details, "status="+strconv.Quote(response.Status))
|
||||
}
|
||||
} else {
|
||||
details = append(details, "status_code=none")
|
||||
}
|
||||
details = append(details, "duration="+elapsed.Round(time.Millisecond).String())
|
||||
if requestErr != nil {
|
||||
details = append(details, "request_error="+strconv.Quote(requestErr.Error()))
|
||||
}
|
||||
if responseBody != "" {
|
||||
details = append(details, "response_body="+strconv.Quote(responseBody))
|
||||
}
|
||||
if responseBodyTruncated {
|
||||
details = append(details, "response_body_truncated=true")
|
||||
}
|
||||
if responseBodyErr != nil {
|
||||
details = append(details, "response_body_error="+strconv.Quote(responseBodyErr.Error()))
|
||||
}
|
||||
|
||||
return "cloud.HandleHeartBeat(): heartbeat request to Kerberos Hub failed: " + strings.Join(details, ", ")
|
||||
}
|
||||
|
||||
func HandleHeartBeat(configuration *models.Configuration, communication *models.Communication, uptimeStart time.Time) {
|
||||
log.Log.Debug("cloud.HandleHeartBeat(): started")
|
||||
|
||||
@@ -644,20 +692,26 @@ loop:
|
||||
|
||||
var jsonStr = []byte(object)
|
||||
buffy := bytes.NewBuffer(jsonStr)
|
||||
req, _ := http.NewRequest("POST", hubURI, buffy)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
resp, err := client.Do(req)
|
||||
if resp != nil {
|
||||
resp.Body.Close()
|
||||
requestStarted := time.Now()
|
||||
req, requestErr := http.NewRequest("POST", hubURI, buffy)
|
||||
var resp *http.Response
|
||||
if requestErr == nil {
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
resp, requestErr = client.Do(req)
|
||||
}
|
||||
if err == nil && resp.StatusCode == 200 {
|
||||
if requestErr == nil && resp != nil && resp.StatusCode == http.StatusOK {
|
||||
if resp.Body != nil {
|
||||
_, _ = io.Copy(io.Discard, resp.Body)
|
||||
resp.Body.Close()
|
||||
}
|
||||
communication.CloudTimestamp.Store(time.Now().Unix())
|
||||
log.Log.Info("cloud.HandleHeartBeat(): (200) Heartbeat received by Kerberos Hub.")
|
||||
} else {
|
||||
responseBody, responseBodyTruncated, responseBodyErr := readHeartbeatResponseBody(resp)
|
||||
if communication.CloudTimestamp != nil && communication.CloudTimestamp.Load() != nil {
|
||||
communication.CloudTimestamp.Store(int64(0))
|
||||
}
|
||||
log.Log.Error("cloud.HandleHeartBeat(): (400) Something went wrong while sending to Kerberos Hub.")
|
||||
log.Log.Error(formatHeartbeatFailureLog(resp, requestErr, responseBody, responseBodyTruncated, responseBodyErr, time.Since(requestStarted)))
|
||||
}
|
||||
} else {
|
||||
log.Log.Error("cloud.HandleHeartBeat(): Disabled as we do not have a public key defined.")
|
||||
|
||||
51
machinery/src/cloud/cloud_test.go
Normal file
51
machinery/src/cloud/cloud_test.go
Normal file
@@ -0,0 +1,51 @@
|
||||
package cloud
|
||||
|
||||
import (
|
||||
"io"
|
||||
"net/http"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestHeartbeatFailureLogIncludesHubResponse(t *testing.T) {
|
||||
response := &http.Response{
|
||||
StatusCode: http.StatusBadRequest,
|
||||
Status: "400 Bad Request",
|
||||
Body: io.NopCloser(strings.NewReader(`{"error":"invalid heartbeat"}`)),
|
||||
}
|
||||
|
||||
responseBody, truncated, err := readHeartbeatResponseBody(response)
|
||||
if err != nil {
|
||||
t.Fatalf("readHeartbeatResponseBody() error = %v", err)
|
||||
}
|
||||
message := formatHeartbeatFailureLog(response, nil, responseBody, truncated, nil, 125*time.Millisecond)
|
||||
|
||||
for _, expected := range []string{
|
||||
"status_code=400",
|
||||
`status="400 Bad Request"`,
|
||||
"duration=125ms",
|
||||
`response_body="{\"error\":\"invalid heartbeat\"}"`,
|
||||
} {
|
||||
if !strings.Contains(message, expected) {
|
||||
t.Errorf("formatHeartbeatFailureLog() = %q, want it to contain %q", message, expected)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestReadHeartbeatResponseBodyTruncatesLargeBody(t *testing.T) {
|
||||
response := &http.Response{
|
||||
Body: io.NopCloser(strings.NewReader(strings.Repeat("x", heartbeatResponseBodyLogLimit+1))),
|
||||
}
|
||||
|
||||
body, truncated, err := readHeartbeatResponseBody(response)
|
||||
if err != nil {
|
||||
t.Fatalf("readHeartbeatResponseBody() error = %v", err)
|
||||
}
|
||||
if !truncated {
|
||||
t.Fatal("readHeartbeatResponseBody() truncated = false, want true")
|
||||
}
|
||||
if len(body) != heartbeatResponseBodyLogLimit {
|
||||
t.Fatalf("len(body) = %d, want %d", len(body), heartbeatResponseBodyLogLimit)
|
||||
}
|
||||
}
|
||||
@@ -35,10 +35,15 @@ func UploadKerberosVault(configuration *models.Configuration, fileName string) (
|
||||
// This can happen when the file was already removed (e.g. cleanup, or an
|
||||
// earlier successful upload). Skip it so the watcher drops the marker
|
||||
// instead of retrying indefinitely.
|
||||
if _, err := os.Stat("data/recordings/" + fileName); err != nil {
|
||||
info, err := os.Stat("data/recordings/" + fileName)
|
||||
if err != nil {
|
||||
log.Log.Info("UploadKerberosVault: skipping " + fileName + ", file doesn't exist anymore")
|
||||
return false, false, nil
|
||||
}
|
||||
if info.Size() == 0 {
|
||||
log.Log.Warning("UploadKerberosVault: skipping " + fileName + ", recording is empty")
|
||||
return false, false, nil
|
||||
}
|
||||
|
||||
// timestamp_microseconds_instanceName_regionCoordinates_numberOfChanges_token
|
||||
// 1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4
|
||||
|
||||
@@ -151,6 +151,13 @@ func publishLiveStreamMoQ(ctx context.Context, config liveMoQConfig) error {
|
||||
}
|
||||
defer client.Close()
|
||||
|
||||
sessionCtx, cancelSessionWatch := context.WithCancel(ctx)
|
||||
defer cancelSessionWatch()
|
||||
sessionClosed := make(chan error, 1)
|
||||
go func() {
|
||||
sessionClosed <- client.Session().Closed(sessionCtx)
|
||||
}()
|
||||
|
||||
broadcast, err := client.CreateBroadcast(config.broadcast)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create broadcast: %w", err)
|
||||
@@ -180,6 +187,12 @@ func publishLiveStreamMoQ(ctx context.Context, config liveMoQConfig) error {
|
||||
var lastDuplicateKeyframeWarning time.Time
|
||||
idle := false
|
||||
for {
|
||||
select {
|
||||
case err := <-sessionClosed:
|
||||
return fmt.Errorf("relay session closed: %w", err)
|
||||
default:
|
||||
}
|
||||
|
||||
packet, err := cursor.ReadPacket()
|
||||
if err != nil {
|
||||
return fmt.Errorf("read packet: %w", err)
|
||||
|
||||
@@ -264,6 +264,36 @@ func testVault(uri string) models.KStorage {
|
||||
}
|
||||
}
|
||||
|
||||
func TestUploadKerberosVaultSkipsEmptyRecording(t *testing.T) {
|
||||
fileName := "1787015373_3-654_office-camera17_0-0-0-0_-1_1960.mp4"
|
||||
withRecording(t, fileName, nil)
|
||||
|
||||
requestCount := 0
|
||||
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
requestCount++
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
}))
|
||||
defer server.Close()
|
||||
|
||||
vault := testVault(server.URL)
|
||||
configuration := &models.Configuration{Config: models.Config{
|
||||
Key: "device-key",
|
||||
KStorage: &vault,
|
||||
KStorageSecondary: &models.KStorage{},
|
||||
}}
|
||||
|
||||
uploaded, configured, err := UploadKerberosVault(configuration, fileName)
|
||||
if err != nil {
|
||||
t.Fatalf("UploadKerberosVault() error = %v", err)
|
||||
}
|
||||
if uploaded || configured {
|
||||
t.Fatalf("UploadKerberosVault() uploaded/configured = %v/%v, want false/false", uploaded, configured)
|
||||
}
|
||||
if requestCount != 0 {
|
||||
t.Fatalf("Vault received %d requests, want 0", requestCount)
|
||||
}
|
||||
}
|
||||
|
||||
func TestUploadVaultResumable_HappyPath(t *testing.T) {
|
||||
srv := newFakeTus()
|
||||
ts := httptest.NewServer(srv)
|
||||
|
||||
@@ -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) {
|
||||
|
||||
209
machinery/src/components/backchannel_test.go
Normal file
209
machinery/src/components/backchannel_test.go
Normal 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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
39
machinery/src/routers/mqtt/main_test.go
Normal file
39
machinery/src/routers/mqtt/main_test.go
Normal 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")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user