Compare commits

..

29 Commits

Author SHA1 Message Date
Cédric Verstraeten
58a79f8278 Merge pull request #294 from kerberos-io/feature/update-readme-turn-info
feature/update-readme-turn-info
2026-06-23 09:27:11 +02:00
Cédric Verstraeten
422279985f Update STUN and TURN server URIs in README 2026-06-23 09:21:27 +02:00
Cédric Verstraeten
13c84a0f36 Merge pull request #292 from kerberos-io/feature/update-turn-uri
Change STUN and TURN URIs in config.json
2026-06-22 21:21:30 +02:00
Cédric Verstraeten
cb6bbe1609 Change STUN and TURN URIs in config.json
Updated STUN and TURN URIs for improved connectivity.
2026-06-22 21:12:52 +02:00
Cédric Verstraeten
476207c1bf Merge pull request #290 from kerberos-io/feature/add-hls-live-streaming
feature/add-hls-live-streaming
2026-06-16 10:15:03 +02:00
Cédric Verstraeten
fcd8ef8ff4 Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-16 10:09:59 +02:00
Cédric Verstraeten
645b6aa0be Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-16 10:09:49 +02:00
Cédric Verstraeten
67ee78dab5 Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-16 10:09:15 +02:00
Cédric Verstraeten
5936c6eaae Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-16 10:09:02 +02:00
Cédric Verstraeten
dafcd06696 Add live HLS streaming support
Introduce live HLS streaming pipeline and wire it into the agent.

- Add cloud/livehls: HandleLiveStreamHLS reads packets, gates session lifetime by viewer keepalives, starts sessions lazily on keyframes and announces ready state over MQTT.
- Add livehls publisher and session (cloud/livehls/{publisher,session}.go) to upload init + media CMAF segments to hub-api using a header-based ingest contract; includes redirect-credential stripping and init-refresh logic.
- Add video/livehls.go: LiveSegmenter converts Annex B video into one init (ftyp+moov) and self-contained CMAF media segments (styp+moof+mdat) keyed by sequence/duration.
- Add tests for publisher and session behavior (cloud/livehls/publisher_test.go, video/livehls_test.go).
- Wire components and routing: add Communication.HandleLiveHLS channel and start HLS handler in RunAgent; add RequestHLSStreamPayload and HandleRequestHLSStream in MQTT router to treat HLS requests as viewer keepalives.

This enables short‑latency HLS streaming (CMAF segments uploaded fire‑and‑forget) with viewer keepalive semantics and minimal changes to the control plane.
2026-06-16 09:33:41 +02:00
Cédric Verstraeten
02d60c71e4 Merge pull request #289 from kerberos-io/fix/dynamic-gopsizes
fix/dynamic-gopsizes
2026-06-15 15:00:47 +02:00
Cédric Verstraeten
52647d7f1d Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-15 14:57:07 +02:00
Cédric Verstraeten
e1fa7d9d7e Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-15 14:56:46 +02:00
Cédric Verstraeten
06e2694763 Improve seam detection and add analysis tools
Refine loop/restart (seam) detection and enhance the mp4 analysis tooling.

- mp4: Replace previous previous-interval-based seam heuristic with a safer approach that (1) tracks the running minimum keyframe interval (MinKeyframeGapMs) as the reference cadence and (2) requires the buffered GOP to be genuinely truncated before dropping it. This avoids false positives on variable-GOP (smart-codec) cameras. Added fields, logging updates, and helper methods: bufferedVideoCount, expectedGopFrames, bufferedVideoFrameDuration. LastKeyframeGapMs is now diagnostic only.

- cmd/mp4analyze: Add -from/-to flags and auto-select a detailed inspection window centred on the largest keyframe gap. Limit per-sample printing to that window (and anomalies), fix sync-sample bit detection, pass window to sliceHeaders, and add a compact SUMMARY health report with medians and checks. Added inspectWindow and median helpers.

- tests: Add mp4_variablegop_test.go to verify variable-GOP streams keep healthy short GOPs and that no frames are dropped by the improved seam logic.

These changes prevent healthy GOPs from being discarded on normal short GOPs that follow long static GOPs and add better diagnostics for debugging artifacts.
2026-06-15 14:44:05 +02:00
Cédric Verstraeten
c0971ca3b2 Merge pull request #288 from kerberos-io/feature/resumable-uploads-to-hub
feature/resumable-uploads-to-hub
2026-06-14 15:21:13 +02:00
Cédric Verstraeten
1a788ebe6c Potential fix for pull request finding 'Writable file handle closed without error handling'
Co-authored-by: Copilot Autofix powered by AI <223894421+github-code-quality[bot]@users.noreply.github.com>
2026-06-13 22:41:50 +02:00
Cédric Verstraeten
a1b4026b4b Remove duplicate uppercase filename entries from git index
Case-only renames had committed both Camera.go and camera.go (identical
blobs) under core.ignorecase=true, causing 'case-insensitive file name
collision' in go build on CI (case-sensitive checkout). Drop the 15 stale
uppercase index entries; the lowercase files are unchanged.
2026-06-13 20:34:41 +00:00
Cédric Verstraeten
9bc9825bb1 Normalize filenames; add tus hub resumable tests
Rename many Go source files to lower_snake_case (e.g. RTSPClient.go -> rtsp_client.go, Server.go -> server.go, etc.) to follow project naming conventions. Enhance machinery/src/cloud/tus_client_test.go: add encoding/base64 import, record incoming requests (recordedRequest + requests slice), provide requestsForMethod helper, add test helpers (testHubConfig, decodeTusMetadata) and two new tests (TestUploadHubResumable_HappyPath and TestUploadHubResumable_Unsupported) that validate hub resumable upload behavior and per-method auth/metadata. Update swag.sh to point to the renamed server.go entry file.
2026-06-13 22:23:53 +02:00
Cédric Verstraeten
e9d2afa228 Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-13 12:23:49 +02:00
Cédric Verstraeten
4b0e0eae9c Trim HubURI when building tus upload URL
Construct the tus upload base URL by trimming any trailing slash from config.HubURI before appending tusUploadPath. This prevents double slashes in the resulting URL when concatenating the base URI and the upload path, avoiding potential request/endpoint errors.
2026-06-13 09:55:49 +02:00
Cédric Verstraeten
e0204e1949 Add tus resumable uploads and Hub support
Prefer and support tus resumable uploads for Kerberos Vault/Hub and fall back to legacy single-POST when not available. Extract runTusUpload and a tusHeaderFunc to share the create/head/patch/terminate state machine between direct vault and hub-proxied uploads. Update tusCreate/tusHead/tusPatch/tusTerminate to accept header injection, implement uploadVaultResumable and uploadHubResumable wrappers, and add setHubTusHeaders. Also update UploadKerberosHub to attempt resumable uploads first and log fallback behavior.
2026-06-13 09:55:15 +02:00
Cédric Verstraeten
3c2a0ce0cf Merge pull request #287 from kerberos-io/feature/add-tus-progress-for-resumable-uploads
feature/add-tus-progress-for-resumable-uploads
2026-06-12 13:33:47 +02:00
Cédric Verstraeten
a5def2ccd8 Log tus resumable upload progress
Add periodic progress logging for tus resumable uploads. Introduces tusProgressBucketPercent (10%) and helper functions tusProgressBucket and logTusUploadProgress to bucket progress into 10% increments, cap at 100%, and avoid repeated logs. Initializes loggedProgressBucket in uploadVaultResumable and calls logTusUploadProgress after each successful PATCH so upload progress is reported concisely (percent and byte offsets) without excessive noise.
2026-06-12 13:29:32 +02:00
Cédric Verstraeten
6ede3c3add Merge pull request #286 from kerberos-io/feature/resumable-uploads-tusd
feature/resumable-uploads-tusd
2026-06-12 11:58:19 +02:00
Cédric Verstraeten
d5de6ae271 Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-12 11:55:02 +02:00
Cédric Verstraeten
2035deaa31 Potential fix for pull request finding
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
2026-06-12 11:54:16 +02:00
Cédric Verstraeten
5973ba025d Add configurable tus chunking and stability fixes
Enable configurable chunked tus uploads and related robustness changes.

- Add AGENT_TUS_CHUNK_SIZE_BYTES env (default 1 MiB, 0 disables chunking) and docs in machinery/.env; ignore machinery/go.work files and set GOWORK=off in VSCode launch to avoid go.work during debugging.
- Implement tusChunkSize() and update uploadVaultResumable to send PATCHes in configurable chunk sizes, checkpoint progress after each chunk, and handle partial failures by refreshing retry budget when progress occurs.
- Change tusPatch to accept an explicit length and return the advanced offset when a PATCH is fully accepted.
- Skip uploads when the recording file no longer exists to avoid infinite retries.
- Add tests exercising chunked uploads, chunking-disabled behavior, and tusChunkSize parsing; extend fakeTus test server to record patch sizes.
- Reduce noisy info logs to debug in AAC transcoder and WebRTC audio processing.

These changes improve resumable upload reliability, allow tuning for proxy/load-balancer limits, and reduce log spam during normal operation.
2026-06-11 23:03:05 +02:00
Cédric Verstraeten
52a54fbae1 Add tus resumable upload client and MP4 analyzer
Introduce a standalone mp4analyze CLI for inspecting fragmented MP4s and add a full-featured tus resumable upload client used by Kerberos Vault uploads.

Changes:
- Add machinery/cmd/mp4analyze/main.go: CLI tool to analyze MP4 structure, fragments, keyframes, NALs and SPS/PPS differences.
- Add machinery/src/cloud/tus_client.go: implements tus 1.0.0 client with create/head/patch/delete, sidecar resume state, metadata encoding, Location resolution, backoff, and helpers (newVaultHTTPClient, setVaultTusHeaders). Resumable uploads are enabled by default and can be disabled via AGENT_DISABLE_RESUMABLE_UPLOAD. Sidecar files are stored under data/tus.
- Add machinery/src/cloud/tus_client_test.go: in-memory fake tus server and unit tests exercising happy path, unsupported servers, finalize-retry and resume-from-sidecar behavior, and helpers.
- Refactor machinery/src/cloud/kerberos_vault.go: replace inlined POST upload logic with sendToVault which attempts resumable uploads first and falls back to legacy single-POST (uploadVaultLegacy). Improve retry semantics so retry counters only advance on definitive vault responses, centralize header setup, and use new HTTP client builder honoring AGENT_TLS_INSECURE.

The change preserves backward compatibility with vaults that do not support tus by transparently falling back to the legacy upload path. Tests cover core tus behaviors and resume semantics.
2026-06-11 21:32:43 +02:00
Cédric Verstraeten
5f828262eb Merge pull request #285 from kerberos-io/fix/githubaction-version-passthrough
fix/githubaction-version-passthrough
2026-06-11 20:00:19 +02:00
40 changed files with 3941 additions and 144 deletions

2
.gitignore vendored
View File

@@ -14,5 +14,7 @@ machinery/test*
machinery/init-dev.sh
machinery/.env.local
machinery/vendor
machinery/go.work
machinery/go.work.sum
deployments/docker/private-docker-compose.yaml
video.mp4

3
.vscode/launch.json vendored
View File

@@ -18,6 +18,9 @@
],
"envFile": "${workspaceFolder}/machinery/.env.local",
"buildFlags": "--tags dynamic",
"env": {
"GOWORK": "off"
},
},
{
"name": "Launch React",

View File

@@ -231,9 +231,9 @@ Next to attaching the configuration file, it is also possible to override the co
| `AGENT_MQTT_PASSWORD` | Password of the MQTT broker. | "" |
| `AGENT_REALTIME_PROCESSING` | If `AGENT_REALTIME_PROCESSING` set to `true`, the agent will send key frames to the topic | "" |
| `AGENT_REALTIME_PROCESSING_TOPIC` | The topic to which keyframes will be sent in base64 encoded format. | "" |
| `AGENT_STUN_URI` | When using WebRTC, you'll need to provide a STUN server. | "stun:turn.kerberos.io:8443" |
| `AGENT_STUN_URI` | When using WebRTC, you'll need to provide a STUN server. | "stun:turn-fra1.kerberos.io:3478"|
| `AGENT_FORCE_TURN` | Force using a TURN server, by generating relay candidates only. | "false" |
| `AGENT_TURN_URI` | When using WebRTC, you'll need to provide a TURN server. | "turn:turn.kerberos.io:8443" |
| `AGENT_TURN_URI` | When using WebRTC, you'll need to provide a TURN server. | "turn:turn-fra1.kerberos.io:348"|
| `AGENT_TURN_USERNAME` | TURN username used for WebRTC. | "username1" |
| `AGENT_TURN_PASSWORD` | TURN password used for WebRTC. | "password1" |
| `AGENT_CLOUD` | Store recordings in Kerberos Hub (s3), Kerberos Vault (kstorage), or Dropbox (dropbox). | "s3" |

View File

@@ -27,5 +27,12 @@ AGENT_KERBEROSVAULT_SECONDARY_DIRECTORY=
AGENT_KERBEROSVAULT_SECONDARY_ACCESS_KEY=
AGENT_KERBEROSVAULT_SECONDARY_SECRET_KEY=
# Resumable (tus) uploads to Kerberos Vault are enabled by default.
# Set to true to fall back to the legacy single-shot POST /storage upload.
#AGENT_DISABLE_RESUMABLE_UPLOAD=true
# Bytes sent per PATCH request (default 1 MiB = 1048576). 0 disables chunking
# and sends the whole file in a single PATCH.
AGENT_TUS_CHUNK_SIZE_BYTES=1048576
# Open telemetry tracing endpoint
OTEL_EXPORTER_OTLP_ENDPOINT=

View File

@@ -0,0 +1,636 @@
package main
import (
"flag"
"fmt"
"os"
"sort"
"github.com/Eyevinn/mp4ff/avc"
mp4ff "github.com/Eyevinn/mp4ff/mp4"
)
func main() {
fromFlag := flag.Int64("from", -1, "start of the detailed inspection window (track timescale units); default auto-detects the largest keyframe gap")
toFlag := flag.Int64("to", -1, "end of the detailed inspection window (track timescale units); default auto-detected")
flag.Parse()
if flag.NArg() < 1 {
fmt.Println("usage: mp4analyze [-from N] [-to N] <file.mp4>")
os.Exit(1)
}
f, err := os.Open(flag.Arg(0))
if err != nil {
panic(err)
}
defer f.Close()
parsed, err := mp4ff.DecodeFile(f)
if err != nil {
panic(err)
}
// Movie-level info
if parsed.Init != nil && parsed.Init.Moov != nil {
moov := parsed.Init.Moov
fmt.Printf("ftyp/moov present. timescale(mvhd)=%d duration(mvhd)=%d\n",
moov.Mvhd.Timescale, moov.Mvhd.Duration)
for _, trak := range moov.Traks {
ts := trak.Mdia.Mdhd.Timescale
fmt.Printf(" trak id=%d handler=%s mdhd.timescale=%d mdhd.duration=%d\n",
trak.Tkhd.TrackID, trak.Mdia.Hdlr.HandlerType, ts, trak.Mdia.Mdhd.Duration)
}
} else {
fmt.Println("no Init/Moov (pure fragmented stream?)")
}
// sidx vs actual segment layout. MSE players use sidx to map presentation
// time -> byte ranges; if sidx references disagree with the real segment
// sizes/durations (e.g. after an early/short flush) the player fetches the
// wrong bytes and fails to decode — a failure that "heals" on seek.
fmt.Println("=== sidx references vs actual segments ===")
var sidxRefs []mp4ff.SidxRef
for _, c := range parsed.Children {
if s, ok := c.(*mp4ff.SidxBox); ok {
fmt.Printf(" sidx: timescale=%d earliestPresTime=%d firstOffset=%d refCount=%d anchor(after sidx)=%d\n",
s.Timescale, s.EarliestPresentationTime, s.FirstOffset, len(s.SidxRefs), s.AnchorPoint)
sidxRefs = s.SidxRefs
}
}
// Actual segment sizes (styp+moof+mdat) and fragment durations.
type segInfo struct {
size uint64
dur uint64
}
var actual []segInfo
for _, seg := range parsed.Segments {
var sz uint64
if seg.Styp != nil {
sz += seg.Styp.Size()
}
if seg.Sidx != nil {
sz += seg.Sidx.Size()
}
var dur uint64
for _, fr := range seg.Fragments {
sz += fr.Moof.Size()
if fr.Mdat != nil {
sz += fr.Mdat.Size()
}
for _, traf := range fr.Moof.Trafs {
if traf.Tfhd.TrackID != 1 {
continue
}
for _, trun := range traf.Truns {
for _, s := range trun.Samples {
dur += uint64(s.Dur)
}
}
}
}
actual = append(actual, segInfo{size: sz, dur: dur})
}
for i := range actual {
refStr := "(no sidx ref)"
if i < len(sidxRefs) {
r := sidxRefs[i]
mark := ""
if uint64(r.ReferencedSize) != actual[i].size {
mark += fmt.Sprintf(" SIZE MISMATCH actual=%d", actual[i].size)
}
if uint64(r.SubSegmentDuration) != actual[i].dur {
mark += fmt.Sprintf(" DUR MISMATCH actual=%d", actual[i].dur)
}
refStr = fmt.Sprintf("sidx.size=%d sidx.dur=%d type=%d sap=%d/%d%s",
r.ReferencedSize, r.SubSegmentDuration, r.ReferenceType, r.StartsWithSAP, r.SAPType, mark)
}
fmt.Printf(" seg%02d actual.size=%d actual.dur=%d | %s\n", i, actual[i].size, actual[i].dur, refStr)
}
fmt.Println("=== fragments ===")
fragIdx := 0
var allKeyGlobal []uint64 // global keyframe decode times (track timescale units)
var prevTfdtEnd = map[uint32]uint64{}
for si, seg := range parsed.Segments {
for _, fr := range seg.Fragments {
for _, traf := range fr.Moof.Trafs {
tid := traf.Tfhd.TrackID
tfdt := traf.Tfdt.BaseMediaDecodeTime()
offset := uint64(0)
var keys []uint64 // keyframe offset-from-tfdt
var durs []uint64
zeroDur := 0
nSamples := 0
for _, trun := range traf.Truns {
for _, s := range trun.Samples {
nSamples++
if (s.Flags>>24)&0x03 == 0x02 { // sample_depends_on==2 => IDR/sync
keys = append(keys, offset)
if tid == 1 {
allKeyGlobal = append(allKeyGlobal, tfdt+offset)
}
}
if s.Dur == 0 {
zeroDur++
}
durs = append(durs, uint64(s.Dur))
offset += uint64(s.Dur)
}
}
cont := ""
if pe, ok := prevTfdtEnd[tid]; ok {
if tfdt != pe {
cont = fmt.Sprintf(" <-- tfdt GAP/JUMP prev_end=%d delta=%d", pe, int64(tfdt)-int64(pe))
}
}
prevTfdtEnd[tid] = tfdt + offset
if tid == 1 {
// in-fragment keyframe gaps
var gaps []int64
for i := 1; i < len(keys); i++ {
gaps = append(gaps, int64(keys[i])-int64(keys[i-1]))
}
fmt.Printf("seg%d frag%d trk%d tfdt=%d dur=%d nSamp=%d zeroDur=%d keys=%v inFragKeyGaps=%v%s\n",
si, fragIdx, tid, tfdt, offset, nSamples, zeroDur, keys, gaps, cont)
}
}
fragIdx++
}
}
fmt.Println("=== global video keyframe decode times & gaps ===")
for i, k := range allKeyGlobal {
gap := int64(0)
if i > 0 {
gap = int64(k) - int64(allKeyGlobal[i-1])
}
seam := ""
if i > 1 {
prevGap := int64(allKeyGlobal[i-1]) - int64(allKeyGlobal[i-2])
if gap > 0 && prevGap > 0 && gap*2 < prevGap {
seam = fmt.Sprintf(" <== SEAM? gap=%d < prevGap/2=%d", gap, prevGap/2)
}
}
fmt.Printf(" kf#%02d dt=%d gap=%d%s\n", i, k, gap, seam)
}
// Choose the detailed-inspection window. By default centre it on the largest
// keyframe gap (the most likely artifact location); -from/-to override.
winLo, winHi := inspectWindow(allKeyGlobal, *fromFlag, *toFlag)
fmt.Printf("=== detailed inspection window: dts %d..%d ===\n", winLo, winHi)
// Full sample timeline: DTS, CTS (=DTS+cto), composition offset, NAL types,
// to detect PTS non-monotonicity / gaps / param-set changes at the seam.
fmt.Println("=== per-sample timeline (full) — checking PTS monotonicity & nal types ===")
var trex *mp4ff.TrexBox
if parsed.Init != nil && parsed.Init.Moov != nil && parsed.Init.Moov.Mvex != nil {
for _, t := range parsed.Init.Moov.Mvex.Trexs {
if t.TrackID == 1 {
trex = t
}
}
}
var lastCTS int64 = -1
var lastDTS int64 = -1
sampIdx := 0
fragIdx = 0
for _, seg := range parsed.Segments {
for _, fr := range seg.Fragments {
fs, err := fr.GetFullSamples(trex)
if err != nil {
fmt.Printf(" frag%d GetFullSamples err: %v\n", fragIdx, err)
fragIdx++
continue
}
for _, s := range fs {
dts := int64(s.DecodeTime)
cts := dts + int64(s.CompositionTimeOffset)
nals := nalTypes(s.Data)
anomaly := ""
if lastCTS >= 0 && cts < lastCTS {
anomaly += fmt.Sprintf(" <== CTS BACKWARDS (prev=%d)", lastCTS)
}
if lastDTS >= 0 && dts < lastDTS {
anomaly += fmt.Sprintf(" <== DTS BACKWARDS (prev=%d)", lastDTS)
}
// sample_is_non_sync_sample is bit 16 (0x00010000); a sync sample
// has it clear and sample_depends_on==2 (i.e. an I-frame).
isSync := s.Flags&0x00010000 == 0 && (s.Flags>>24)&0x03 == 0x02
// Only print inside the inspection window and any anomalies, to keep output small.
near := dts >= winLo && dts <= winHi
if near || anomaly != "" {
fmt.Printf(" s%04d frag%d dts=%d cts=%d cto=%d dur=%d size=%d sync=%v nal=%v%s\n",
sampIdx, fragIdx, dts, cts, s.CompositionTimeOffset, s.Dur, len(s.Data), isSync, nals, anomaly)
}
lastCTS = cts
lastDTS = dts
sampIdx++
}
fragIdx++
}
}
// Compare parameter sets: avcC (in moov) vs inline SPS/PPS at every IDR.
// A looping source that restarts may re-emit SPS/PPS that differ from the
// ones the player configured its decoder with from avcC — a classic cause
// of a freeze that "heals" when you seek past the seam.
fmt.Println("=== parameter set comparison (avcC vs inline IDR) ===")
var avccSPS, avccPPS [][]byte
if parsed.Init != nil && parsed.Init.Moov != nil {
for _, trak := range parsed.Init.Moov.Traks {
if trak.Mdia == nil || trak.Mdia.Minf == nil || trak.Mdia.Minf.Stbl == nil {
continue
}
stsd := trak.Mdia.Minf.Stbl.Stsd
if stsd == nil || stsd.AvcX == nil || stsd.AvcX.AvcC == nil {
continue
}
avccSPS = stsd.AvcX.AvcC.SPSnalus
avccPPS = stsd.AvcX.AvcC.PPSnalus
}
}
for i, s := range avccSPS {
fmt.Printf(" avcC SPS[%d] = %x\n", i, s)
}
for i, p := range avccPPS {
fmt.Printf(" avcC PPS[%d] = %x\n", i, p)
}
fragIdx = 0
sampIdx = 0
var baseSPS, basePPS []byte
if len(avccSPS) > 0 {
baseSPS = avccSPS[0]
}
if len(avccPPS) > 0 {
basePPS = avccPPS[0]
}
for _, seg := range parsed.Segments {
for _, fr := range seg.Fragments {
fs, err := fr.GetFullSamples(trex)
if err != nil {
fragIdx++
continue
}
for _, s := range fs {
spsList := nalsByType(s.Data, 7)
ppsList := nalsByType(s.Data, 8)
if len(spsList) > 0 || len(ppsList) > 0 {
dts := int64(s.DecodeTime)
note := ""
if len(spsList) > 0 {
if baseSPS == nil {
baseSPS = spsList[0]
} else if !bytesEqual(baseSPS, spsList[0]) {
note += " <== SPS CHANGED vs base/avcC"
}
}
if len(ppsList) > 0 {
if basePPS == nil {
basePPS = ppsList[0]
} else if !bytesEqual(basePPS, ppsList[0]) {
note += " <== PPS CHANGED vs base/avcC"
}
}
var spsHex, ppsHex string
if len(spsList) > 0 {
spsHex = fmt.Sprintf("%x", spsList[0])
}
if len(ppsList) > 0 {
ppsHex = fmt.Sprintf("%x", ppsList[0])
}
fmt.Printf(" IDR s%04d frag%d dts=%d SPS=%s PPS=%s%s\n",
sampIdx, fragIdx, dts, spsHex, ppsHex, note)
}
sampIdx++
}
fragIdx++
}
}
sliceHeaders(parsed, trex, winLo, winHi)
summary(parsed, trex)
}
func sliceHeaders(parsed *mp4ff.File, trex *mp4ff.TrexBox, winLo, winHi int64) {
// Build SPS/PPS maps from avcC.
spsMap := map[uint32]*avc.SPS{}
ppsMap := map[uint32]*avc.PPS{}
if parsed.Init != nil && parsed.Init.Moov != nil {
for _, trak := range parsed.Init.Moov.Traks {
if trak.Mdia == nil || trak.Mdia.Minf == nil || trak.Mdia.Minf.Stbl == nil {
continue
}
stsd := trak.Mdia.Minf.Stbl.Stsd
if stsd == nil || stsd.AvcX == nil || stsd.AvcX.AvcC == nil {
continue
}
for _, s := range stsd.AvcX.AvcC.SPSnalus {
if sps, err := avc.ParseSPSNALUnit(s, true); err == nil {
spsMap[uint32(sps.ParameterID)] = sps
}
}
for _, p := range stsd.AvcX.AvcC.PPSnalus {
if pps, err := avc.ParsePPSNALUnit(p, spsMap); err == nil {
ppsMap[pps.PicParameterSetID] = pps
}
}
}
}
fmt.Println("=== slice headers in inspection window (frame_num / poc / idr_pic_id) ===")
fragIdx := 0
sampIdx := 0
for _, seg := range parsed.Segments {
for _, fr := range seg.Fragments {
fs, err := fr.GetFullSamples(trex)
if err != nil {
fragIdx++
continue
}
for _, s := range fs {
dts := int64(s.DecodeTime)
if dts < winLo || dts > winHi {
sampIdx++
continue
}
for _, nal := range splitAVCC(s.Data) {
t := nal[0] & 0x1f
if t == 1 || t == 5 { // non-IDR or IDR slice
sh, err := avc.ParseSliceHeader(nal, spsMap, ppsMap)
if err != nil {
fmt.Printf(" s%04d frag%d dts=%d nalType=%d sliceHeader ERR: %v\n", sampIdx, fragIdx, dts, t, err)
break
}
fmt.Printf(" s%04d frag%d dts=%d nalType=%d sliceType=%v frameNum=%d idrPicId=%d pocLsb=%d\n",
sampIdx, fragIdx, dts, t, sh.SliceType, sh.FrameNum, sh.IDRPicID, sh.PicOrderCntLsb)
break
}
}
sampIdx++
}
fragIdx++
}
}
}
// splitAVCC splits a length-prefixed (4-byte) AVCC buffer into NAL units.
func splitAVCC(b []byte) [][]byte {
var out [][]byte
i := 0
for i+4 <= len(b) {
n := int(uint32(b[i])<<24 | uint32(b[i+1])<<16 | uint32(b[i+2])<<8 | uint32(b[i+3]))
i += 4
if n <= 0 || i+n > len(b) {
break
}
out = append(out, b[i:i+n])
i += n
}
return out
}
func bytesEqual(a, b []byte) bool {
if len(a) != len(b) {
return false
}
for i := range a {
if a[i] != b[i] {
return false
}
}
return true
}
// nalTypes returns the list of H.264 NAL unit types present in an AVCC
// (length-prefixed) sample buffer.
func nalTypes(b []byte) []int {
var out []int
i := 0
for i+4 <= len(b) {
n := int(uint32(b[i])<<24 | uint32(b[i+1])<<16 | uint32(b[i+2])<<8 | uint32(b[i+3]))
i += 4
if n <= 0 || i+n > len(b) {
break
}
out = append(out, int(b[i]&0x1f))
i += n
}
return out
}
// nalsByType returns the raw NAL payloads (without length prefix) of the given
// type from an AVCC (length-prefixed) sample buffer.
func nalsByType(b []byte, want int) [][]byte {
var out [][]byte
i := 0
for i+4 <= len(b) {
n := int(uint32(b[i])<<24 | uint32(b[i+1])<<16 | uint32(b[i+2])<<8 | uint32(b[i+3]))
i += 4
if n <= 0 || i+n > len(b) {
break
}
if int(b[i]&0x1f) == want {
nal := make([]byte, n)
copy(nal, b[i:i+n])
out = append(out, nal)
}
i += n
}
return out
}
// inspectWindow returns the [lo,hi] decode-time range (track timescale units)
// for which sample-level detail is printed. Explicit -from/-to win; otherwise
// the window auto-centres on the largest gap between consecutive video
// keyframes — the most likely location of a visible artifact — with a margin on
// each side so the frames leading into and out of the gap are shown too.
func inspectWindow(keyDecodeTimes []uint64, from, to int64) (int64, int64) {
if from >= 0 || to >= 0 {
if from < 0 {
from = 0
}
if to < 0 {
to = from + 2000
}
return from, to
}
if len(keyDecodeTimes) < 2 {
return 0, 1 << 62
}
worstIdx, worstGap := 1, uint64(0)
for i := 1; i < len(keyDecodeTimes); i++ {
if g := keyDecodeTimes[i] - keyDecodeTimes[i-1]; g > worstGap {
worstGap = g
worstIdx = i
}
}
const margin = 500
lo := int64(keyDecodeTimes[worstIdx-1]) - margin
if lo < 0 {
lo = 0
}
return lo, int64(keyDecodeTimes[worstIdx]) + margin
}
// summary prints a compact, generic health report so a recording can be
// validated at a glance without reading the full per-sample dump above.
func summary(parsed *mp4ff.File, trex *mp4ff.TrexBox) {
fmt.Println("=== SUMMARY (health checks) ===")
videoTracks, audioTracks := 0, 0
var videoTimescale uint64 = 1
if parsed.Init != nil && parsed.Init.Moov != nil {
for _, trak := range parsed.Init.Moov.Traks {
switch trak.Mdia.Hdlr.HandlerType {
case "vide":
videoTracks++
if trak.Mdia.Mdhd.Timescale != 0 {
videoTimescale = uint64(trak.Mdia.Mdhd.Timescale)
}
case "soun":
audioTracks++
}
}
}
fmt.Printf(" tracks: %d video, %d audio\n", videoTracks, audioTracks)
if audioTracks == 0 {
fmt.Println(" note: no audio track is embedded in this file")
}
type fragStat struct {
idx int
tfdt uint64
dur uint64
nSamp int
nKeys int
zeroDur int
fps float64
}
var stats []fragStat
var keyTimes []uint64
var fpsArr []float64
tfdtGaps := 0
var prevEnd uint64
havePrev := false
fi := 0
for _, seg := range parsed.Segments {
for _, fr := range seg.Fragments {
for _, traf := range fr.Moof.Trafs {
if traf.Tfhd.TrackID != 1 {
continue
}
st := fragStat{idx: fi, tfdt: traf.Tfdt.BaseMediaDecodeTime()}
off := uint64(0)
for _, trun := range traf.Truns {
for _, s := range trun.Samples {
st.nSamp++
if (s.Flags>>24)&0x03 == 0x02 {
st.nKeys++
keyTimes = append(keyTimes, st.tfdt+off)
}
if s.Dur == 0 {
st.zeroDur++
}
off += uint64(s.Dur)
}
}
st.dur = off
d := st.dur
if d == 0 {
d = 1
}
st.fps = float64(st.nSamp) * float64(videoTimescale) / float64(d)
fpsArr = append(fpsArr, st.fps)
if havePrev && st.tfdt != prevEnd {
tfdtGaps++
}
prevEnd = st.tfdt + st.dur
havePrev = true
stats = append(stats, st)
}
fi++
}
}
medFps := medianFloat(fpsArr)
fmt.Printf(" fragments: %d (video timescale=%d, median %.1f fps)\n", len(stats), videoTimescale, medFps)
lowFps := 0
totalZero := 0
for _, st := range stats {
totalZero += st.zeroDur
flagStr := ""
if medFps > 0 && st.fps < medFps*0.9 {
lowFps++
flagStr = " <== LOW FRAME RATE — likely dropped frames"
}
fmt.Printf(" frag%02d tfdt=%-6d dur=%-5d samples=%-3d keyframes=%d zeroDur=%d fps=%.1f%s\n",
st.idx, st.tfdt, st.dur, st.nSamp, st.nKeys, st.zeroDur, st.fps, flagStr)
}
var gaps []uint64
for i := 1; i < len(keyTimes); i++ {
gaps = append(gaps, keyTimes[i]-keyTimes[i-1])
}
irregular := 0
if len(gaps) > 0 {
med := medianUint(gaps)
mn, mx := gaps[0], gaps[0]
for _, g := range gaps {
if g < mn {
mn = g
}
if g > mx {
mx = g
}
// Flag intervals that deviate by more than ~50% from the median GOP.
if med > 0 && (g*2 > med*3 || g*2 < med) {
irregular++
}
}
fmt.Printf(" keyframe gaps: min=%d median=%d max=%d irregular=%d/%d\n", mn, med, mx, irregular, len(gaps))
}
fmt.Printf(" tfdt discontinuities: %d\n", tfdtGaps)
fmt.Printf(" zero-duration samples: %d\n", totalZero)
fmt.Println(" verdict:")
clean := true
if audioTracks == 0 {
fmt.Println(" - no audio track (expected if this recording is video-only)")
}
if lowFps > 0 {
clean = false
fmt.Printf(" - %d fragment(s) have a reduced frame rate (dropped frames) — likely source of the artifacts\n", lowFps)
}
if irregular > 0 {
clean = false
fmt.Printf(" - %d irregular keyframe interval(s)\n", irregular)
}
if tfdtGaps > 0 {
clean = false
fmt.Printf(" - %d timeline (tfdt) discontinuity(ies)\n", tfdtGaps)
}
if totalZero > 0 {
clean = false
fmt.Printf(" - %d zero-duration sample(s)\n", totalZero)
}
if clean {
fmt.Println(" - container structure looks healthy")
}
}
func medianUint(v []uint64) uint64 {
if len(v) == 0 {
return 0
}
c := append([]uint64(nil), v...)
sort.Slice(c, func(i, j int) bool { return c[i] < c[j] })
return c[len(c)/2]
}
func medianFloat(v []float64) float64 {
if len(v) == 0 {
return 0
}
c := append([]float64(nil), v...)
sort.Float64s(c)
return c[len(c)/2]
}

View File

@@ -106,9 +106,9 @@
"mqtturi": "tcp://mqtt.kerberos.io:1883",
"mqtt_username": "",
"mqtt_password": "",
"stunuri": "stun:turn.kerberos.io:8443",
"turn_force": "false",
"turnuri": "turn:turn.kerberos.io:8443",
"stunuri": "stun:turn-fra1.kerberos.io:3478",
"turnuri": "turn:turn-fra1.kerberos.io:3478",
"turn_username": "username1",
"turn_password": "password1",
"heartbeaturi": "",

View File

@@ -46,7 +46,11 @@ func UploadDropbox(configuration *models.Configuration, fileName string) (bool,
file, err := os.OpenFile(fullname, os.O_RDWR, 0755)
if file != nil {
defer file.Close()
defer func() {
if cerr := file.Close(); cerr != nil {
log.Log.Error("UploadDropbox: Error closing file: " + cerr.Error())
}
}()
}
if err == nil {

View File

@@ -34,6 +34,29 @@ func UploadKerberosHub(configuration *models.Configuration, fileName string) (bo
log.Log.Info("UploadKerberosHub: Uploading to Kerberos Hub (" + config.HubURI + ")")
log.Log.Info("UploadKerberosHub: Upload started for " + fileName)
// Prefer the resumable (tus) upload when enabled (the default). Kerberos Hub
// authenticates the agent with its Hub public/private key and proxies the
// resumable upload to the Kerberos Vault. When Hub does not expose a tus
// endpoint (older deployments) we transparently fall back to the legacy
// single-POST upload below.
if resumableUploadsEnabled() {
uploaded, _, supported, body, rerr := uploadHubResumable(&config, fileName, "UploadKerberosHub", "hub")
if supported {
if uploaded {
log.Log.Info("UploadKerberosHub: Upload Finished (resumable), " + body)
return true, true, nil
}
if rerr != nil {
log.Log.Info("UploadKerberosHub: resumable upload failed, " + rerr.Error())
} else {
log.Log.Info("UploadKerberosHub: resumable upload incomplete, " + body)
}
return false, true, rerr
}
log.Log.Info("UploadKerberosHub: resumable (tus) endpoint not available, falling back to legacy upload")
}
fullname := "data/recordings/" + fileName
// Check if we still have the file otherwise we abort the request.

View File

@@ -30,6 +30,15 @@ func UploadKerberosVault(configuration *models.Configuration, fileName string) (
return false, false, errors.New(err)
}
// If the recording no longer exists on disk there is nothing to upload.
// 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 {
log.Log.Info("UploadKerberosVault: skipping " + fileName + ", file doesn't exist anymore")
return false, false, nil
}
// timestamp_microseconds_instanceName_regionCoordinates_numberOfChanges_token
// 1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4
// - Timestamp
@@ -41,17 +50,6 @@ func UploadKerberosVault(configuration *models.Configuration, fileName string) (
// KerberosCloud, this means storage is disabled and proxy enabled.
log.Log.Info("UploadKerberosVault: Uploading to Kerberos Vault (" + config.KStorage.URI + ")")
log.Log.Info("UploadKerberosVault: Upload started for " + fileName)
fullname := "data/recordings/" + fileName
file, err := os.OpenFile(fullname, os.O_RDWR, 0755)
if file != nil {
defer file.Close()
}
if err != nil {
err := "UploadKerberosVault: Upload Failed, file doesn't exists anymore"
log.Log.Info(err)
return false, false, errors.New(err)
}
publicKey := config.KStorage.CloudKey
if config.HubKey != "" {
@@ -60,62 +58,30 @@ func UploadKerberosVault(configuration *models.Configuration, fileName string) (
// We need to check if we are in a retry timeout.
if kstorageRetryTimeout <= time.Now().Unix() {
uploaded, responded, body, err := sendToVault(*config.KStorage, publicKey, config.Key, fileName, "UploadKerberosVault", "primary")
if uploaded {
kstorageRetryCount = 0
log.Log.Info("UploadKerberosVault: Upload Finished, " + body)
return true, true, nil
}
req, err := http.NewRequest("POST", config.KStorage.URI+"/storage", file)
if err != nil {
errorMessage := "UploadKerberosVault: error reading request, " + config.KStorage.URI + "/storage: " + err.Error()
log.Log.Error(errorMessage)
return false, true, errors.New(errorMessage)
}
req.Header.Set("Content-Type", "video/mp4")
req.Header.Set("X-Kerberos-Storage-CloudKey", publicKey)
req.Header.Set("X-Kerberos-Storage-AccessKey", config.KStorage.AccessKey)
req.Header.Set("X-Kerberos-Storage-SecretAccessKey", config.KStorage.SecretAccessKey)
req.Header.Set("X-Kerberos-Storage-Provider", config.KStorage.Provider)
req.Header.Set("X-Kerberos-Storage-FileName", fileName)
req.Header.Set("X-Kerberos-Storage-Device", config.Key)
req.Header.Set("X-Kerberos-Storage-Capture", "IPCamera")
req.Header.Set("X-Kerberos-Storage-Directory", config.KStorage.Directory)
var client *http.Client
if os.Getenv("AGENT_TLS_INSECURE") == "true" {
tr := &http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
}
client = &http.Client{Transport: tr}
} else {
client = &http.Client{}
}
resp, err := client.Do(req)
if resp != nil {
defer resp.Body.Close()
}
if err == nil {
if resp != nil {
body, err := io.ReadAll(resp.Body)
if err == nil {
if resp.StatusCode == 200 {
kstorageRetryCount = 0
log.Log.Info("UploadKerberosVault: Upload Finished, " + resp.Status + ", " + string(body))
return true, true, nil
} else {
// We increase the retry count, and set the timeout.
// If we have reached the retry policy, we set the timeout.
// This means we will not retry for the next 5 minutes.
if kstorageRetryCount < config.KStorage.MaxRetries {
kstorageRetryCount = (kstorageRetryCount + 1)
}
if kstorageRetryCount == config.KStorage.MaxRetries {
kstorageRetryTimeout = time.Now().Add(time.Duration(config.KStorage.Timeout) * time.Second).Unix()
}
log.Log.Info("UploadKerberosVault: Upload Failed, " + resp.Status + ", " + string(body))
}
}
}
} else {
log.Log.Info("UploadKerberosVault: Upload Failed, " + err.Error())
} else {
log.Log.Info("UploadKerberosVault: Upload Failed, " + body)
}
// We only advance the retry policy when the vault gave a definitive
// response (mirroring the original behaviour where transient network
// errors did not consume retries). When the retry count reaches the
// configured maximum we back off for the configured timeout.
if responded {
if kstorageRetryCount < config.KStorage.MaxRetries {
kstorageRetryCount = (kstorageRetryCount + 1)
}
if kstorageRetryCount == config.KStorage.MaxRetries {
kstorageRetryTimeout = time.Now().Add(time.Duration(config.KStorage.Timeout) * time.Second).Unix()
}
}
}
@@ -134,61 +100,116 @@ func UploadKerberosVault(configuration *models.Configuration, fileName string) (
log.Log.Info("UploadKerberosVault (Secondary): Uploading to Secondary Kerberos Vault (" + config.KStorageSecondary.URI + ")")
file, err = os.OpenFile(fullname, os.O_RDWR, 0755)
if file != nil {
defer file.Close()
}
if err != nil {
err := "UploadKerberosVault (Secondary): Upload Failed, file doesn't exists anymore"
log.Log.Info(err)
return false, false, errors.New(err)
uploaded, _, body, err := sendToVault(*config.KStorageSecondary, publicKey, config.Key, fileName, "UploadKerberosVault (Secondary)", "secondary")
if uploaded {
log.Log.Info("UploadKerberosVault (Secondary): Upload Finished to secondary, " + body)
return true, true, nil
}
req, err := http.NewRequest("POST", config.KStorageSecondary.URI+"/storage", file)
if err != nil {
errorMessage := "UploadKerberosVault (Secondary): error reading request, " + config.KStorageSecondary.URI + "/storage: " + err.Error()
log.Log.Error(errorMessage)
return false, true, errors.New(errorMessage)
}
req.Header.Set("Content-Type", "video/mp4")
req.Header.Set("X-Kerberos-Storage-CloudKey", publicKey)
req.Header.Set("X-Kerberos-Storage-AccessKey", config.KStorageSecondary.AccessKey)
req.Header.Set("X-Kerberos-Storage-SecretAccessKey", config.KStorageSecondary.SecretAccessKey)
req.Header.Set("X-Kerberos-Storage-Provider", config.KStorageSecondary.Provider)
req.Header.Set("X-Kerberos-Storage-FileName", fileName)
req.Header.Set("X-Kerberos-Storage-Device", config.Key)
req.Header.Set("X-Kerberos-Storage-Capture", "IPCamera")
req.Header.Set("X-Kerberos-Storage-Directory", config.KStorageSecondary.Directory)
var client *http.Client
if os.Getenv("AGENT_TLS_INSECURE") == "true" {
tr := &http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
}
client = &http.Client{Transport: tr}
log.Log.Info("UploadKerberosVault (Secondary): Upload Failed to secondary, " + err.Error())
} else {
client = &http.Client{}
}
resp, err := client.Do(req)
if resp != nil {
defer resp.Body.Close()
}
if err == nil {
if resp != nil {
body, err := io.ReadAll(resp.Body)
if err == nil {
if resp.StatusCode == 200 {
log.Log.Info("UploadKerberosVault (Secondary): Upload Finished to secondary, " + resp.Status + ", " + string(body))
return true, true, nil
} else {
log.Log.Info("UploadKerberosVault (Secondary): Upload Failed to secondary, " + resp.Status + ", " + string(body))
}
}
}
log.Log.Info("UploadKerberosVault (Secondary): Upload Failed to secondary, " + body)
}
}
return false, true, nil
}
// sendToVault uploads a single recording to one Kerberos Vault. When resumable
// uploads are enabled (the default) it attempts the tus protocol first and, if
// the vault does not expose a tus endpoint (older deployments), transparently
// falls back to the legacy single-shot POST.
//
// It returns whether the upload succeeded, whether the vault gave a definitive
// HTTP response (so the caller can advance its retry policy), a short message
// for logging, and a transport error if any.
func sendToVault(vault models.KStorage, publicKey, deviceKey, fileName, label, slot string) (bool, bool, string, error) {
if resumableUploadsEnabled() {
uploaded, responded, supported, body, err := uploadVaultResumable(vault, publicKey, deviceKey, fileName, label, slot)
if supported {
return uploaded, responded, body, err
}
log.Log.Info(label + ": resumable (tus) endpoint not available, falling back to legacy upload")
}
return uploadVaultLegacy(vault, publicKey, deviceKey, fileName, label)
}
// uploadVaultLegacy performs the original single-request upload: the whole file
// is sent as the body of a POST to {URI}/storage. Kept for backwards
// compatibility with vault deployments that do not support resumable uploads.
func uploadVaultLegacy(vault models.KStorage, publicKey, deviceKey, fileName, label string) (bool, bool, string, error) {
fullname := "data/recordings/" + fileName
file, err := os.Open(fullname)
if file != nil {
defer file.Close()
}
if err != nil {
msg := label + ": Upload Failed, file doesn't exists anymore"
log.Log.Info(msg)
return false, false, "", errors.New(msg)
}
uri := vault.URI
for len(uri) > 0 && uri[len(uri)-1] == '/' {
uri = uri[:len(uri)-1]
}
req, err := http.NewRequest("POST", uri+"/storage", file)
if err != nil {
errorMessage := label + ": error reading request, " + uri + "/storage: " + err.Error()
log.Log.Error(errorMessage)
return false, false, "", errors.New(errorMessage)
}
req.Header.Set("Content-Type", "video/mp4")
setVaultHeaders(req.Header, vault, publicKey, deviceKey, fileName)
client := newVaultHTTPClient(0)
resp, err := client.Do(req)
if resp != nil {
defer resp.Body.Close()
}
if err != nil {
return false, false, "", err
}
body, rerr := io.ReadAll(resp.Body)
if rerr != nil {
return false, false, "", rerr
}
if resp.StatusCode == 200 {
return true, true, resp.Status + ", " + string(body), nil
}
return false, true, resp.Status + ", " + string(body), nil
}
// setVaultHeaders sets the standard Kerberos Vault headers used by the legacy
// single-POST upload.
func setVaultHeaders(h http.Header, vault models.KStorage, publicKey, deviceKey, fileName string) {
h.Set("X-Kerberos-Storage-CloudKey", publicKey)
h.Set("X-Kerberos-Storage-AccessKey", vault.AccessKey)
h.Set("X-Kerberos-Storage-SecretAccessKey", vault.SecretAccessKey)
h.Set("X-Kerberos-Storage-Provider", vault.Provider)
h.Set("X-Kerberos-Storage-FileName", fileName)
h.Set("X-Kerberos-Storage-Device", deviceKey)
h.Set("X-Kerberos-Storage-Capture", "IPCamera")
h.Set("X-Kerberos-Storage-Directory", vault.Directory)
}
// newVaultHTTPClient builds an HTTP client honouring the AGENT_TLS_INSECURE
// escape hatch. A timeout of 0 disables the client-level timeout, which is
// required for streaming large upload bodies.
func newVaultHTTPClient(timeout time.Duration) *http.Client {
client := &http.Client{}
if os.Getenv("AGENT_TLS_INSECURE") == "true" {
client.Transport = &http.Transport{
TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
}
}
if timeout > 0 {
client.Timeout = timeout
}
return client
}

View File

@@ -0,0 +1,177 @@
package cloud
import (
"time"
mqtt "github.com/eclipse/paho.mqtt.golang"
"github.com/kerberos-io/agent/machinery/src/capture"
"github.com/kerberos-io/agent/machinery/src/cloud/livehls"
"github.com/kerberos-io/agent/machinery/src/log"
"github.com/kerberos-io/agent/machinery/src/models"
"github.com/kerberos-io/agent/machinery/src/packets"
)
// hlsViewerTimeoutSeconds is how long the agent keeps shipping live HLS segments
// after the last viewer keepalive. It is a few seconds longer than the segment
// duration so a viewer whose keepalive is briefly delayed does not cause the
// session to flap. When it lapses the session is torn down to stop wasting
// upload bandwidth when nobody is watching.
const hlsViewerTimeoutSeconds = 8
// hlsReadyReannounceSeconds throttles how often the agent re-announces an
// already-ready session over MQTT in response to viewer keepalives. The initial
// "receive-hls-ready" is a one-shot fired when the first segment lands; a viewer
// that connects or hard-refreshes after that (while the session is still alive)
// missed it, so we re-announce on subsequent keepalives. Viewers dedupe by
// session id, so a re-announce for a session they already play is a no-op. ~2s
// gets a refreshed viewer playing well within its connection timeout without
// spamming the control plane.
const hlsReadyReannounceSeconds = 2
// HandleLiveStreamHLS drives the live HLS producer. It mirrors HandleLiveStreamSD:
// it reads the camera's packet stream from a Latest() cursor, and while a viewer
// is active (kept alive via communication.HandleLiveHLS) it muxes the packets
// into CMAF segments and ships them to hub-api, which stores each segment in an
// ephemeral, short-TTL live window and serves the rolling playlist to viewers.
//
// A session is created lazily on the first keyframe seen while a viewer is active
// and torn down once viewers go away, so an idle camera produces no live traffic.
func HandleLiveStreamHLS(livestreamCursor *packets.QueueCursor, configuration *models.Configuration, communication *models.Communication, mqttClient mqtt.Client, _ capture.RTSPClient) {
log.Log.Debug("cloud.HandleLiveStreamHLS(): started")
config := configuration.Config
if config.Offline == "true" {
log.Log.Debug("cloud.HandleLiveStreamHLS(): stopping as Offline is enabled.")
return
}
if config.Capture.Liveview == "false" {
log.Log.Debug("cloud.HandleLiveStreamHLS(): stopping as Liveview is disabled.")
return
}
if config.HubURI == "" || config.HubKey == "" {
log.Log.Debug("cloud.HandleLiveStreamHLS(): stopping as the Hub is not configured (HubURI/HubKey).")
return
}
hubKey := config.HubKey
deviceId := config.Key
region := ""
if config.S3 != nil {
region = config.S3.Region
}
publisher := livehls.NewPublisher(livehls.PublisherConfig{
HubURI: config.HubURI,
HubKey: config.HubKey,
HubPrivateKey: config.HubPrivateKey,
Region: region,
DeviceKey: deviceId,
})
// Encoded dimensions are only needed for the avcC fallback path (an SPS that
// mp4ff's strict parser rejects); the main stream dimensions are a safe value.
width := uint16(config.Capture.IPCamera.Width)
height := uint16(config.Capture.IPCamera.Height)
var session *livehls.Session
lastViewerRequest := int64(0)
lastReadyAnnounce := int64(0)
var cursorError error
var pkt packets.Packet
for cursorError == nil {
pkt, cursorError = livestreamCursor.ReadPacket()
now := time.Now().Unix()
select {
case <-communication.HandleLiveHLS:
lastViewerRequest = now
// A keepalive may come from a viewer that just connected or hard-
// refreshed and therefore missed the one-shot readiness announcement
// fired when this session's first segment landed. Re-announce (throttled)
// so late/refreshed viewers learn the active session id; the frontend
// dedupes by session id, so this is a no-op for viewers already playing.
if session != nil && session.IsReady() && now-lastReadyAnnounce >= hlsReadyReannounceSeconds {
publishHLSReady(configuration, mqttClient, hubKey, deviceId, session.SessionID())
lastReadyAnnounce = now
}
default:
}
viewerActive := now-lastViewerRequest <= hlsViewerTimeoutSeconds
if !viewerActive {
// No viewer: stop and discard the session so we stop shipping segments.
if session != nil {
_ = session.Close()
log.Log.Info("cloud.HandleLiveStreamHLS(): no active viewers, stopped live HLS session " + session.SessionID())
session = nil
}
continue
}
if len(pkt.Data) == 0 || !pkt.IsVideo {
continue
}
// Start a session lazily, but only on a keyframe so the first segment opens
// on a random-access point.
if session == nil {
if !pkt.IsKeyFrame {
continue
}
session = livehls.NewSession(publisher, livehls.SessionOptions{
Codec: pkt.Codec,
SPSNALUs: config.Capture.IPCamera.SPSNALUs,
PPSNALUs: config.Capture.IPCamera.PPSNALUs,
VPSNALUs: config.Capture.IPCamera.VPSNALUs,
Width: width,
Height: height,
})
session.SetOnReady(func(sessionID string) {
log.Log.Info("cloud.HandleLiveStreamHLS(): live HLS session ready, announcing " + sessionID)
publishHLSReady(configuration, mqttClient, hubKey, deviceId, sessionID)
lastReadyAnnounce = time.Now().Unix()
})
log.Log.Info("cloud.HandleLiveStreamHLS(): started live HLS session " + session.SessionID())
}
if err := session.WritePacket(pkt); err != nil {
log.Log.Error("cloud.HandleLiveStreamHLS(): " + err.Error())
}
}
if session != nil {
_ = session.Close()
}
log.Log.Debug("cloud.HandleLiveStreamHLS(): finished")
}
// publishHLSReady announces, over MQTT, that a live HLS session is available so
// viewers can load the rolling playlist hub-api serves for {device}/{session}.
func publishHLSReady(configuration *models.Configuration, mqttClient mqtt.Client, hubKey, deviceId, sessionID string) {
valueMap := map[string]interface{}{
"session": sessionID,
"device": deviceId,
}
message := models.Message{
Payload: models.Payload{
Action: "receive-hls-ready",
DeviceId: deviceId,
Value: valueMap,
},
}
payload, err := models.PackageMQTTMessage(configuration, message)
if err == nil {
mqttClient.Publish("kerberos/hub/"+hubKey, 0, false, payload)
log.Log.Info("cloud.HandleLiveStreamHLS(): announced live HLS session " + sessionID)
} else {
log.Log.Error("cloud.HandleLiveStreamHLS(): failed to package receive-hls-ready message: " + err.Error())
}
}

View File

@@ -0,0 +1,202 @@
// Package livehls implements the agent-side producer for live HLS streaming.
//
// It complements the recording pipeline: where recordings are muxed into one
// fragmented MP4 and uploaded resumably (TUS) when complete, live HLS ships a
// continuous series of small, independently-decodable CMAF segments to hub-api
// the instant each is produced, so a browser can play a near-live HLS stream
// without WebRTC/TURN (outbound HTTPS only).
//
// The wire contract (agent -> hub-api) intentionally mirrors the existing
// header-based storage convention (X-Kerberos-Storage-Device / -FileName, plus
// the Hub public/private key auth headers). hub-api authenticates the agent and
// stores each segment in an ephemeral, short-TTL live window keyed by
// {device}/{session}, which it serves straight back to the browser. The live
// window is deliberately kept out of the vault and the recordings collection;
// durable archival/DVR is a separate, later concern.
//
// Unlike recordings, live segments are NOT uploaded resumably: a 1-2s segment
// that fails to upload is stale by the time a retry would land, so the publisher
// is fire-and-forget and drops on failure (logged) rather than blocking the live
// pipeline behind a retry/handshake.
package livehls
import (
"bytes"
"context"
"fmt"
"net/http"
"strconv"
"strings"
"time"
"github.com/kerberos-io/agent/machinery/src/log"
"github.com/kerberos-io/agent/machinery/src/video"
)
const (
// liveIngestPath is the hub-api endpoint that accepts a single live segment
// (or the init segment) and stores it in the ephemeral live window. hub-api
// distinguishes init vs media segment and the object name via the
// X-Kerberos-Live-* headers below, keeping a single route (mirrors the
// existing /storage/upload convention).
liveIngestPath = "/storage/live"
// Object names within a session. The init segment (ftyp+moov) is uploaded
// once per session; media segments are seg-<sequence>.m4s.
initObjectName = "init.mp4"
contentTypeInit = "video/mp4"
contentTypeSegment = "video/iso.segment"
// Header names for the live ingest contract.
headerHubPublicKey = "X-Kerberos-Hub-PublicKey"
headerHubPrivateKey = "X-Kerberos-Hub-PrivateKey"
headerHubRegion = "X-Kerberos-Hub-Region"
headerStorageDevice = "X-Kerberos-Storage-Device"
headerLiveSession = "X-Kerberos-Live-Session"
headerLiveName = "X-Kerberos-Live-Name"
headerLiveSequence = "X-Kerberos-Live-Sequence"
headerLiveDuration = "X-Kerberos-Live-Duration"
// defaultPublishTimeout bounds a single segment upload. A live segment that
// cannot be delivered within roughly its own duration is stale, so the upload
// is abandoned (dropped) rather than allowed to back up the pipeline.
defaultPublishTimeout = 4 * time.Second
)
// PublisherConfig carries the hub endpoint and credentials needed to ship live
// segments. It is populated from the agent's models.Config (HubURI/HubKey/...).
type PublisherConfig struct {
HubURI string // base hub-api URL, e.g. https://api.hub.example.com
HubKey string // Hub public key (X-Kerberos-Hub-PublicKey)
HubPrivateKey string // Hub private key (X-Kerberos-Hub-PrivateKey)
Region string // storage region (X-Kerberos-Hub-Region), may be empty
DeviceKey string // device/camera key (X-Kerberos-Storage-Device)
// Timeout optionally overrides defaultPublishTimeout (used by tests).
Timeout time.Duration
// HTTPClient optionally injects a client (used by tests). When nil a
// redirect-credential-stripping client is created.
HTTPClient *http.Client
}
// Publisher ships init and media segments to hub-api over plain HTTP POST.
//
// It is safe for sequential use from a single live-stream goroutine. Methods are
// fire-and-forget: they return an error for the caller to log, but the caller is
// expected to continue (drop-on-fail) rather than retry.
type Publisher struct {
cfg PublisherConfig
client *http.Client
}
// NewPublisher builds a Publisher. The HTTP client strips the Hub credential
// headers on a cross-host redirect (net/http does this for standard auth headers
// but not custom-named ones), matching the recording upload client.
func NewPublisher(cfg PublisherConfig) *Publisher {
client := cfg.HTTPClient
if client == nil {
timeout := cfg.Timeout
if timeout <= 0 {
timeout = defaultPublishTimeout
}
client = &http.Client{
Timeout: timeout,
CheckRedirect: stripHubCredentialsOnCrossHostRedirect,
}
}
return &Publisher{cfg: cfg, client: client}
}
// PublishInit uploads the session's init segment (ftyp+moov). It must be called
// (and succeed) before the player can use any media segment, so the caller
// should treat a failure here as "session not yet established" and retry on the
// next init opportunity rather than shipping media segments blindly.
func (p *Publisher) PublishInit(ctx context.Context, sessionID string, data []byte) error {
return p.post(ctx, postParams{
sessionID: sessionID,
name: initObjectName,
contentType: contentTypeInit,
body: data,
})
}
// PublishSegment uploads one media segment (styp+moof+mdat). The segment's
// sequence number and duration travel in headers so hub-api can update the
// rolling playlist window without parsing the box structure.
func (p *Publisher) PublishSegment(ctx context.Context, sessionID string, seg video.LiveSegment) error {
return p.post(ctx, postParams{
sessionID: sessionID,
name: fmt.Sprintf("seg-%d.m4s", seg.SequenceNumber),
sequence: seg.SequenceNumber,
durationMs: seg.DurationMs,
hasSegment: true,
contentType: contentTypeSegment,
body: seg.Data,
})
}
type postParams struct {
sessionID string
name string
sequence uint32
durationMs uint64
hasSegment bool
contentType string
body []byte
}
// post performs a single fire-and-forget upload to the live ingest endpoint.
func (p *Publisher) post(ctx context.Context, params postParams) error {
if p.cfg.HubURI == "" {
return fmt.Errorf("livehls: HubURI not configured")
}
if params.sessionID == "" {
return fmt.Errorf("livehls: empty session id")
}
url := strings.TrimRight(p.cfg.HubURI, "/") + liveIngestPath
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(params.body))
if err != nil {
return fmt.Errorf("livehls: build request: %w", err)
}
req.Header.Set("Content-Type", params.contentType)
req.Header.Set(headerStorageDevice, p.cfg.DeviceKey)
req.Header.Set(headerLiveSession, params.sessionID)
req.Header.Set(headerLiveName, params.name)
if params.hasSegment {
req.Header.Set(headerLiveSequence, strconv.FormatUint(uint64(params.sequence), 10))
req.Header.Set(headerLiveDuration, strconv.FormatUint(params.durationMs, 10))
}
req.Header.Set(headerHubPublicKey, p.cfg.HubKey)
req.Header.Set(headerHubPrivateKey, p.cfg.HubPrivateKey)
req.Header.Set(headerHubRegion, p.cfg.Region)
resp, err := p.client.Do(req)
if err != nil {
return fmt.Errorf("livehls: upload %s: %w", params.name, err)
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("livehls: upload %s rejected: %s", params.name, resp.Status)
}
log.Log.Debug("livehls.Publisher.post(): shipped " + params.name + " for session " + params.sessionID)
return nil
}
// stripHubCredentialsOnCrossHostRedirect removes the Hub credential headers when
// a redirect crosses to a different host. net/http strips standard sensitive
// headers on a cross-host redirect but not custom-named ones, so without this the
// Hub keys could leak to a redirect target.
func stripHubCredentialsOnCrossHostRedirect(req *http.Request, via []*http.Request) error {
if len(via) == 0 {
return nil
}
if req.URL.Host != via[0].URL.Host {
req.Header.Del(headerHubPrivateKey)
req.Header.Del(headerHubPublicKey)
}
return nil
}

View File

@@ -0,0 +1,312 @@
package livehls
import (
"context"
"io"
"net/http"
"net/http/httptest"
"sync"
"testing"
"time"
"github.com/kerberos-io/agent/machinery/src/packets"
"github.com/kerberos-io/agent/machinery/src/video"
)
// captured records one received upload for assertions.
type captured struct {
path string
method string
contentType string
device string
session string
name string
sequence string
duration string
hubPublic string
hubPrivate string
region string
body []byte
}
// newCapturingServer returns an httptest server that records every upload and
// replies with the given status code.
func newCapturingServer(t *testing.T, status int) (*httptest.Server, *[]captured, *sync.Mutex) {
t.Helper()
var mu sync.Mutex
var got []captured
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
body, _ := io.ReadAll(r.Body)
mu.Lock()
got = append(got, captured{
path: r.URL.Path,
method: r.Method,
contentType: r.Header.Get("Content-Type"),
device: r.Header.Get(headerStorageDevice),
session: r.Header.Get(headerLiveSession),
name: r.Header.Get(headerLiveName),
sequence: r.Header.Get(headerLiveSequence),
duration: r.Header.Get(headerLiveDuration),
hubPublic: r.Header.Get(headerHubPublicKey),
hubPrivate: r.Header.Get(headerHubPrivateKey),
region: r.Header.Get(headerHubRegion),
body: body,
})
mu.Unlock()
w.WriteHeader(status)
}))
t.Cleanup(srv.Close)
return srv, &got, &mu
}
func testPublisher(hubURI string) *Publisher {
return NewPublisher(PublisherConfig{
HubURI: hubURI,
HubKey: "pub-key",
HubPrivateKey: "priv-key",
Region: "eu-west",
DeviceKey: "cam-1",
Timeout: 2 * time.Second,
})
}
func TestPublisherPublishInitSendsContractHeaders(t *testing.T) {
srv, got, mu := newCapturingServer(t, http.StatusOK)
p := testPublisher(srv.URL)
if err := p.PublishInit(context.Background(), "sess-1", []byte("INITBYTES")); err != nil {
t.Fatalf("PublishInit: %v", err)
}
mu.Lock()
defer mu.Unlock()
if len(*got) != 1 {
t.Fatalf("server received %d requests, want 1", len(*got))
}
c := (*got)[0]
if c.method != http.MethodPost {
t.Errorf("method=%s, want POST", c.method)
}
if c.path != liveIngestPath {
t.Errorf("path=%s, want %s", c.path, liveIngestPath)
}
if c.contentType != contentTypeInit {
t.Errorf("content-type=%s, want %s", c.contentType, contentTypeInit)
}
if c.device != "cam-1" {
t.Errorf("device=%s, want cam-1", c.device)
}
if c.session != "sess-1" {
t.Errorf("session=%s, want sess-1", c.session)
}
if c.name != initObjectName {
t.Errorf("name=%s, want %s", c.name, initObjectName)
}
if c.hubPublic != "pub-key" || c.hubPrivate != "priv-key" || c.region != "eu-west" {
t.Errorf("auth headers wrong: pub=%q priv=%q region=%q", c.hubPublic, c.hubPrivate, c.region)
}
if string(c.body) != "INITBYTES" {
t.Errorf("body=%q, want INITBYTES", string(c.body))
}
// init must NOT carry segment-only headers.
if c.sequence != "" || c.duration != "" {
t.Errorf("init should not send sequence/duration, got seq=%q dur=%q", c.sequence, c.duration)
}
}
func TestPublisherPublishSegmentSendsSequenceAndDuration(t *testing.T) {
srv, got, mu := newCapturingServer(t, http.StatusOK)
p := testPublisher(srv.URL)
seg := video.LiveSegment{SequenceNumber: 7, DurationMs: 1960, Data: []byte("SEGMENT")}
if err := p.PublishSegment(context.Background(), "sess-9", seg); err != nil {
t.Fatalf("PublishSegment: %v", err)
}
mu.Lock()
defer mu.Unlock()
c := (*got)[0]
if c.contentType != contentTypeSegment {
t.Errorf("content-type=%s, want %s", c.contentType, contentTypeSegment)
}
if c.name != "seg-7.m4s" {
t.Errorf("name=%s, want seg-7.m4s", c.name)
}
if c.sequence != "7" {
t.Errorf("sequence=%s, want 7", c.sequence)
}
if c.duration != "1960" {
t.Errorf("duration=%s, want 1960", c.duration)
}
if string(c.body) != "SEGMENT" {
t.Errorf("body=%q, want SEGMENT", string(c.body))
}
}
func TestPublisherReturnsErrorOnNon2xx(t *testing.T) {
srv, _, _ := newCapturingServer(t, http.StatusInternalServerError)
p := testPublisher(srv.URL)
err := p.PublishSegment(context.Background(), "s", video.LiveSegment{SequenceNumber: 1, Data: []byte("x")})
if err == nil {
t.Fatal("expected an error on 500 response")
}
}
func TestPublisherErrorsWithoutHubURI(t *testing.T) {
p := NewPublisher(PublisherConfig{DeviceKey: "cam"})
if err := p.PublishInit(context.Background(), "s", []byte("x")); err == nil {
t.Fatal("expected error when HubURI is empty")
}
}
// makeAnnexBVideoPacket builds a synthetic capture packet carrying one Annex B
// H.264 access unit at the given decode time (ms).
func makeAnnexBVideoPacket(isKey bool, timeMs int64) packets.Packet {
nalType := byte(0x01)
if isKey {
nalType = 0x65
}
data := []byte{0x00, 0x00, 0x00, 0x01, nalType}
for i := 0; i < 80; i++ {
data = append(data, byte(i))
}
return packets.Packet{
IsVideo: true,
IsKeyFrame: isKey,
Codec: "H264",
Data: data,
TimeLegacy: time.Duration(timeMs) * time.Millisecond,
}
}
func TestSessionShipsInitThenSegmentsAndFiresReady(t *testing.T) {
srv, got, mu := newCapturingServer(t, http.StatusOK)
p := testPublisher(srv.URL)
sess := NewSession(p, SessionOptions{
Codec: "H264",
SPSNALUs: [][]byte{liveTestSPSForSession()},
PPSNALUs: [][]byte{{0x68, 0xce, 0x38, 0x80}},
Width: 640,
Height: 480,
TargetSegmentMs: 2000,
})
var readyCalls int
var readySession string
sess.SetOnReady(func(id string) {
readyCalls++
readySession = id
})
// 4 GOPs of 25 frames @ 40ms = 1s GOPs => with 2s target, 2 segments emitted
// during streaming and a final one on Close.
const gopFrames, gops = 25, 4
for i := 0; i < gopFrames*gops; i++ {
isKey := i%gopFrames == 0
pkt := makeAnnexBVideoPacket(isKey, int64(i*40))
if err := sess.WritePacket(pkt); err != nil {
t.Fatalf("WritePacket(%d): %v", i, err)
}
}
// A non-video packet must be ignored.
if err := sess.WritePacket(packets.Packet{IsAudio: true, Data: []byte{1, 2, 3}}); err != nil {
t.Fatalf("WritePacket(audio): %v", err)
}
if err := sess.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
mu.Lock()
defer mu.Unlock()
var initCount, segCount int
for _, c := range *got {
if c.name == initObjectName {
initCount++
if string(c.body[4:8]) != "ftyp" {
t.Errorf("init body is not an ftyp box: % x", c.body[:12])
}
} else {
segCount++
if c.session != sess.SessionID() {
t.Errorf("segment session=%s, want %s", c.session, sess.SessionID())
}
}
}
if initCount != 1 {
t.Errorf("init uploaded %d times, want exactly 1", initCount)
}
if segCount < 2 {
t.Errorf("got %d segment uploads, want >= 2", segCount)
}
if readyCalls != 1 {
t.Errorf("OnReady fired %d times, want exactly 1", readyCalls)
}
if readySession != sess.SessionID() {
t.Errorf("OnReady session=%s, want %s", readySession, sess.SessionID())
}
}
func TestSessionRetriesInitWhenFirstAttemptFails(t *testing.T) {
// Server fails the first N requests, then succeeds. This proves init is
// re-attempted (not dropped) so the session can still establish.
var mu sync.Mutex
var inits, segs int
failFirst := 1
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
mu.Lock()
defer mu.Unlock()
name := r.Header.Get(headerLiveName)
if name == initObjectName {
inits++
if inits <= failFirst {
w.WriteHeader(http.StatusBadGateway)
return
}
} else {
segs++
}
w.WriteHeader(http.StatusOK)
}))
t.Cleanup(srv.Close)
sess := NewSession(testPublisher(srv.URL), SessionOptions{
Codec: "H264",
SPSNALUs: [][]byte{liveTestSPSForSession()},
PPSNALUs: [][]byte{{0x68, 0xce, 0x38, 0x80}},
Width: 640,
Height: 480,
})
var ready int
sess.SetOnReady(func(string) { ready++ })
for i := 0; i < 60; i++ {
isKey := i%25 == 0
if err := sess.WritePacket(makeAnnexBVideoPacket(isKey, int64(i*40))); err != nil {
t.Fatalf("WritePacket(%d): %v", i, err)
}
}
if err := sess.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
mu.Lock()
defer mu.Unlock()
if inits < 2 {
t.Errorf("init attempted %d times, want >= 2 (first failed then retried)", inits)
}
if segs < 1 {
t.Errorf("no segments delivered after init recovered (segs=%d)", segs)
}
if ready != 1 {
t.Errorf("OnReady fired %d times, want 1", ready)
}
}
// liveTestSPSForSession is the known-good baseline SPS reused across tests.
func liveTestSPSForSession() []byte {
return []byte{0x67, 0x42, 0xc0, 0x1e, 0xd9, 0x00, 0xa0, 0x47, 0xfe, 0xc8}
}

View File

@@ -0,0 +1,255 @@
package livehls
import (
"context"
"crypto/rand"
"encoding/hex"
"fmt"
"sync"
"time"
"github.com/kerberos-io/agent/machinery/src/log"
"github.com/kerberos-io/agent/machinery/src/packets"
"github.com/kerberos-io/agent/machinery/src/video"
)
// DefaultTargetSegmentMs is the nominal live segment length. ~2s keeps standard
// HLS latency reasonable (a player typically buffers ~3 segments) while staying
// large enough that per-segment HTTP overhead is negligible.
const DefaultTargetSegmentMs = 2000
// Session ties a video.LiveSegmenter to a Publisher: it converts capture packets
// into CMAF segments and ships each one to hub-api. Exactly one init segment is
// delivered per session (re-attempted until it lands), after which media
// segments are published and the OnReady signal fires once so the control plane
// (MQTT) can tell viewers the live playlist exists.
//
// A Session is driven from a single goroutine (the live-stream loop); its methods
// are not safe for concurrent use except SessionID, which is immutable.
type Session struct {
id string
publisher *Publisher
segmenter *video.LiveSegmenter
// newContext produces the per-upload context (timeout). Overridable in tests.
newContext func() (context.Context, context.CancelFunc)
mu sync.Mutex
initBytes []byte
initPublished bool
// lastInitAt is when the init segment was last (re)uploaded. The init is
// re-sent periodically so its short TTL in the hub live window never lapses
// mid-session; see refreshInitIfStale.
lastInitAt time.Time
readyFired bool
onReady func(sessionID string)
}
// SessionOptions configures a live HLS session.
type SessionOptions struct {
Codec string // "H264" or "H265"
SPSNALUs [][]byte // parameter sets (raw or Annex B)
PPSNALUs [][]byte //
VPSNALUs [][]byte // H.265 only
Width uint16 // encoded width (for the avcC fallback path)
Height uint16 // encoded height
TargetSegmentMs uint64 // 0 => DefaultTargetSegmentMs
}
// NewSession builds a session with a fresh random id and wires the segmenter's
// init/segment callbacks to the publisher.
func NewSession(publisher *Publisher, opts SessionOptions) *Session {
target := opts.TargetSegmentMs
if target == 0 {
target = DefaultTargetSegmentMs
}
seg := video.NewLiveSegmenter(opts.Codec, opts.SPSNALUs, opts.PPSNALUs, opts.VPSNALUs, target)
seg.SetDimensions(opts.Width, opts.Height)
s := &Session{
id: newSessionID(),
publisher: publisher,
segmenter: seg,
newContext: func() (context.Context, context.CancelFunc) {
return context.WithTimeout(context.Background(), defaultPublishTimeout)
},
}
// The segmenter emits the init segment exactly once; capture it and try to
// ship it. Failures here are non-fatal - publishInitIfNeeded re-attempts
// before the next media segment so a transient hub hiccup at startup does not
// permanently break the session.
seg.OnInit = func(initBytes []byte) error {
s.mu.Lock()
s.initBytes = append([]byte(nil), initBytes...)
s.mu.Unlock()
s.publishInitIfNeeded()
return nil
}
// Each completed media segment is shipped. We only publish a segment once the
// init segment has landed (a media segment is useless without it), and we fire
// OnReady after the first successfully shipped segment.
seg.OnSegment = func(segment video.LiveSegment) error {
if !s.publishInitIfNeeded() {
log.Log.Warning("livehls.Session: dropping segment " +
fmt.Sprintf("%d", segment.SequenceNumber) + " because init has not been delivered yet")
return nil
}
ctx, cancel := s.newContext()
defer cancel()
if err := s.publisher.PublishSegment(ctx, s.id, segment); err != nil {
log.Log.Warning("livehls.Session: " + err.Error())
return nil
}
s.fireReadyOnce()
// Keep the (write-once) init segment from ageing out of the live window
// while the session is still producing media.
s.refreshInitIfStale()
return nil
}
return s
}
// SessionID returns the immutable session identifier used in object keys and the
// MQTT ready signal.
func (s *Session) SessionID() string { return s.id }
// IsReady reports whether the session has delivered its init segment and at
// least one media segment, i.e. the playlist hub-api serves is now playable. It
// lets the live-stream loop re-announce "receive-hls-ready" to viewers that join
// or hard-refresh after the initial one-shot signal (which they would otherwise
// never receive, leaving the stream blank until the session is recreated).
func (s *Session) IsReady() bool {
s.mu.Lock()
defer s.mu.Unlock()
return s.readyFired
}
// SetOnReady registers a callback fired exactly once, after the first media
// segment has been successfully delivered. Used to publish the MQTT
// "receive-hls-ready" signal so viewers can load the playlist.
func (s *Session) SetOnReady(fn func(sessionID string)) {
s.mu.Lock()
s.onReady = fn
s.mu.Unlock()
}
// WritePacket feeds one capture packet into the segmenter. Non-video packets are
// ignored (the spike is video-only). The decode timestamp is derived exactly as
// the recording muxer does: DTS = PTS - compositionOffset, with the composition
// offset forwarded for correct B-frame presentation order.
func (s *Session) WritePacket(pkt packets.Packet) error {
if !pkt.IsVideo {
return nil
}
pts := uint64(pkt.TimeLegacy.Milliseconds())
compositionOffset := pkt.CompositionTime
dts := pts
if compositionOffset > 0 && uint64(compositionOffset) <= pts {
dts = pts - uint64(compositionOffset)
} else if compositionOffset < 0 || uint64(compositionOffset) > pts {
// Guard against invalid offsets to avoid producing a CTS (DTS+CTO) jump.
compositionOffset = 0
}
return s.segmenter.WriteSample(pkt.IsKeyFrame, pkt.Data, dts, int32(compositionOffset))
}
// Close flushes any buffered sample and ships the final segment.
func (s *Session) Close() error {
return s.segmenter.Close()
}
// publishInitIfNeeded ensures the init segment has been delivered, attempting an
// upload if it has not. Returns true once init is known to be published.
func (s *Session) publishInitIfNeeded() bool {
s.mu.Lock()
if s.initPublished {
s.mu.Unlock()
return true
}
initBytes := s.initBytes
s.mu.Unlock()
if len(initBytes) == 0 {
return false
}
ctx, cancel := s.newContext()
defer cancel()
if err := s.publisher.PublishInit(ctx, s.id, initBytes); err != nil {
log.Log.Warning("livehls.Session: init upload failed, will retry: " + err.Error())
return false
}
s.mu.Lock()
s.initPublished = true
s.lastInitAt = time.Now()
s.mu.Unlock()
log.Log.Info("livehls.Session: init segment delivered for session " + s.id)
return true
}
// initRefreshInterval is how often the init segment is re-uploaded so its TTL in
// the hub-api live window never lapses mid-session. The init segment is otherwise
// written only once per session; because the live window expires objects after a
// short TTL (LiveSegmentTTLSeconds, 45s on the hub) the init would age out after
// ~1 minute and the playlist's #EXT-X-MAP would start 404ing, stalling playback.
// Re-uploading well inside that TTL keeps the init alive for the life of the
// session while still letting it expire naturally once the session ends.
const initRefreshInterval = 15 * time.Second
// refreshInitIfStale re-uploads the init segment if it has not been refreshed
// within initRefreshInterval, keeping its created_at (and thus its TTL) current
// for as long as the session is producing segments. It is a no-op until the init
// has first been published. Failures are non-fatal: the next segment retries.
func (s *Session) refreshInitIfStale() {
s.mu.Lock()
if !s.initPublished || time.Since(s.lastInitAt) < initRefreshInterval {
s.mu.Unlock()
return
}
initBytes := s.initBytes
s.mu.Unlock()
if len(initBytes) == 0 {
return
}
ctx, cancel := s.newContext()
defer cancel()
if err := s.publisher.PublishInit(ctx, s.id, initBytes); err != nil {
log.Log.Warning("livehls.Session: init refresh failed, will retry: " + err.Error())
return
}
s.mu.Lock()
s.lastInitAt = time.Now()
s.mu.Unlock()
log.Log.Debug("livehls.Session: refreshed init segment TTL for session " + s.id)
}
// fireReadyOnce invokes the OnReady callback the first time it is called.
func (s *Session) fireReadyOnce() {
s.mu.Lock()
if s.readyFired || s.onReady == nil {
s.mu.Unlock()
return
}
s.readyFired = true
fn := s.onReady
s.mu.Unlock()
fn(s.id)
}
// newSessionID returns a short, unique, URL-safe session identifier of the form
// <unix-seconds>-<random-hex>.
func newSessionID() string {
b := make([]byte, 4)
if _, err := rand.Read(b); err != nil {
// rand.Read essentially never fails; fall back to a time-only id.
return fmt.Sprintf("%d", time.Now().UnixNano())
}
return fmt.Sprintf("%d-%s", time.Now().Unix(), hex.EncodeToString(b))
}

View File

@@ -0,0 +1,546 @@
package cloud
import (
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"net/url"
"os"
"path/filepath"
"sort"
"strconv"
"strings"
"time"
"github.com/kerberos-io/agent/machinery/src/log"
"github.com/kerberos-io/agent/machinery/src/models"
)
// tusResumableVersion is the tus protocol version implemented by this client.
const tusResumableVersion = "1.0.0"
// tusUploadPath is appended to the configured Kerberos Vault URI to reach the
// resumable upload endpoint. It mirrors how the legacy uploader appends
// "/storage".
const tusUploadPath = "/storage/tus/"
// tusResumeState is persisted in a sidecar file next to the agent data so an
// interrupted upload can be resumed across retries and even agent restarts.
type tusResumeState struct {
UploadURL string `json:"upload_url"`
VaultURI string `json:"vault_uri"`
Size int64 `json:"size"`
}
// resumableUploadsEnabled reports whether the resumable (tus) upload path should
// be attempted. It is enabled by default and can be disabled (falling back to
// the legacy single POST) by setting AGENT_DISABLE_RESUMABLE_UPLOAD=true.
func resumableUploadsEnabled() bool {
return os.Getenv("AGENT_DISABLE_RESUMABLE_UPLOAD") != "true"
}
// tusDefaultChunkSize is the number of bytes uploaded per PATCH request when no
// explicit size is configured. Splitting the upload into chunks keeps each HTTP
// request small enough for intermediary proxies/load balancers and checkpoints
// progress frequently, so an interruption resumes with minimal re-upload.
const tusDefaultChunkSize int64 = 1 << 20 // 1 MiB
const tusProgressBucketPercent int64 = 10
// tusChunkSize returns the number of bytes to send per PATCH request. It
// defaults to tusDefaultChunkSize (1 MiB) and can be overridden with the
// AGENT_TUS_CHUNK_SIZE_BYTES environment variable. A value of 0 (or negative)
// disables chunking and sends the remaining bytes in a single PATCH.
func tusChunkSize() int64 {
v := os.Getenv("AGENT_TUS_CHUNK_SIZE_BYTES")
if v == "" {
return tusDefaultChunkSize
}
n, err := strconv.ParseInt(v, 10, 64)
if err != nil {
return tusDefaultChunkSize
}
if n <= 0 {
return 0 // chunking disabled: send everything in one PATCH
}
return n
}
func tusProgressBucket(offset, size int64) int64 {
if size <= 0 {
return 100
}
percent := (offset * 100) / size
if percent > 100 {
percent = 100
}
return percent / tusProgressBucketPercent
}
func logTusUploadProgress(label string, offset, size int64, loggedBucket *int64) {
bucket := tusProgressBucket(offset, size)
if bucket <= *loggedBucket {
return
}
*loggedBucket = bucket
percent := bucket * tusProgressBucketPercent
if percent > 100 {
percent = 100
}
log.Log.Infof("%s: resumable upload progress %d%% (%d/%d bytes)", label, percent, offset, size)
}
// tusHeaderFunc sets the authentication and routing headers required on every
// tus request for a particular upload target (Kerberos Vault directly, or
// Kerberos Hub which proxies to a vault). fileName is only meaningful on the
// creation request; it is empty on HEAD/PATCH/DELETE.
type tusHeaderFunc func(h http.Header, fileName string)
// runTusUpload performs a resumable (tus) upload of data/recordings/<fileName>
// to baseURL, sending target-specific authentication/routing headers via
// setHeaders on every request. It encapsulates the create/resume/chunk/finalize
// state machine shared by the Kerberos Vault (direct) and Kerberos Hub (proxied)
// upload paths.
//
// Return values:
// - uploaded: the recording was fully received and persisted by the server.
// - responded: the server returned a definitive HTTP response (used by the
// caller to advance its retry/secondary-failover policy).
// - supported: the server exposes a tus endpoint. When false, the caller
// should fall back to the legacy single-POST upload (older deployments).
// - body: a short message for logging.
func runTusUpload(baseURL, metadata, fileName, label, slot string, setHeaders tusHeaderFunc) (uploaded bool, responded bool, supported bool, body string, err error) {
fullname := "data/recordings/" + fileName
file, ferr := os.Open(fullname)
if file != nil {
defer file.Close()
}
if ferr != nil {
msg := label + ": resumable upload failed, file doesn't exist anymore"
log.Log.Info(msg)
// The file is gone, so the legacy path cannot help either. Report it as
// "supported" to avoid a pointless fallback attempt.
return false, false, true, "", errors.New(msg)
}
info, serr := file.Stat()
if serr != nil {
return false, false, true, "", serr
}
size := info.Size()
client := newVaultHTTPClient(0)
client.CheckRedirect = func(req *http.Request, via []*http.Request) error {
if len(via) == 0 {
return nil
}
if req.URL.Host != via[0].URL.Host {
for k := range req.Header {
if strings.HasPrefix(http.CanonicalHeaderKey(k), "X-Kerberos-") {
req.Header.Del(k)
}
}
}
return nil
}
sidecar := tusSidecarPath(fileName, slot)
uploadURL := loadTusResumeState(sidecar, baseURL)
const maxAttempts = 4
restartedAfterComplete := false
for attempt := 0; attempt < maxAttempts; attempt++ {
// (1) Ensure we have an active upload URL, creating one if needed.
if uploadURL == "" {
created, status, cerr := tusCreate(client, baseURL, size, metadata, setHeaders, fileName)
if cerr != nil {
if status == http.StatusNotFound || status == http.StatusMethodNotAllowed || status == http.StatusNotImplemented {
// The vault does not implement tus; let the caller fall back.
return false, false, false, "", cerr
}
log.Log.Info(label + ": resumable create failed, " + cerr.Error())
tusBackoff(attempt)
continue
}
uploadURL = created
saveTusResumeState(sidecar, tusResumeState{UploadURL: uploadURL, VaultURI: baseURL, Size: size})
}
// (2) Query the current server-side offset.
offset, status, herr := tusHead(client, uploadURL, setHeaders)
if herr != nil {
if status == http.StatusNotFound || status == http.StatusGone {
// The upload expired/was removed server-side; start over.
removeTusResumeState(sidecar)
uploadURL = ""
continue
}
log.Log.Info(label + ": resumable head failed, " + herr.Error())
tusBackoff(attempt)
continue
}
// (3) All bytes are present but the upload was not finalized (e.g. the
// completion hook failed). A completed tus upload cannot be re-finalized
// with another PATCH, so delete it and re-upload to force a clean finalize.
if offset >= size {
if restartedAfterComplete {
return false, true, true, "resumable finalize did not complete", errors.New(label + ": resumable finalize did not complete")
}
tusTerminate(client, uploadURL, setHeaders)
removeTusResumeState(sidecar)
uploadURL = ""
restartedAfterComplete = true
continue
}
// (4) Stream the remaining bytes to the vault via PATCH, reading directly
// from disk so the recording is never fully buffered in memory. When a chunk
// size is configured the data is sent across several PATCH requests,
// checkpointing the offset after each one so an interruption resumes from the
// last completed chunk instead of re-uploading everything.
chunkSize := tusChunkSize()
progressed := false
patchFailed := false
var lastBody string
loggedProgressBucket := tusProgressBucket(offset, size)
for offset < size {
// Re-seek every chunk so the on-disk position always matches the
// server-acknowledged offset, even if a PATCH was partially accepted.
if _, sErr := file.Seek(offset, io.SeekStart); sErr != nil {
return false, false, true, "", sErr
}
patchLen := size - offset
if chunkSize > 0 && chunkSize < patchLen {
patchLen = chunkSize
}
newOffset, status, respBody, perr := tusPatch(client, uploadURL, offset, patchLen, file, setHeaders)
if perr != nil {
if status >= 400 {
// Definitive rejection (e.g. provider push failed during finalize).
// Re-evaluate via HEAD on the next iteration to decide retry/restart.
log.Log.Info(label + ": resumable patch rejected, " + perr.Error())
} else {
log.Log.Info(label + ": resumable patch failed, " + perr.Error())
}
tusBackoff(attempt)
patchFailed = true
break
}
if newOffset > offset {
progressed = true
}
offset = newOffset
lastBody = respBody
logTusUploadProgress(label, offset, size, &loggedProgressBucket)
if offset < size {
// Partial progress: persist so a later retry resumes from here.
saveTusResumeState(sidecar, tusResumeState{UploadURL: uploadURL, VaultURI: baseURL, Size: size})
}
}
if patchFailed {
if progressed {
// Forward progress refreshes the retry budget: maxAttempts bounds the
// number of consecutive failures, not the number of chunks needed for
// a large recording.
attempt = -1
}
continue
}
// All declared bytes have been sent and acknowledged: the upload is done.
removeTusResumeState(sidecar)
return true, true, true, lastBody, nil
}
return false, true, true, "resumable upload did not complete after retries", errors.New(label + ": resumable upload did not complete after retries")
}
// uploadVaultResumable uploads a recording directly to a Kerberos Vault using
// the tus resumable upload protocol. Credentials travel in the
// X-Kerberos-Storage-* headers on every request and routing (directory/provider)
// is additionally carried in the tus Upload-Metadata.
func uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName, label, slot string) (bool, bool, bool, string, error) {
baseURL := strings.TrimRight(vault.URI, "/") + tusUploadPath
metadata := encodeTusMetadata(map[string]string{
"filename": fileName,
"device": deviceKey,
"directory": vault.Directory,
"provider": vault.Provider,
"capture": "IPCamera",
"cloudkey": publicKey,
})
setHeaders := func(h http.Header, fn string) {
setVaultTusHeaders(h, vault, publicKey, deviceKey, fn)
}
return runTusUpload(baseURL, metadata, fileName, label, slot, setHeaders)
}
// uploadHubResumable uploads a recording to Kerberos Hub's tus endpoint, which
// authenticates the agent with its Hub public/private key and proxies the
// resumable upload to the Kerberos Vault on the agent's behalf. The vault
// directory and provider are resolved and injected by Kerberos Hub, so they are
// intentionally omitted from the metadata here.
func uploadHubResumable(config *models.Config, fileName, label, slot string) (bool, bool, bool, string, error) {
baseURL := strings.TrimRight(config.HubURI, "/") + tusUploadPath
metadata := encodeTusMetadata(map[string]string{
"filename": fileName,
"device": config.Key,
"capture": "IPCamera",
})
setHeaders := func(h http.Header, fn string) {
setHubTusHeaders(h, config, fn)
}
return runTusUpload(baseURL, metadata, fileName, label, slot, setHeaders)
}
// tusCreate performs the tus "creation" request (POST). On success it returns
// the resolved upload URL the agent should use for subsequent HEAD/PATCH calls.
func tusCreate(client *http.Client, baseURL string, size int64, metadata string, setHeaders tusHeaderFunc, fileName string) (string, int, error) {
req, err := http.NewRequest("POST", baseURL, nil)
if err != nil {
return "", 0, err
}
req.Header.Set("Tus-Resumable", tusResumableVersion)
req.Header.Set("Upload-Length", strconv.FormatInt(size, 10))
if metadata != "" {
req.Header.Set("Upload-Metadata", metadata)
}
setHeaders(req.Header, fileName)
resp, err := client.Do(req)
if resp != nil {
defer resp.Body.Close()
}
if err != nil {
return "", 0, err
}
io.Copy(io.Discard, resp.Body)
if resp.StatusCode != http.StatusCreated {
return "", resp.StatusCode, fmt.Errorf("unexpected status creating upload: %s", resp.Status)
}
location := resp.Header.Get("Location")
if location == "" {
return "", resp.StatusCode, errors.New("missing Location header in create response")
}
return resolveTusLocation(baseURL, location), resp.StatusCode, nil
}
// tusHead performs the tus "offset" request (HEAD) and returns the current
// server-side upload offset.
func tusHead(client *http.Client, uploadURL string, setHeaders tusHeaderFunc) (int64, int, error) {
req, err := http.NewRequest("HEAD", uploadURL, nil)
if err != nil {
return 0, 0, err
}
req.Header.Set("Tus-Resumable", tusResumableVersion)
setHeaders(req.Header, "")
resp, err := client.Do(req)
if resp != nil {
defer resp.Body.Close()
}
if err != nil {
return 0, 0, err
}
io.Copy(io.Discard, resp.Body)
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusNoContent {
return 0, resp.StatusCode, fmt.Errorf("unexpected status on HEAD: %s", resp.Status)
}
offsetStr := resp.Header.Get("Upload-Offset")
offset, perr := strconv.ParseInt(offsetStr, 10, 64)
if perr != nil {
return 0, resp.StatusCode, fmt.Errorf("invalid Upload-Offset header: %q", offsetStr)
}
return offset, resp.StatusCode, nil
}
// tusPatch streams up to length bytes of the file (starting at offset) to the
// upload URL using a single PATCH request. The body is read straight from the
// *os.File, so the recording is never fully buffered in memory.
func tusPatch(client *http.Client, uploadURL string, offset, length int64, file io.Reader, setHeaders tusHeaderFunc) (int64, int, string, error) {
req, err := http.NewRequest("PATCH", uploadURL, io.LimitReader(file, length))
if err != nil {
return offset, 0, "", err
}
req.ContentLength = length
req.Header.Set("Tus-Resumable", tusResumableVersion)
req.Header.Set("Content-Type", "application/offset+octet-stream")
req.Header.Set("Upload-Offset", strconv.FormatInt(offset, 10))
setHeaders(req.Header, "")
resp, err := client.Do(req)
if resp != nil {
defer resp.Body.Close()
}
if err != nil {
return offset, 0, "", err
}
bodyBytes, _ := io.ReadAll(resp.Body)
respBody := string(bodyBytes)
if resp.StatusCode != http.StatusNoContent {
return offset, resp.StatusCode, respBody, fmt.Errorf("unexpected status on PATCH: %s, %s", resp.Status, respBody)
}
newOffsetStr := resp.Header.Get("Upload-Offset")
newOffset, perr := strconv.ParseInt(newOffsetStr, 10, 64)
if perr != nil {
// A 204 without a parseable offset means this PATCH was fully accepted.
return offset + length, resp.StatusCode, respBody, nil
}
return newOffset, resp.StatusCode, respBody, nil
}
// tusTerminate best-effort deletes an upload server-side (DELETE).
func tusTerminate(client *http.Client, uploadURL string, setHeaders tusHeaderFunc) {
req, err := http.NewRequest("DELETE", uploadURL, nil)
if err != nil {
return
}
req.Header.Set("Tus-Resumable", tusResumableVersion)
setHeaders(req.Header, "")
resp, derr := client.Do(req)
if resp != nil {
io.Copy(io.Discard, resp.Body)
resp.Body.Close()
}
_ = derr
}
// setVaultTusHeaders sets the Kerberos Vault authentication and routing headers
// on every tus request. Credentials are sent on each request (and never stored
// server-side in the upload metadata). When fileName is empty it is omitted, as
// it is only useful on the creation request (routing also travels in the tus
// Upload-Metadata).
func setVaultTusHeaders(h http.Header, vault models.KStorage, publicKey, deviceKey, fileName string) {
h.Set("X-Kerberos-Storage-CloudKey", publicKey)
h.Set("X-Kerberos-Storage-AccessKey", vault.AccessKey)
h.Set("X-Kerberos-Storage-SecretAccessKey", vault.SecretAccessKey)
h.Set("X-Kerberos-Storage-Provider", vault.Provider)
h.Set("X-Kerberos-Storage-Device", deviceKey)
h.Set("X-Kerberos-Storage-Directory", vault.Directory)
h.Set("X-Kerberos-Storage-Capture", "IPCamera")
if fileName != "" {
h.Set("X-Kerberos-Storage-FileName", fileName)
}
}
// setHubTusHeaders sets the Kerberos Hub authentication headers on every tus
// request of a hub-proxied resumable upload. The agent authenticates with its
// Hub public/private key (exactly as the legacy single-POST hub upload does);
// Kerberos Hub validates the subscription and injects the vault credentials and
// directory/provider on the agent's behalf.
func setHubTusHeaders(h http.Header, config *models.Config, fileName string) {
h.Set("X-Kerberos-Hub-PublicKey", config.HubKey)
h.Set("X-Kerberos-Hub-PrivateKey", config.HubPrivateKey)
h.Set("X-Kerberos-Hub-Region", config.S3.Region)
h.Set("X-Kerberos-Storage-Device", config.Key)
h.Set("X-Kerberos-Storage-Capture", "IPCamera")
if fileName != "" {
h.Set("X-Kerberos-Storage-FileName", fileName)
}
}
// encodeTusMetadata serializes a map into the tus Upload-Metadata header format:
// a comma separated list of "key base64(value)" pairs. Keys are sorted for a
// deterministic header value. Empty values are skipped.
func encodeTusMetadata(pairs map[string]string) string {
parts := make([]string, 0, len(pairs))
for k, v := range pairs {
if v == "" {
continue
}
parts = append(parts, k+" "+base64.StdEncoding.EncodeToString([]byte(v)))
}
sort.Strings(parts)
return strings.Join(parts, ",")
}
// resolveTusLocation turns the Location header returned by the create request
// into an absolute URL. To keep talking to the agent's configured vault host
// (and avoid issues when the vault sits behind a proxy that rewrites the host),
// it keeps the configured base URL and only appends the server-assigned upload
// id taken from the Location.
func resolveTusLocation(baseURL, location string) string {
if ref, err := url.Parse(location); err == nil {
trimmed := strings.Trim(ref.Path, "/")
if trimmed != "" {
segments := strings.Split(trimmed, "/")
id := segments[len(segments)-1]
if id != "" {
return strings.TrimRight(baseURL, "/") + "/" + id
}
}
}
// Fallback: resolve the reference against the base URL as-is.
if base, err := url.Parse(baseURL); err == nil {
if ref, err := url.Parse(location); err == nil {
return base.ResolveReference(ref).String()
}
}
return location
}
// tusSidecarDir is the directory where resume state files are kept. It is
// intentionally separate from data/cloud (which is scanned for recordings to
// upload) so the sidecar files are never mistaken for recordings.
func tusSidecarDir() string {
return "data/tus"
}
func tusSidecarPath(fileName, slot string) string {
safe := strings.ReplaceAll(fileName, "/", "_")
safe = strings.ReplaceAll(safe, string(os.PathSeparator), "_")
return filepath.Join(tusSidecarDir(), safe+"."+slot+".json")
}
// loadTusResumeState returns a previously stored upload URL for the given
// sidecar, but only if it was created against the same vault base URL. Any
// mismatch or read/parse error yields an empty string (start fresh).
func loadTusResumeState(path, baseURL string) string {
b, err := os.ReadFile(path)
if err != nil {
return ""
}
var state tusResumeState
if err := json.Unmarshal(b, &state); err != nil {
return ""
}
if state.UploadURL == "" || state.VaultURI != baseURL {
return ""
}
return state.UploadURL
}
func saveTusResumeState(path string, state tusResumeState) {
if err := os.MkdirAll(filepath.Dir(path), 0o755); err != nil {
return
}
b, err := json.Marshal(state)
if err != nil {
return
}
_ = os.WriteFile(path, b, 0o644)
}
func removeTusResumeState(path string) {
_ = os.Remove(path)
}
// tusBackoff sleeps for an exponentially increasing duration (capped) between
// resume attempts to avoid hammering a temporarily unavailable vault.
func tusBackoff(attempt int) {
delay := time.Duration(500*(1<<uint(attempt))) * time.Millisecond
if delay > 3*time.Second {
delay = 3 * time.Second
}
time.Sleep(delay)
}

View File

@@ -0,0 +1,608 @@
package cloud
import (
"bytes"
"encoding/base64"
"fmt"
"io"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strconv"
"strings"
"sync"
"testing"
"github.com/kerberos-io/agent/machinery/src/models"
)
// fakeUpload tracks the state of a single resumable upload on the fake server.
type fakeUpload struct {
size int64
offset int64
}
// recordedRequest captures the method and headers of a request received by the
// fake tus server, so tests can assert the client's per-method auth headers.
type recordedRequest struct {
method string
header http.Header
}
// fakeTus is a tiny in-memory implementation of the tus 1.0.0 server protocol,
// sufficient to exercise the agent's resumable client.
type fakeTus struct {
mu sync.Mutex
uploads map[string]*fakeUpload
counter int
creates int
lastPatchBytes int64
patchSizes []int64
// unsupported makes the creation endpoint return 404, simulating an older
// vault without a tus endpoint.
unsupported bool
// failFinalize causes the next N completing PATCH requests to return 502
// after storing the bytes, simulating a failed completion hook.
failFinalize int
// requests records the headers of every received request (in order) so
// tests can assert which auth/routing headers the client sent per method.
requests []recordedRequest
}
func newFakeTus() *fakeTus {
return &fakeTus{uploads: map[string]*fakeUpload{}}
}
func (s *fakeTus) seed(size, offset int64) string {
s.mu.Lock()
defer s.mu.Unlock()
s.counter++
id := fmt.Sprintf("seed-%d", s.counter)
s.uploads[id] = &fakeUpload{size: size, offset: offset}
return id
}
func (s *fakeTus) totalBytes() int64 {
s.mu.Lock()
defer s.mu.Unlock()
var total int64
for _, u := range s.uploads {
total += u.offset
}
return total
}
func (s *fakeTus) lastPatch() int64 {
s.mu.Lock()
defer s.mu.Unlock()
return s.lastPatchBytes
}
// patchCounts returns the number of PATCH requests received and the size of each.
func (s *fakeTus) patchCounts() (int, []int64) {
s.mu.Lock()
defer s.mu.Unlock()
sizes := make([]int64, len(s.patchSizes))
copy(sizes, s.patchSizes)
return len(s.patchSizes), sizes
}
func (s *fakeTus) createCount() int {
s.mu.Lock()
defer s.mu.Unlock()
return s.creates
}
// requestsForMethod returns the recorded requests for the given HTTP method.
func (s *fakeTus) requestsForMethod(method string) []recordedRequest {
s.mu.Lock()
defer s.mu.Unlock()
var out []recordedRequest
for _, req := range s.requests {
if req.method == method {
out = append(out, req)
}
}
return out
}
func (s *fakeTus) ServeHTTP(w http.ResponseWriter, r *http.Request) {
id := strings.TrimPrefix(r.URL.Path, tusUploadPath)
w.Header().Set("Tus-Resumable", tusResumableVersion)
s.mu.Lock()
s.requests = append(s.requests, recordedRequest{method: r.Method, header: r.Header.Clone()})
s.mu.Unlock()
switch r.Method {
case http.MethodPost:
if s.unsupported {
w.WriteHeader(http.StatusNotFound)
return
}
length, _ := strconv.ParseInt(r.Header.Get("Upload-Length"), 10, 64)
s.mu.Lock()
s.counter++
s.creates++
newID := fmt.Sprintf("up-%d", s.counter)
s.uploads[newID] = &fakeUpload{size: length}
s.mu.Unlock()
w.Header().Set("Location", tusUploadPath+newID)
w.WriteHeader(http.StatusCreated)
case http.MethodHead:
s.mu.Lock()
u, ok := s.uploads[id]
s.mu.Unlock()
if !ok {
w.WriteHeader(http.StatusNotFound)
return
}
w.Header().Set("Upload-Offset", strconv.FormatInt(u.offset, 10))
w.Header().Set("Upload-Length", strconv.FormatInt(u.size, 10))
w.WriteHeader(http.StatusOK)
case http.MethodPatch:
s.mu.Lock()
u, ok := s.uploads[id]
s.mu.Unlock()
if !ok {
w.WriteHeader(http.StatusNotFound)
return
}
n, _ := io.Copy(io.Discard, r.Body)
s.mu.Lock()
u.offset += n
s.lastPatchBytes = n
s.patchSizes = append(s.patchSizes, n)
complete := u.offset >= u.size
failNow := complete && s.failFinalize > 0
if failNow {
s.failFinalize--
}
offset := u.offset
s.mu.Unlock()
w.Header().Set("Upload-Offset", strconv.FormatInt(offset, 10))
if failNow {
// Bytes are stored but the (simulated) completion hook failed.
w.WriteHeader(http.StatusBadGateway)
return
}
w.WriteHeader(http.StatusNoContent)
case http.MethodDelete:
s.mu.Lock()
delete(s.uploads, id)
s.mu.Unlock()
w.WriteHeader(http.StatusNoContent)
default:
w.WriteHeader(http.StatusMethodNotAllowed)
}
}
// withRecording switches into a fresh temp working directory containing a
// recording at data/recordings/<fileName>. The working directory is restored on
// cleanup. Tests using this helper must not run in parallel.
func withRecording(t *testing.T, fileName string, payload []byte) {
t.Helper()
dir := t.TempDir()
old, err := os.Getwd()
if err != nil {
t.Fatalf("getwd: %v", err)
}
if err := os.Chdir(dir); err != nil {
t.Fatalf("chdir: %v", err)
}
t.Cleanup(func() { _ = os.Chdir(old) })
if err := os.MkdirAll("data/recordings", 0o755); err != nil {
t.Fatalf("mkdir recordings: %v", err)
}
if err := os.WriteFile(filepath.Join("data/recordings", fileName), payload, 0o644); err != nil {
t.Fatalf("write recording: %v", err)
}
}
func testVault(uri string) models.KStorage {
return models.KStorage{
URI: uri,
AccessKey: "ak",
SecretAccessKey: "sk",
Provider: "gcp",
Directory: "dir",
}
}
func TestUploadVaultResumable_HappyPath(t *testing.T) {
srv := newFakeTus()
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4"
payload := bytes.Repeat([]byte("x"), 4096)
withRecording(t, fileName, payload)
uploaded, responded, supported, _, err := uploadVaultResumable(testVault(ts.URL), "pk", "dev", fileName, "test", "primary")
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !uploaded || !responded || !supported {
t.Fatalf("uploaded/responded/supported = %v/%v/%v, want all true", uploaded, responded, supported)
}
if got := srv.totalBytes(); got != int64(len(payload)) {
t.Fatalf("server received %d bytes, want %d", got, len(payload))
}
if _, err := os.Stat(tusSidecarPath(fileName, "primary")); !os.IsNotExist(err) {
t.Fatalf("expected sidecar to be removed after success, stat err = %v", err)
}
}
func TestUploadVaultResumable_Chunked(t *testing.T) {
srv := newFakeTus()
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4"
// 10 KiB payload uploaded in 4 KiB chunks => 3 PATCH requests (4096+4096+2048).
payload := bytes.Repeat([]byte("c"), 10240)
withRecording(t, fileName, payload)
t.Setenv("AGENT_TUS_CHUNK_SIZE_BYTES", "4096")
uploaded, _, supported, _, err := uploadVaultResumable(testVault(ts.URL), "pk", "dev", fileName, "test", "primary")
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !uploaded || !supported {
t.Fatalf("expected chunked upload success, got uploaded=%v supported=%v", uploaded, supported)
}
if got := srv.totalBytes(); got != int64(len(payload)) {
t.Fatalf("server received %d bytes, want %d", got, len(payload))
}
count, sizes := srv.patchCounts()
if count != 3 {
t.Fatalf("expected 3 chunked PATCH requests, got %d (sizes=%v)", count, sizes)
}
want := []int64{4096, 4096, 2048}
for i, w := range want {
if sizes[i] != w {
t.Fatalf("chunk %d size = %d, want %d (sizes=%v)", i, sizes[i], w, sizes)
}
}
if _, err := os.Stat(tusSidecarPath(fileName, "primary")); !os.IsNotExist(err) {
t.Fatalf("expected sidecar removed after success, stat err = %v", err)
}
}
func TestUploadVaultResumable_ChunkingDisabled(t *testing.T) {
srv := newFakeTus()
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4"
payload := bytes.Repeat([]byte("d"), 10240)
withRecording(t, fileName, payload)
// 0 disables chunking: the whole file should go out in a single PATCH.
t.Setenv("AGENT_TUS_CHUNK_SIZE_BYTES", "0")
uploaded, _, supported, _, err := uploadVaultResumable(testVault(ts.URL), "pk", "dev", fileName, "test", "primary")
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !uploaded || !supported {
t.Fatalf("expected success, got uploaded=%v supported=%v", uploaded, supported)
}
count, sizes := srv.patchCounts()
if count != 1 {
t.Fatalf("expected a single PATCH when chunking is disabled, got %d (sizes=%v)", count, sizes)
}
if sizes[0] != int64(len(payload)) {
t.Fatalf("single PATCH size = %d, want %d", sizes[0], len(payload))
}
}
func TestTusChunkSize(t *testing.T) {
cases := []struct {
name string
env string
set bool
want int64
}{
{name: "default when unset", set: false, want: tusDefaultChunkSize},
{name: "default on invalid", env: "notanumber", set: true, want: tusDefaultChunkSize},
{name: "explicit value", env: "65536", set: true, want: 65536},
{name: "zero disables", env: "0", set: true, want: 0},
{name: "negative disables", env: "-5", set: true, want: 0},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if tc.set {
t.Setenv("AGENT_TUS_CHUNK_SIZE_BYTES", tc.env)
} else {
t.Setenv("AGENT_TUS_CHUNK_SIZE_BYTES", "")
}
if got := tusChunkSize(); got != tc.want {
t.Fatalf("tusChunkSize() = %d, want %d", got, tc.want)
}
})
}
}
func TestUploadVaultResumable_Unsupported(t *testing.T) {
srv := newFakeTus()
srv.unsupported = true
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "f.mp4"
withRecording(t, fileName, []byte("hello"))
uploaded, _, supported, _, _ := uploadVaultResumable(testVault(ts.URL), "pk", "dev", fileName, "test", "primary")
if uploaded {
t.Fatal("expected uploaded=false against a vault without a tus endpoint")
}
if supported {
t.Fatal("expected supported=false so the caller falls back to the legacy upload")
}
}
func TestUploadVaultResumable_FinalizeRetry(t *testing.T) {
srv := newFakeTus()
srv.failFinalize = 1
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4"
payload := bytes.Repeat([]byte("y"), 2048)
withRecording(t, fileName, payload)
uploaded, _, supported, _, err := uploadVaultResumable(testVault(ts.URL), "pk", "dev", fileName, "test", "primary")
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !uploaded || !supported {
t.Fatalf("expected success after a failed finalize + restart, got uploaded=%v supported=%v", uploaded, supported)
}
if got := srv.createCount(); got < 2 {
t.Fatalf("expected at least 2 create requests (restart after failed finalize), got %d", got)
}
}
func TestUploadVaultResumable_ResumeFromSidecar(t *testing.T) {
srv := newFakeTus()
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4"
total := 8192
half := 4096
payload := bytes.Repeat([]byte("z"), total)
withRecording(t, fileName, payload)
// Simulate a previous run that uploaded half the file before being interrupted.
id := srv.seed(int64(total), int64(half))
baseURL := strings.TrimRight(ts.URL, "/") + tusUploadPath
saveTusResumeState(tusSidecarPath(fileName, "primary"), tusResumeState{
UploadURL: strings.TrimRight(baseURL, "/") + "/" + id,
VaultURI: baseURL,
Size: int64(total),
})
uploaded, _, supported, _, err := uploadVaultResumable(testVault(ts.URL), "pk", "dev", fileName, "test", "primary")
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !uploaded || !supported {
t.Fatalf("expected resume success, got uploaded=%v supported=%v", uploaded, supported)
}
if got := srv.lastPatch(); got != int64(total-half) {
t.Fatalf("resume should only send the remaining %d bytes, sent %d", total-half, got)
}
if srv.createCount() != 0 {
t.Fatalf("resume should not create a new upload, got %d creates", srv.createCount())
}
}
func testHubConfig(hubURI string) *models.Config {
return &models.Config{
Key: "device-key",
HubURI: hubURI,
HubKey: "hubpub",
HubPrivateKey: "hubpriv",
S3: &models.S3{Region: "eu-west"},
}
}
// decodeTusMetadata parses a tus Upload-Metadata header value ("key b64,key b64")
// back into a map of decoded key/value pairs.
func decodeTusMetadata(meta string) map[string]string {
out := map[string]string{}
if meta == "" {
return out
}
for _, pair := range strings.Split(meta, ",") {
parts := strings.SplitN(strings.TrimSpace(pair), " ", 2)
if parts[0] == "" {
continue
}
val := ""
if len(parts) == 2 {
if b, err := base64.StdEncoding.DecodeString(parts[1]); err == nil {
val = string(b)
}
}
out[parts[0]] = val
}
return out
}
func TestUploadHubResumable_HappyPath(t *testing.T) {
srv := newFakeTus()
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "1564859471_6-474162_oprit_577-283-727-375_1153_27.mp4"
payload := bytes.Repeat([]byte("h"), 4096)
withRecording(t, fileName, payload)
uploaded, _, supported, _, err := uploadHubResumable(testHubConfig(ts.URL), fileName, "test", "hub")
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if !uploaded || !supported {
t.Fatalf("uploaded/supported = %v/%v, want both true", uploaded, supported)
}
if got := srv.totalBytes(); got != int64(len(payload)) {
t.Fatalf("server received %d bytes, want %d", got, len(payload))
}
// The Hub auth headers must be present on every request type (POST/HEAD/PATCH),
// because Kerberos Hub validates them on each proxied request. Conversely the
// vault credentials/routing are injected by Kerberos Hub on the agent's behalf
// and must never be sent by the agent on the hub path.
for _, method := range []string{http.MethodPost, http.MethodHead, http.MethodPatch} {
reqs := srv.requestsForMethod(method)
if len(reqs) == 0 {
t.Fatalf("expected at least one %s request", method)
}
for _, req := range reqs {
if got := req.header.Get("X-Kerberos-Hub-PublicKey"); got != "hubpub" {
t.Errorf("%s: X-Kerberos-Hub-PublicKey = %q, want %q", method, got, "hubpub")
}
if got := req.header.Get("X-Kerberos-Hub-PrivateKey"); got != "hubpriv" {
t.Errorf("%s: X-Kerberos-Hub-PrivateKey = %q, want %q", method, got, "hubpriv")
}
if got := req.header.Get("X-Kerberos-Hub-Region"); got != "eu-west" {
t.Errorf("%s: X-Kerberos-Hub-Region = %q, want %q", method, got, "eu-west")
}
if got := req.header.Get("X-Kerberos-Storage-Device"); got != "device-key" {
t.Errorf("%s: X-Kerberos-Storage-Device = %q, want %q", method, got, "device-key")
}
for _, h := range []string{
"X-Kerberos-Storage-AccessKey",
"X-Kerberos-Storage-SecretAccessKey",
"X-Kerberos-Storage-CloudKey",
"X-Kerberos-Storage-Provider",
"X-Kerberos-Storage-Directory",
} {
if got := req.header.Get(h); got != "" {
t.Errorf("%s: %s should be empty on the hub path, got %q", method, h, got)
}
}
}
}
// The creation request carries the upload metadata; on the hub path it must
// omit directory/provider/cloudkey (Hub resolves those) but include
// filename/device/capture. The filename header is also set on create.
posts := srv.requestsForMethod(http.MethodPost)
if got := posts[0].header.Get("X-Kerberos-Storage-FileName"); got != fileName {
t.Errorf("POST X-Kerberos-Storage-FileName = %q, want %q", got, fileName)
}
meta := decodeTusMetadata(posts[0].header.Get("Upload-Metadata"))
for _, omitted := range []string{"directory", "provider", "cloudkey"} {
if _, ok := meta[omitted]; ok {
t.Errorf("hub metadata must omit %q, got %v", omitted, meta)
}
}
if meta["filename"] != fileName {
t.Errorf("hub metadata filename = %q, want %q", meta["filename"], fileName)
}
if meta["device"] != "device-key" {
t.Errorf("hub metadata device = %q, want %q", meta["device"], "device-key")
}
if meta["capture"] != "IPCamera" {
t.Errorf("hub metadata capture = %q, want %q", meta["capture"], "IPCamera")
}
}
func TestUploadHubResumable_Unsupported(t *testing.T) {
srv := newFakeTus()
srv.unsupported = true
ts := httptest.NewServer(srv)
defer ts.Close()
fileName := "f.mp4"
withRecording(t, fileName, []byte("hello"))
uploaded, _, supported, _, _ := uploadHubResumable(testHubConfig(ts.URL), fileName, "test", "hub")
if uploaded {
t.Fatal("expected uploaded=false against a hub without a tus endpoint")
}
if supported {
t.Fatal("expected supported=false so the caller falls back to the legacy upload")
}
}
func TestEncodeTusMetadata(t *testing.T) {
got := encodeTusMetadata(map[string]string{
"b": "2",
"a": "1",
"empty": "",
})
// keys sorted, empty values skipped, values base64-encoded.
want := "a MQ==,b Mg=="
if got != want {
t.Fatalf("encodeTusMetadata = %q, want %q", got, want)
}
}
func TestResolveTusLocation(t *testing.T) {
cases := []struct {
name string
base string
location string
want string
}{
{
name: "absolute path location",
base: "http://host/storage/tus/",
location: "/storage/tus/abc",
want: "http://host/storage/tus/abc",
},
{
name: "absolute url keeps configured host",
base: "http://host/storage/tus/",
location: "http://internal:8080/storage/tus/xyz",
want: "http://host/storage/tus/xyz",
},
{
name: "relative id",
base: "http://host/api/storage/tus/",
location: "abc",
want: "http://host/api/storage/tus/abc",
},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := resolveTusLocation(tc.base, tc.location); got != tc.want {
t.Fatalf("resolveTusLocation(%q, %q) = %q, want %q", tc.base, tc.location, got, tc.want)
}
})
}
}
func TestTusResumeStateRoundTrip(t *testing.T) {
dir := t.TempDir()
old, _ := os.Getwd()
if err := os.Chdir(dir); err != nil {
t.Fatalf("chdir: %v", err)
}
defer os.Chdir(old)
path := tusSidecarPath("file.mp4", "primary")
state := tusResumeState{UploadURL: "http://host/storage/tus/abc", VaultURI: "http://host/storage/tus/", Size: 123}
saveTusResumeState(path, state)
if got := loadTusResumeState(path, state.VaultURI); got != state.UploadURL {
t.Fatalf("loadTusResumeState = %q, want %q", got, state.UploadURL)
}
// A mismatched vault URI must not be reused.
if got := loadTusResumeState(path, "http://other/storage/tus/"); got != "" {
t.Fatalf("loadTusResumeState with mismatched vault = %q, want empty", got)
}
}

View File

@@ -72,6 +72,7 @@ func Bootstrap(ctx context.Context, configDirectory string, configuration *model
communication.HandleLiveSD = make(chan int64, 1)
communication.HandleLiveHDKeepalive = make(chan string, 1)
communication.HandleLiveHDPeers = make(chan string, 1)
communication.HandleLiveHLS = make(chan int64, 1)
communication.IsConfiguring = abool.New()
cameraSettings := &models.Camera{}
@@ -304,6 +305,18 @@ func RunAgent(configDirectory string, configuration *models.Configuration, commu
go cloud.HandleLiveStreamSD(livestreamCursor, configuration, communication, mqttClient, rtspClient)
}
// Handle livestream HLS (adaptive segments over HTTP via hub-api -> vault).
// Uses the sub stream when available (lower bitrate, browser-friendly), else
// the main stream. Like SD it is viewer-keepalive gated and produces no
// traffic while nobody is watching.
if subStreamEnabled {
livestreamHLSCursor := subQueue.Latest()
go cloud.HandleLiveStreamHLS(livestreamHLSCursor, configuration, communication, mqttClient, rtspSubClient)
} else {
livestreamHLSCursor := queue.Latest()
go cloud.HandleLiveStreamHLS(livestreamHLSCursor, configuration, communication, mqttClient, rtspClient)
}
// Handle livestream HD (high resolution over WEBRTC)
communication.HandleLiveHDHandshake = make(chan models.LiveHDHandshake, 100)
if subStreamEnabled {

View File

@@ -40,6 +40,7 @@ type Communication struct {
HandleLiveHDKeepalive chan string
HandleLiveHDHandshake chan LiveHDHandshake
HandleLiveHDPeers chan string
HandleLiveHLS chan int64
HandleONVIF chan OnvifAction
IsConfiguring *abool.AtomicBool
Queue *packets.Queue

View File

@@ -173,6 +173,13 @@ type RequestSDStreamPayload struct {
Timestamp int64 `json:"timestamp"` // timestamp
}
// We received a live HLS stream request. Like SD it is a simple viewer
// keepalive: the agent owns the live HLS session, so the request only needs to
// signal "a viewer is watching" to keep the segment pipeline alive.
type RequestHLSStreamPayload struct {
Timestamp int64 `json:"timestamp"` // timestamp
}
// We received a request HD stream request
type RequestHDStreamPayload struct {
Timestamp int64 `json:"timestamp"` // timestamp

View File

@@ -344,6 +344,8 @@ func MQTTListenerHandler(mqttClient mqtt.Client, hubKey string, configDirectory
go HandleRequestSDStream(mqttClient, hubKey, payload, configuration, communication)
case "request-hd-stream":
go HandleRequestHDStream(mqttClient, hubKey, payload, configuration, communication)
case "request-hls-stream":
go HandleRequestHLSStream(mqttClient, hubKey, payload, configuration, communication)
case "receive-hd-candidates":
go HandleReceiveHDCandidates(mqttClient, hubKey, payload, configuration, communication)
case "trigger-relay":
@@ -559,6 +561,30 @@ func HandleRequestSDStream(mqttClient mqtt.Client, hubKey string, payload models
}
}
// HandleRequestHLSStream is the viewer keepalive for live HLS. Like the SD
// stream it simply signals that a viewer is watching; the agent owns the live
// HLS session, so a single non-zero timestamp on the channel keeps the segment
// pipeline alive (see cloud.HandleLiveStreamHLS). Viewers republish this
// periodically; when the keepalives stop, the agent tears the session down.
func HandleRequestHLSStream(mqttClient mqtt.Client, hubKey string, payload models.Payload, configuration *models.Configuration, communication *models.Communication) {
value := payload.Value
jsonData, _ := json.Marshal(value)
var requestHLSStreamPayload models.RequestHLSStreamPayload
json.Unmarshal(jsonData, &requestHLSStreamPayload)
if requestHLSStreamPayload.Timestamp != 0 {
if communication.CameraConnected {
select {
case communication.HandleLiveHLS <- time.Now().Unix():
default:
}
log.Log.Info("routers.mqtt.main.HandleRequestHLSStream(): received request to livestream over HLS.")
} else {
log.Log.Info("routers.mqtt.main.HandleRequestHLSStream(): received request to livestream over HLS, but camera is not connected.")
}
}
}
func HandleRequestHDStream(mqttClient mqtt.Client, hubKey string, payload models.Payload, configuration *models.Configuration, communication *models.Communication) {
value := payload.Value
// Convert map[string]interface{} to RequestHDStreamPayload

View File

@@ -0,0 +1,367 @@
package video
import (
"bytes"
"fmt"
mp4ff "github.com/Eyevinn/mp4ff/mp4"
"github.com/kerberos-io/agent/machinery/src/log"
)
// LiveSegmenter turns a live stream of Annex B video samples into HLS-ready
// fragmented-MP4 (CMAF) output: ONE init segment (ftyp+moov) followed by a
// series of INDEPENDENT media segments (styp+moof+mdat), each beginning with a
// keyframe and carrying its own tfdt. This is the building block for the live
// HLS pipeline (agent -> hub-api -> vault -> hub-frontend) and is intentionally
// kept separate from the recording muxer in mp4.go:
//
// - mp4.go writes ONE fragmented MP4 per recording (free-box placeholder up
// front, back-filled on Close). That layout is great for archived files but
// useless for live, where each segment must be shippable the instant it is
// produced and must decode on its own after the init segment.
// - LiveSegmenter emits discrete, self-contained segments via callbacks, so
// the transport (single-POST to hub-api, drop-on-failure) never has to wait
// for the recording to finish.
//
// Both producers use the SAME mp4ff fragment format, so live and archived video
// share one toolchain on the player side (hls.js #EXT-X-MAP + byte-range parts).
//
// The spike scope is video-only H.264/H.265. Audio and multi-track interleaving
// can be layered on later by adding tracks to the init segment and a second trun
// to each fragment; nothing here precludes that.
type LiveSegmenter struct {
// codec is "H264"/"H265" (case handled in buildInit).
codec string
// timescale is the media timescale used in the init segment. The agent's
// capture path feeds presentation timestamps in milliseconds, so a 1000-tick
// timescale keeps sample durations exact with no rescaling.
timescale uint32
// targetSegmentMs is the minimum amount of media a segment accumulates before
// the next keyframe is allowed to start a fresh segment. Keeping segments
// keyframe-aligned is what makes each one independently decodable.
targetSegmentMs uint64
spsNALUs [][]byte
ppsNALUs [][]byte
vpsNALUs [][]byte
// width/height are written into the visual sample entry. They are optional:
// on a successful strict SPS parse mp4ff derives them, but the manual avcC
// fallback (used for SPS that mp4ff cannot parse) needs them supplied.
width uint16
height uint16
videoTrackID uint32
initSegment *mp4ff.InitSegment
initBytes []byte
initEmitted bool
seg *mp4ff.MediaSegment
frag *mp4ff.Fragment
seqNr uint32
// started becomes true once the first segment has been opened.
started bool
// segStartPTS is the decode time (ms) of the first sample in the open
// segment; elapsed media is measured against it to decide segment cuts.
segStartPTS uint64
// segDurationMs accumulates the committed sample durations of the open
// segment so the playlist can advertise an accurate #EXTINF.
segDurationMs uint64
// pending holds the most recently received sample. Its duration is only known
// once the NEXT sample arrives (duration = nextPTS - thisPTS), mirroring the
// pending-sample pattern used by the recording muxer.
pending *mp4ff.FullSample
// lastDurationMs is the previous committed duration, reused to close out the
// final pending sample (and to bridge non-monotonic timestamps).
lastDurationMs uint64
// OnInit is invoked exactly once with the encoded init segment bytes before
// the first media segment is emitted. Optional.
OnInit func(initBytes []byte) error
// OnSegment is invoked once per completed media segment. Optional.
OnSegment func(seg LiveSegment) error
}
// LiveSegment is one independently-decodable CMAF media segment.
type LiveSegment struct {
// SequenceNumber is the monotonically increasing fragment sequence number
// (also used as the moof sequence number and the seg-N.m4s index).
SequenceNumber uint32
// DurationMs is the summed sample duration of the segment, for #EXTINF.
DurationMs uint64
// Data is the complete styp+moof+mdat segment, ready to append after the init
// segment and hand to hls.js / a vault object.
Data []byte
}
// Sample-entry flags matching the recording muxer so live and archived fragments
// describe random access points identically.
//
// keyframe 0x02000000 = sampleDependsOn=2 (depends on nothing), sync sample
// non-keyframe 0x01010000 = sampleDependsOn=1, sampleIsNonSyncSample=1
const (
liveSyncSampleFlags uint32 = 0x02000000
liveNonSyncSampleFlags uint32 = 0x01010000
// liveFallbackDurationMs is used when a duration cannot be derived (first
// frame at Close, or non-monotonic timestamps) and no prior duration exists.
// ~33 ms approximates 30 fps and is only ever a single-frame nicety.
liveFallbackDurationMs uint64 = 33
)
// NewLiveSegmenter creates a video-only live segmenter for the given codec.
// spsNALUs/ppsNALUs (and vpsNALUs for H.265) may be raw NAL units or Annex B
// blobs with start codes; both are normalized. targetSegmentMs is clamped to a
// sane floor so a misconfiguration cannot produce one-frame segments.
func NewLiveSegmenter(codec string, spsNALUs, ppsNALUs, vpsNALUs [][]byte, targetSegmentMs uint64) *LiveSegmenter {
if targetSegmentMs < 500 {
targetSegmentMs = 500
}
return &LiveSegmenter{
codec: codec,
timescale: 1000,
targetSegmentMs: targetSegmentMs,
spsNALUs: spsNALUs,
ppsNALUs: ppsNALUs,
vpsNALUs: vpsNALUs,
}
}
// SetDimensions records the encoded video width/height in pixels. They are
// written into the avc1/hvc1 visual sample entry and are required for the manual
// descriptor fallback path (SPS that mp4ff's strict parser rejects).
func (ls *LiveSegmenter) SetDimensions(width, height uint16) {
ls.width = width
ls.height = height
}
// InitSegment returns the encoded init segment bytes, building them on demand.
// Useful for tests and for serving the #EXT-X-MAP target without waiting for the
// first media segment.
func (ls *LiveSegmenter) InitSegment() ([]byte, error) {
if ls.initBytes == nil {
if err := ls.buildInit(); err != nil {
return nil, err
}
}
return ls.initBytes, nil
}
// buildInit constructs the ftyp+moov init segment from the parameter sets.
func (ls *LiveSegmenter) buildInit() error {
init := mp4ff.CreateEmptyInit()
init.AddEmptyTrack(ls.timescale, "video", "und")
trak := init.Moov.Traks[0]
switch ls.codec {
case "H264", "h264", "AVC", "avc", "AVC1", "avc1":
sps, pps := normalizeH264ParameterSets(ls.spsNALUs, ls.ppsNALUs)
if len(sps) == 0 || len(pps) == 0 {
return fmt.Errorf("livehls: missing H264 SPS/PPS (sps=%d pps=%d)", len(sps), len(pps))
}
// includePS=true stores SPS/PPS in the avcC so segments need not carry
// in-band parameter sets - browsers read them from the init segment. Some
// camera SPS variants trip mp4ff's strict parser (e.g. unusual VUI/SAR);
// fall back to a manually built avcC just like the recording muxer does so
// those cameras still produce a valid init segment.
if err := trak.SetAVCDescriptor("avc1", sps, pps, true); err != nil {
log.Log.Warning("livehls: SetAVCDescriptor failed, using manual avcC fallback: " + err.Error())
if fbErr := addAVCDescriptorFallback(trak, sps, pps, ls.width, ls.height); fbErr != nil {
return fmt.Errorf("livehls: AVC descriptor fallback: %w", fbErr)
}
}
case "H265", "h265", "HEVC", "hevc", "HVC1", "hvc1":
vps, sps, pps := normalizeH265ParameterSets(ls.vpsNALUs, ls.spsNALUs, ls.ppsNALUs)
if len(vps) == 0 || len(sps) == 0 || len(pps) == 0 {
return fmt.Errorf("livehls: missing H265 VPS/SPS/PPS (vps=%d sps=%d pps=%d)", len(vps), len(sps), len(pps))
}
if err := trak.SetHEVCDescriptor("hvc1", vps, sps, pps, [][]byte{}, true); err != nil {
return fmt.Errorf("livehls: SetHEVCDescriptor: %w", err)
}
default:
return fmt.Errorf("livehls: unsupported codec %q", ls.codec)
}
// Record the encoded dimensions in the track header when known.
if ls.width > 0 && ls.height > 0 {
trak.Tkhd.Width = mp4ff.Fixed32(uint32(ls.width) << 16)
trak.Tkhd.Height = mp4ff.Fixed32(uint32(ls.height) << 16)
}
// mdhd.Duration MUST be 0 for fragmented MP4 so players derive duration from
// the fragments rather than a (here unknown) total.
trak.Mdia.Mdhd.Duration = 0
ls.videoTrackID = trak.Tkhd.TrackID
var buf bytes.Buffer
if err := init.Encode(&buf); err != nil {
return fmt.Errorf("livehls: encode init: %w", err)
}
ls.initSegment = init
ls.initBytes = buf.Bytes()
return nil
}
// WriteSample feeds one Annex B access unit with its decode timestamp (DTS) in
// milliseconds. The first sample of a session MUST be a keyframe; a non-keyframe
// first sample is dropped (it could not be decoded without a preceding IDR).
//
// compositionOffsetMs is the CTS offset (PTS-DTS, for B-frame reordering) in
// timescale ticks; pass 0 for streams without B-frames.
func (ls *LiveSegmenter) WriteSample(isKeyframe bool, annexB []byte, ptsMs uint64, compositionOffsetMs int32) error {
// Lazily build + emit the init segment on the first accepted sample.
if ls.initBytes == nil {
if err := ls.buildInit(); err != nil {
return err
}
}
if !ls.initEmitted {
ls.initEmitted = true
if ls.OnInit != nil {
if err := ls.OnInit(ls.initBytes); err != nil {
return err
}
}
}
// A session must open on a random-access point; otherwise the first segment
// would reference frames that never arrived.
if !ls.started && !isKeyframe {
log.Log.Debug("LiveSegmenter.WriteSample(): dropping leading non-keyframe before first IDR")
return nil
}
lengthPrefixed, err := annexBToLengthPrefixed(annexB)
if err != nil {
return fmt.Errorf("livehls: convert AnnexB: %w", err)
}
// The previous sample's duration is the gap to this sample's PTS. Commit it
// to the (still open) current fragment before we consider rolling segments,
// because the pending sample always precedes this one in decode order.
if ls.pending != nil {
dur := ls.lastDurationMs
if ptsMs > ls.pending.DecodeTime {
dur = ptsMs - ls.pending.DecodeTime
}
if dur == 0 {
dur = liveFallbackDurationMs
}
ls.lastDurationMs = dur
ls.pending.Sample.Dur = uint32(dur)
if err := ls.commitPending(); err != nil {
return err
}
}
// At every keyframe, decide whether enough media has accumulated to close the
// open segment and start a new one. Cutting only on keyframes guarantees each
// segment is independently decodable.
if isKeyframe {
shouldCut := !ls.started || (ptsMs-ls.segStartPTS) >= ls.targetSegmentMs
if shouldCut {
if ls.started {
if err := ls.emitSegment(); err != nil {
return err
}
}
ls.openSegment(ptsMs)
}
}
// Stage this sample; its duration is filled in when the next sample arrives
// (or at Close()).
flags := liveNonSyncSampleFlags
if isKeyframe {
flags = liveSyncSampleFlags
}
ls.pending = &mp4ff.FullSample{
Sample: mp4ff.Sample{
Flags: flags,
Size: uint32(len(lengthPrefixed)),
CompositionTimeOffset: compositionOffsetMs,
},
DecodeTime: ptsMs,
Data: lengthPrefixed,
}
return nil
}
// openSegment starts a fresh media segment (with CMAF styp) and an empty
// single-track fragment whose moof sequence number is the segment index.
func (ls *LiveSegmenter) openSegment(startPTS uint64) {
ls.seqNr++
ls.seg = mp4ff.NewMediaSegment() // includes a CMAF styp box by default
frag, err := mp4ff.CreateFragment(ls.seqNr, ls.videoTrackID)
if err != nil {
log.Log.Error("LiveSegmenter.openSegment(): CreateFragment failed: " + err.Error())
return
}
ls.seg.AddFragment(frag)
ls.frag = frag
ls.segStartPTS = startPTS
ls.segDurationMs = 0
ls.started = true
}
// commitPending appends the staged sample to the open fragment. The first sample
// of a fragment seeds the tfdt baseMediaDecodeTime from its absolute DecodeTime,
// which is what makes the segment independently seekable/decodable.
func (ls *LiveSegmenter) commitPending() error {
if ls.pending == nil {
return nil
}
if ls.frag == nil {
// No open segment yet (e.g. pending set before the first keyframe cut). The
// keyframe path always opens a segment before staging, so this only guards
// against logic drift; drop rather than panic.
ls.pending = nil
return nil
}
if err := ls.frag.AddFullSampleToTrack(*ls.pending, ls.videoTrackID); err != nil {
return fmt.Errorf("livehls: AddFullSampleToTrack: %w", err)
}
ls.segDurationMs += uint64(ls.pending.Sample.Dur)
ls.pending = nil
return nil
}
// emitSegment encodes the open segment and hands it to OnSegment.
func (ls *LiveSegmenter) emitSegment() error {
if ls.seg == nil {
return nil
}
var buf bytes.Buffer
if err := ls.seg.Encode(&buf); err != nil {
return fmt.Errorf("livehls: encode segment %d: %w", ls.seqNr, err)
}
out := LiveSegment{
SequenceNumber: ls.seqNr,
DurationMs: ls.segDurationMs,
Data: buf.Bytes(),
}
ls.seg = nil
ls.frag = nil
if ls.OnSegment != nil {
return ls.OnSegment(out)
}
return nil
}
// Close flushes the final pending sample and emits the last open segment. Call
// once when the live session ends so no trailing media is lost.
func (ls *LiveSegmenter) Close() error {
if ls.pending != nil {
dur := ls.lastDurationMs
if dur == 0 {
dur = liveFallbackDurationMs
}
ls.pending.Sample.Dur = uint32(dur)
if err := ls.commitPending(); err != nil {
return err
}
}
return ls.emitSegment()
}

View File

@@ -0,0 +1,371 @@
package video
import (
"bytes"
"fmt"
"math"
"os"
"path/filepath"
"strings"
"testing"
mp4ff "github.com/Eyevinn/mp4ff/mp4"
)
// Known-good minimal H.264 baseline parameter sets (640x480), reused from the
// recording-muxer tests so the live segmenter is exercised against the exact
// SPS/PPS mp4ff is already known to parse into an avcC descriptor.
var (
liveTestSPS = []byte{0x67, 0x42, 0xc0, 0x1e, 0xd9, 0x00, 0xa0, 0x47, 0xfe, 0xc8}
liveTestPPS = []byte{0x68, 0xce, 0x38, 0x80}
)
// makeAnnexBFrame builds a single-NALU Annex B access unit: a 4-byte start code,
// the NAL header (IDR=0x65 for keyframes, non-IDR=0x01 otherwise) and padding.
func makeAnnexBFrame(isKey bool) []byte {
nalType := byte(0x01)
if isKey {
nalType = 0x65
}
frame := []byte{0x00, 0x00, 0x00, 0x01, nalType}
for i := 0; i < 100; i++ {
frame = append(frame, byte(i))
}
return frame
}
// isSyncSample reports whether a parsed sample is a random-access point
// (sample_depends_on == 2 => "depends on nothing" => IDR/sync).
func isSyncSample(s mp4ff.Sample) bool {
return (s.Flags>>24)&0x03 == 0x02
}
// TestLiveSegmenterProducesIndependentCMAFSegments feeds a synthetic H.264
// stream (25 fps, 1s GOPs) through the live segmenter and asserts that:
// - exactly one init segment (ftyp+moov, single avc1 video track) is produced;
// - segments are cut on keyframe boundaries honoring the target duration;
// - every media segment carries a CMAF styp + exactly one moof+mdat fragment;
// - each segment begins with a sync sample and its tfdt equals the absolute
// decode time of that first sample (the property that makes it independently
// decodable after the init segment);
// - sample counts and durations are preserved end to end.
func TestLiveSegmenterProducesIndependentCMAFSegments(t *testing.T) {
const (
frameDurMs = uint64(40) // 25 fps
gopFrames = 25 // keyframe every 1000 ms
numGOPs = 6
numFrames = gopFrames * numGOPs // 150 frames, 6000 ms
targetMs = uint64(2000) // 2s segments => 2 GOPs each
)
seg := NewLiveSegmenter("H264", [][]byte{liveTestSPS}, [][]byte{liveTestPPS}, nil, targetMs)
seg.SetDimensions(640, 480)
var initBytes []byte
var initCalls int
var segments []LiveSegment
seg.OnInit = func(b []byte) error {
initCalls++
initBytes = append([]byte(nil), b...)
return nil
}
seg.OnSegment = func(s LiveSegment) error {
segments = append(segments, s)
return nil
}
for i := 0; i < numFrames; i++ {
isKey := i%gopFrames == 0
pts := uint64(i) * frameDurMs
if err := seg.WriteSample(isKey, makeAnnexBFrame(isKey), pts, 0); err != nil {
t.Fatalf("WriteSample(frame=%d): %v", i, err)
}
}
if err := seg.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
// --- Init segment: emitted exactly once, well-formed, single video track. ---
if initCalls != 1 {
t.Fatalf("OnInit called %d times, want 1", initCalls)
}
if len(initBytes) == 0 {
t.Fatal("init segment is empty")
}
parsedInit, err := mp4ff.DecodeFile(bytes.NewReader(initBytes))
if err != nil {
t.Fatalf("decode init: %v", err)
}
if parsedInit.Init == nil || parsedInit.Init.Ftyp == nil || parsedInit.Init.Moov == nil {
t.Fatal("init segment missing ftyp/moov")
}
if got := len(parsedInit.Init.Moov.Traks); got != 1 {
t.Fatalf("init moov has %d traks, want 1", got)
}
// --- Segment cut cadence: 6 GOPs at 2s target => 3 segments of 2 GOPs each. ---
const wantSegments = 3
if len(segments) != wantSegments {
t.Fatalf("got %d media segments, want %d", len(segments), wantSegments)
}
for i, s := range segments {
if want := uint32(i + 1); s.SequenceNumber != want {
t.Errorf("segment %d: SequenceNumber=%d, want %d", i, s.SequenceNumber, want)
}
if s.DurationMs != targetMs {
t.Errorf("segment %d: DurationMs=%d, want %d", i, s.DurationMs, targetMs)
}
}
// --- Each segment must decode INDEPENDENTLY after the init segment. ---
// Parsing init+oneSegment in isolation mirrors exactly what hls.js does with
// an #EXT-X-MAP init and a single media part.
var totalSamples, totalSync int
wantTFDT := []uint64{0, 2000, 4000}
for i, s := range segments {
standalone := append(append([]byte(nil), initBytes...), s.Data...)
parsed, err := mp4ff.DecodeFile(bytes.NewReader(standalone))
if err != nil {
t.Fatalf("segment %d: decode init+segment: %v", i, err)
}
if len(parsed.Segments) != 1 {
t.Fatalf("segment %d: parsed %d media segments, want 1", i, len(parsed.Segments))
}
mseg := parsed.Segments[0]
if mseg.Styp == nil {
t.Errorf("segment %d: missing CMAF styp box", i)
}
if len(mseg.Fragments) != 1 {
t.Fatalf("segment %d: %d fragments, want 1", i, len(mseg.Fragments))
}
fr := mseg.Fragments[0]
if got := fr.Moof.Mfhd.SequenceNumber; got != s.SequenceNumber {
t.Errorf("segment %d: moof sequence=%d, want %d", i, got, s.SequenceNumber)
}
traf := fr.Moof.Traf
if traf.Tfhd.TrackID != 1 {
t.Errorf("segment %d: track id=%d, want 1", i, traf.Tfhd.TrackID)
}
if got := traf.Tfdt.BaseMediaDecodeTime(); got != wantTFDT[i] {
t.Errorf("segment %d: tfdt baseMediaDecodeTime=%d, want %d", i, got, wantTFDT[i])
}
var samples []mp4ff.Sample
for _, trun := range traf.Truns {
samples = append(samples, trun.Samples...)
}
if len(samples) == 0 {
t.Fatalf("segment %d: no samples", i)
}
if !isSyncSample(samples[0]) {
t.Errorf("segment %d: first sample is not a keyframe/sync sample", i)
}
var segDur uint64
for j, smp := range samples {
totalSamples++
if isSyncSample(smp) {
totalSync++
}
segDur += uint64(smp.Dur)
if smp.Size == 0 {
t.Errorf("segment %d sample %d: zero size", i, j)
}
}
if segDur != s.DurationMs {
t.Errorf("segment %d: summed sample dur=%d, reported DurationMs=%d", i, segDur, s.DurationMs)
}
}
if totalSamples != numFrames {
t.Errorf("total samples across segments=%d, want %d", totalSamples, numFrames)
}
if totalSync != numGOPs {
t.Errorf("total sync samples=%d, want %d (one per GOP)", totalSync, numGOPs)
}
}
// TestLiveSegmenterDropsLeadingNonKeyframe verifies a session cannot open on a
// non-IDR frame (which would reference frames that never arrived); such leading
// samples are dropped until the first keyframe.
func TestLiveSegmenterDropsLeadingNonKeyframe(t *testing.T) {
seg := NewLiveSegmenter("H264", [][]byte{liveTestSPS}, [][]byte{liveTestPPS}, nil, 1000)
seg.SetDimensions(640, 480)
var segments []LiveSegment
seg.OnSegment = func(s LiveSegment) error { segments = append(segments, s); return nil }
// Two P-frames before any IDR must be ignored.
if err := seg.WriteSample(false, makeAnnexBFrame(false), 0, 0); err != nil {
t.Fatalf("WriteSample(p0): %v", err)
}
if err := seg.WriteSample(false, makeAnnexBFrame(false), 40, 0); err != nil {
t.Fatalf("WriteSample(p1): %v", err)
}
// First IDR opens the session at decode time 0.
for i := 0; i < 25; i++ {
isKey := i == 0
if err := seg.WriteSample(isKey, makeAnnexBFrame(isKey), uint64(i)*40, 0); err != nil {
t.Fatalf("WriteSample(%d): %v", i, err)
}
}
if err := seg.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
if len(segments) == 0 {
t.Fatal("expected at least one segment after the first IDR")
}
initBytes, err := seg.InitSegment()
if err != nil {
t.Fatalf("InitSegment: %v", err)
}
standalone := append(append([]byte(nil), initBytes...), segments[0].Data...)
parsed, err := mp4ff.DecodeFile(bytes.NewReader(standalone))
if err != nil {
t.Fatalf("decode: %v", err)
}
traf := parsed.Segments[0].Fragments[0].Moof.Traf
if got := traf.Tfdt.BaseMediaDecodeTime(); got != 0 {
t.Errorf("first segment tfdt=%d, want 0 (session opens on the IDR)", got)
}
var first mp4ff.Sample
for _, trun := range traf.Truns {
if len(trun.Samples) > 0 {
first = trun.Samples[0]
break
}
}
if !isSyncSample(first) {
t.Error("first committed sample must be the IDR, not a dropped P-frame")
}
}
// renderLiveMediaPlaylist renders a live (no #EXT-X-ENDLIST) fMP4 HLS media
// playlist for the given segments. This mirrors the shape hub-api will serve for
// live streams: an #EXT-X-MAP init segment followed by one #EXTINF per CMAF part.
// In production hub-api emits a sliding WINDOW of the most recent segments and
// advances #EXT-X-MEDIA-SEQUENCE; here we list the whole synthetic capture for a
// self-contained, inspectable bundle.
func renderLiveMediaPlaylist(initURI string, segs []LiveSegment, mediaSequence uint32) string {
var maxDurMs uint64
for _, s := range segs {
if s.DurationMs > maxDurMs {
maxDurMs = s.DurationMs
}
}
target := uint64(math.Ceil(float64(maxDurMs) / 1000.0))
if target == 0 {
target = 1
}
var b strings.Builder
b.WriteString("#EXTM3U\n")
b.WriteString("#EXT-X-VERSION:7\n")
fmt.Fprintf(&b, "#EXT-X-TARGETDURATION:%d\n", target)
fmt.Fprintf(&b, "#EXT-X-MEDIA-SEQUENCE:%d\n", mediaSequence)
b.WriteString("#EXT-X-INDEPENDENT-SEGMENTS\n")
fmt.Fprintf(&b, "#EXT-X-MAP:URI=%q\n", initURI)
for _, s := range segs {
fmt.Fprintf(&b, "#EXTINF:%.3f,\n", float64(s.DurationMs)/1000.0)
fmt.Fprintf(&b, "seg-%d.m4s\n", s.SequenceNumber)
}
// NOTE: deliberately no #EXT-X-ENDLIST - its absence is what marks the
// playlist as live so hls.js keeps polling for new segments.
return b.String()
}
// TestLiveSegmenterWritesHLSBundle runs the segmenter over a synthetic stream and
// writes a complete on-disk fMP4 HLS bundle (init.mp4 + seg-N.m4s + a live
// stream.m3u8). It validates the playlist shape and that every referenced file
// exists, then logs the output directory so the structure can be eyeballed.
//
// Set LIVEHLS_OUT=/some/dir to keep the bundle for manual inspection (e.g. serve
// it and point hls.js at stream.m3u8); otherwise a temp dir is used and removed.
//
// The frames here are synthetic (valid fMP4 boxing, non-decodable payloads), so
// this validates CONTAINER/playlist structure, not pixel decode - the round-trip
// assertions in TestLiveSegmenterProducesIndependentCMAFSegments cover decodable
// box layout.
func TestLiveSegmenterWritesHLSBundle(t *testing.T) {
const (
frameDurMs = uint64(40)
gopFrames = 25
numGOPs = 6
numFrames = gopFrames * numGOPs
targetMs = uint64(2000)
)
outDir := os.Getenv("LIVEHLS_OUT")
if outDir == "" {
outDir = t.TempDir()
} else {
if err := os.MkdirAll(outDir, 0o755); err != nil {
t.Fatalf("mkdir %s: %v", outDir, err)
}
}
seg := NewLiveSegmenter("H264", [][]byte{liveTestSPS}, [][]byte{liveTestPPS}, nil, targetMs)
seg.SetDimensions(640, 480)
var segments []LiveSegment
seg.OnInit = func(b []byte) error {
return os.WriteFile(filepath.Join(outDir, "init.mp4"), b, 0o644)
}
seg.OnSegment = func(s LiveSegment) error {
segments = append(segments, s)
name := fmt.Sprintf("seg-%d.m4s", s.SequenceNumber)
return os.WriteFile(filepath.Join(outDir, name), s.Data, 0o644)
}
for i := 0; i < numFrames; i++ {
isKey := i%gopFrames == 0
if err := seg.WriteSample(isKey, makeAnnexBFrame(isKey), uint64(i)*frameDurMs, 0); err != nil {
t.Fatalf("WriteSample(%d): %v", i, err)
}
}
if err := seg.Close(); err != nil {
t.Fatalf("Close: %v", err)
}
if len(segments) == 0 {
t.Fatal("no segments produced")
}
playlist := renderLiveMediaPlaylist("init.mp4", segments, segments[0].SequenceNumber)
if err := os.WriteFile(filepath.Join(outDir, "stream.m3u8"), []byte(playlist), 0o644); err != nil {
t.Fatalf("write playlist: %v", err)
}
// --- Validate the live playlist shape. ---
mustContain := []string{
"#EXTM3U",
"#EXT-X-VERSION:7",
"#EXT-X-TARGETDURATION:2",
"#EXT-X-MEDIA-SEQUENCE:1",
`#EXT-X-MAP:URI="init.mp4"`,
"#EXT-X-INDEPENDENT-SEGMENTS",
}
for _, tag := range mustContain {
if !strings.Contains(playlist, tag) {
t.Errorf("playlist missing %q\n---\n%s", tag, playlist)
}
}
if strings.Contains(playlist, "#EXT-X-ENDLIST") {
t.Error("live playlist must NOT contain #EXT-X-ENDLIST")
}
if got, want := strings.Count(playlist, "#EXTINF:"), len(segments); got != want {
t.Errorf("playlist has %d #EXTINF entries, want %d", got, want)
}
// --- Every referenced file must exist on disk. ---
if _, err := os.Stat(filepath.Join(outDir, "init.mp4")); err != nil {
t.Errorf("init.mp4 missing: %v", err)
}
for _, s := range segments {
name := fmt.Sprintf("seg-%d.m4s", s.SequenceNumber)
if _, err := os.Stat(filepath.Join(outDir, name)); err != nil {
t.Errorf("%s missing: %v", name, err)
}
}
t.Logf("wrote HLS bundle to %s (%d segments)\n%s", outDir, len(segments), playlist)
}

View File

@@ -33,13 +33,18 @@ const MacEpochOffset uint64 = 2082844800
const FragmentDurationMs = 3000
// SeamGapDivisor controls loop-seam detection. A keyframe is treated as an
// upstream loop/restart seam when it arrives in less than (previous keyframe
// interval / SeamGapDivisor) — i.e. far sooner than the established keyframe
// cadence. Comparing against the *previous* interval (rather than a fixed
// millisecond threshold) makes the check scale automatically with the camera's
// configured GOP size: it works the same whether keyframes are 0.5s, 1s, 2s or
// more apart, and does not misfire on legitimately short-GOP or all-intra
// streams (where every interval is similar, so none looks anomalously short).
// upstream loop/restart seam when it arrives in less than (smallest normal
// keyframe interval / SeamGapDivisor) — i.e. far sooner than the camera's
// tightest established keyframe cadence.
//
// The reference is the running *minimum* keyframe interval, NOT the immediately
// preceding one. Variable-GOP ("smart codec") cameras lengthen the GOP during
// static scenes and shorten it again on motion, so consecutive intervals differ
// wildly (e.g. 2000 ms then 500 ms). Comparing against the previous interval
// then flags every normal short GOP that happens to follow a long static GOP as
// a seam and drops healthy video. Comparing against the minimum cadence instead
// scales with any configured GOP size (0.5s, 1s, 2s, ...) yet never mistakes the
// camera's own normal cadence for a premature seam IDR.
const SeamGapDivisor = 2
type MP4 struct {
@@ -85,7 +90,8 @@ type MP4 struct {
FragmentKeyframeCount int // Keyframes in the current fragment
PendingSampleIsKeyframe bool // Whether the pending video sample is a keyframe
LastKeyframeRawPTS uint64 // Raw PTS of the most recently seen keyframe (across fragments)
LastKeyframeGapMs uint64 // Interval (ms) between the two most recent keyframes; reference cadence for seam detection
LastKeyframeGapMs uint64 // Interval (ms) between the two most recent keyframes (diagnostic only)
MinKeyframeGapMs uint64 // Smallest keyframe interval (ms) seen so far; the camera's tightest cadence and the reference for seam detection
gopBuffer []bufferedSample // Current, not-yet-committed GOP (video frames + interleaved audio), held so a loop-seam GOP can be dropped before it reaches the file
}
@@ -333,23 +339,40 @@ func (mp4 *MP4) AddSampleToTrack(trackID uint32, isKeyframe bool, data []byte, p
// buffered GOP is genuine (commit it) or the truncated tail GOP at an upstream
// loop/restart seam (drop it).
//
// The GOP size is configurable per camera, so we do NOT compare against a
// fixed millisecond threshold. Instead we compare this keyframe interval to
// the previous one and only flag a *sudden* shortening: a seam IDR arrives in
// less than (previous interval / SeamGapDivisor). Deriving the threshold from
// the observed cadence keeps detection correct for any configured GOP (0.5s,
// 1s, 2s, ...) and avoids false positives on steady short-GOP / all-intra
// streams (where consecutive intervals are similar, so none looks anomalously
// short). Because the reference is the immediately preceding interval, a burst
// of close keyframes only drops a single GOP instead of cascading.
// A genuine loop/restart seam has TWO signatures that must BOTH hold; we never
// drop a GOP on the interval alone, because variable-GOP ("smart codec")
// cameras legitimately shorten the GOP on motion:
//
// 1. The new keyframe arrives much sooner than the camera's tightest normal
// cadence: gap*SeamGapDivisor < MinKeyframeGapMs (the running MINIMUM
// interval). Using the minimum — not the previous interval — means a
// normal short GOP that merely follows a long static GOP (2000 ms -> 500 ms)
// is NOT flagged, while a true premature restart still is.
// 2. The GOP we just buffered is actually TRUNCATED — far shorter than a full
// GOP. A real seam cuts a GOP off mid-stream, leaving only a handful of
// frames; a healthy GOP (even a legitimately short one) is left intact and
// must be committed in full. We require the buffered tail to be under half
// the minimum normal GOP length to qualify as truncated.
//
// Deriving both thresholds from the observed cadence keeps detection correct
// for any configured GOP size (0.5s, 1s, 2s, ...) and stops the heuristic from
// discarding healthy video.
seam := false
if mp4.LastKeyframeRawPTS > 0 && pts > mp4.LastKeyframeRawPTS {
gap := pts - mp4.LastKeyframeRawPTS
if mp4.LastKeyframeGapMs > 0 && gap*SeamGapDivisor < mp4.LastKeyframeGapMs {
bufferedVideo := mp4.bufferedVideoCount()
// Frames a full GOP at the tightest normal cadence would contain.
fullGopFrames := mp4.expectedGopFrames(gap)
closeKeyframe := mp4.MinKeyframeGapMs > 0 && gap*SeamGapDivisor < mp4.MinKeyframeGapMs
truncatedTail := fullGopFrames > 0 && bufferedVideo*2 < fullGopFrames
if closeKeyframe && truncatedTail {
seam = true
log.Log.Warning(fmt.Sprintf("mp4.AddSampleToTrack(): dropping truncated GOP at unexpectedly close keyframe (interval=%d ms, previous interval=%d ms, buffered samples=%d) - likely upstream loop/restart discontinuity", gap, mp4.LastKeyframeGapMs, len(mp4.gopBuffer)))
log.Log.Warning(fmt.Sprintf("mp4.AddSampleToTrack(): dropping truncated GOP at premature keyframe (interval=%d ms, min interval=%d ms, buffered video frames=%d of ~%d) - likely upstream loop/restart discontinuity", gap, mp4.MinKeyframeGapMs, bufferedVideo, fullGopFrames))
}
mp4.LastKeyframeGapMs = gap
if !seam && (mp4.MinKeyframeGapMs == 0 || gap < mp4.MinKeyframeGapMs) {
mp4.MinKeyframeGapMs = gap
}
}
mp4.LastKeyframeRawPTS = pts
@@ -372,6 +395,70 @@ func (mp4 *MP4) AddSampleToTrack(trackID uint32, isKeyframe bool, data []byte, p
return nil
}
// bufferedVideoCount returns how many video-track samples are currently held in
// the GOP buffer (interleaved audio samples are ignored). It measures how
// complete the buffered GOP is, used to tell a truncated seam tail from a
// healthy — possibly legitimately short — GOP.
func (mp4 *MP4) bufferedVideoCount() uint64 {
var n uint64
for _, s := range mp4.gopBuffer {
if s.trackID == uint32(mp4.VideoTrack) {
n++
}
}
return n
}
// expectedGopFrames estimates how many video frames a full GOP at the camera's
// tightest normal cadence (MinKeyframeGapMs) would contain, using the video
// frame interval inferred from the buffered GOP. gap is the current keyframe
// interval, used as a fallback frame-duration source. Returns 0 when there is
// not yet enough information to judge (so callers must not treat a GOP as
// truncated without a reliable estimate).
func (mp4 *MP4) expectedGopFrames(gap uint64) uint64 {
cadence := mp4.MinKeyframeGapMs
if cadence == 0 {
return 0
}
frameDur := mp4.bufferedVideoFrameDuration()
if frameDur == 0 {
// Fall back to deriving a per-frame duration from the buffered tail across
// the current interval; if that is unavailable too, we cannot estimate.
if n := mp4.bufferedVideoCount(); n > 0 && gap > 0 {
frameDur = gap / n
}
}
if frameDur == 0 {
return 0
}
return cadence / frameDur
}
// bufferedVideoFrameDuration returns the average per-frame duration (in PTS
// units) of the video samples currently buffered, derived from the PTS deltas
// between consecutive video frames. Returns 0 when fewer than two video frames
// are buffered.
func (mp4 *MP4) bufferedVideoFrameDuration() uint64 {
var prev uint64
havePrev := false
var sum, count uint64
for _, s := range mp4.gopBuffer {
if s.trackID != uint32(mp4.VideoTrack) {
continue
}
if havePrev && s.pts > prev {
sum += s.pts - prev
count++
}
prev = s.pts
havePrev = true
}
if count == 0 {
return 0
}
return sum / count
}
// commitBufferedGOP writes every sample currently held in gopBuffer to the file
// in arrival order, then clears the buffer. Committing in arrival order
// preserves the original audio/video interleave and lets commitSampleToTrack's

View File

@@ -0,0 +1,129 @@
package video
import (
"os"
"testing"
mp4ff "github.com/Eyevinn/mp4ff/mp4"
"github.com/kerberos-io/agent/machinery/src/models"
)
// TestMP4VariableGOPKeepsHealthyShortGOP reproduces the adam-drive regression:
// a variable-GOP ("smart codec") camera lengthens its keyframe interval during a
// static scene (e.g. 500ms -> 1500/2000ms) and then drops back to its normal
// 500ms cadence on motion. That normal, FULL 500ms GOP arrives much sooner than
// the immediately preceding (long, static) GOP.
//
// The previous heuristic compared the new keyframe interval against the *previous*
// interval and dropped the GOP whenever gap < previousInterval/2 — so every normal
// 500ms keyframe following a long static GOP was misclassified as a premature
// loop/restart seam and a whole healthy GOP (~15 frames) was discarded. In the
// field this silently deleted ~0.5s of video on virtually every recording from
// such cameras, producing a freeze/jump artifact.
//
// After the fix the seam check compares against the running MINIMUM cadence and
// additionally requires the buffered GOP to be genuinely truncated, so a full
// healthy GOP is always kept regardless of how long the preceding GOP was. This
// test asserts that NO frames are dropped for a pure variable-GOP stream.
func TestMP4VariableGOPKeepsHealthyShortGOP(t *testing.T) {
tmpFile, err := os.CreateTemp("", "test_variable_gop_*.mp4")
if err != nil {
t.Fatalf("create temp: %v", err)
}
tmpFile.Close()
defer os.Remove(tmpFile.Name())
sps := []byte{0x67, 0x42, 0xc0, 0x1e, 0xd9, 0x00, 0xa0, 0x47, 0xfe, 0xc8}
pps := []byte{0x68, 0xce, 0x38, 0x80}
mp4Video := NewMP4(tmpFile.Name(), [][]byte{sps}, [][]byte{pps}, nil, 60)
mp4Video.SetWidth(1920)
mp4Video.SetHeight(1080)
v := mp4Video.AddVideoTrack("H264")
mk := func(k bool) []byte {
nt := byte(0x01)
if k {
nt = 0x65
}
f := []byte{0, 0, 0, 1, nt}
for i := 0; i < 200; i++ {
f = append(f, byte(i))
}
return f
}
const frameDur = uint64(33)
pts := uint64(0)
emitFrame := func(isKey bool) {
mp4Video.AddSampleToTrack(v, isKey, mk(isKey), pts, 0)
pts += frameDur
}
// emitGOP emits a complete GOP of exactly frames frames: a leading keyframe
// followed by frames-1 P-frames. Every GOP here is healthy and complete; only
// its length varies, exactly as a smart-codec camera varies the GOP.
emitGOP := func(frames int) {
emitFrame(true)
for i := 0; i < frames-1; i++ {
emitFrame(false)
}
}
// Normal cadence is 15 frames (~500ms). The camera then lengthens the GOP for
// several static scenes (45 and 60 frames, ~1500ms and ~2000ms) before
// dropping back to the normal 15-frame GOP on motion — the transition the old
// heuristic wrongly treated as a seam. The whole sequence is then repeated to
// cover multiple long->short transitions.
gopLengths := []int{15, 15, 45, 15, 60, 15, 15, 45, 15, 15, 60, 15}
totalEmittedFrames := 0
emittedKeyframes := 0
for _, n := range gopLengths {
emitGOP(n)
totalEmittedFrames += n
emittedKeyframes++
}
mp4Video.Close(&models.Config{Signing: &models.Signing{PrivateKey: ""}})
f, err := os.Open(tmpFile.Name())
if err != nil {
t.Fatalf("open: %v", err)
}
defer f.Close()
parsed, err := mp4ff.DecodeFile(f)
if err != nil {
t.Fatalf("decode: %v", err)
}
totalSamples := 0
totalSync := 0
for _, seg := range parsed.Segments {
for _, fr := range seg.Fragments {
for _, traf := range fr.Moof.Trafs {
if traf.Tfhd.TrackID != 1 {
continue
}
for _, trun := range traf.Truns {
for _, s := range trun.Samples {
totalSamples++
// sample_depends_on == 2 => "does not depend on others" => IDR/sync.
if (s.Flags>>24)&0x03 == 0x02 {
totalSync++
}
}
}
}
}
}
// Every GOP is healthy, so nothing must be dropped: all keyframes and all
// frames must survive. A shortfall means a normal variable-GOP keyframe was
// misclassified as a seam.
if totalSync != emittedKeyframes {
t.Errorf("got %d keyframes in output, want %d - a healthy variable-GOP keyframe was wrongly dropped as a seam",
totalSync, emittedKeyframes)
}
if totalSamples != totalEmittedFrames {
t.Errorf("got %d video samples in output, want %d - a healthy variable-GOP GOP was wrongly dropped as a seam",
totalSamples, totalEmittedFrames)
}
}

View File

@@ -96,7 +96,7 @@ func NewAACTranscoder() (*AACTranscoder, error) {
buffered := t.outBuf.Len()
t.outMu.Unlock()
if buffered <= 8192 || buffered%16000 == 0 {
log.Log.Info("webrtc.aac_transcoder: ffmpeg produced PCMU bytes, buffered=" + strconv.Itoa(buffered))
log.Log.Debug("webrtc.aac_transcoder: ffmpeg produced PCMU bytes, buffered=" + strconv.Itoa(buffered))
}
}
if readErr != nil {
@@ -129,14 +129,14 @@ func (t *AACTranscoder) Transcode(adtsData []byte) ([]byte, error) {
return nil, err
}
if len(adtsData) <= 512 || len(adtsData)%1024 == 0 {
log.Log.Info("webrtc.aac_transcoder: wrote AAC bytes to ffmpeg, input=" + strconv.Itoa(len(adtsData)))
log.Log.Debug("webrtc.aac_transcoder: wrote AAC bytes to ffmpeg, input=" + strconv.Itoa(len(adtsData)))
}
deadline := time.Now().Add(75 * time.Millisecond)
for {
data := t.readAvailable()
if len(data) > 0 {
log.Log.Info("webrtc.aac_transcoder: returning PCMU bytes=" + strconv.Itoa(len(data)))
log.Log.Debug("webrtc.aac_transcoder: returning PCMU bytes=" + strconv.Itoa(len(data)))
return data, nil
}
@@ -144,7 +144,7 @@ func (t *AACTranscoder) Transcode(adtsData []byte) ([]byte, error) {
if stderr := t.stderrString(); stderr != "" {
log.Log.Warning("webrtc.aac_transcoder: no output before deadline, ffmpeg stderr: " + stderr)
} else {
log.Log.Info("webrtc.aac_transcoder: no PCMU output before deadline")
log.Log.Debug("webrtc.aac_transcoder: no PCMU output before deadline")
}
return nil, nil
}

View File

@@ -988,7 +988,7 @@ func processAudioPacket(pkt packets.Packet, state *streamState, audioBroadcaster
if len(pcmu) == 0 {
state.aacNoOutput++
if state.aacNoOutput <= 5 || state.aacNoOutput%100 == 0 {
log.Log.Info(fmt.Sprintf("webrtc.main.processAudioPacket(): AAC packet produced no PCMU output yet (aac_packets=%d, no_output=%d, input_bytes=%d)", state.aacPacketsSeen, state.aacNoOutput, len(pkt.Data)))
log.Log.Debug(fmt.Sprintf("webrtc.main.processAudioPacket(): AAC packet produced no PCMU output yet (aac_packets=%d, no_output=%d, input_bytes=%d)", state.aacPacketsSeen, state.aacNoOutput, len(pkt.Data)))
}
return // decoder still buffering
}
@@ -1004,7 +1004,7 @@ func processAudioPacket(pkt packets.Packet, state *streamState, audioBroadcaster
state.lastAudioSample.Duration = sampleDuration(pkt, state.lastAudioSample.PacketTimestamp, 20*time.Millisecond)
state.audioSamplesSent++
if state.audioSamplesSent <= 5 || state.audioSamplesSent%100 == 0 {
log.Log.Info(fmt.Sprintf("webrtc.main.processAudioPacket(): queueing audio sample (samples=%d, codec=%s, bytes=%d, duration_ms=%d, peers=%d)", state.audioSamplesSent, pkt.Codec, len(state.lastAudioSample.Data), state.lastAudioSample.Duration.Milliseconds(), audioBroadcaster.PeerCount()))
log.Log.Debug(fmt.Sprintf("webrtc.main.processAudioPacket(): queueing audio sample (samples=%d, codec=%s, bytes=%d, duration_ms=%d, peers=%d)", state.audioSamplesSent, pkt.Codec, len(state.lastAudioSample.Data), state.lastAudioSample.Duration.Milliseconds(), audioBroadcaster.PeerCount()))
}
audioBroadcaster.WriteSample(*state.lastAudioSample)
}

View File

@@ -1,2 +1,2 @@
#!/bin/bash
swag init -g ./src/routers/http/Server.go
swag init -g ./src/routers/http/server.go