diff --git a/machinery/src/cloud/cloud.go b/machinery/src/cloud/cloud.go index 5f628a1..21fac4c 100644 --- a/machinery/src/cloud/cloud.go +++ b/machinery/src/cloud/cloud.go @@ -2,6 +2,7 @@ package cloud import ( "bytes" + "context" "crypto/tls" "encoding/base64" "encoding/json" @@ -21,6 +22,7 @@ import ( "time" "github.com/kerberos-io/agent/machinery/src/capture" + "github.com/kerberos-io/agent/machinery/src/cloud/livesnapshot" "github.com/kerberos-io/agent/machinery/src/encryption" "github.com/kerberos-io/agent/machinery/src/log" "github.com/kerberos-io/agent/machinery/src/models" @@ -684,7 +686,28 @@ func HandleLiveStreamSD(livestreamCursor *packets.QueueCursor, configuration *mo hubKey = config.HubKey } - lastLivestreamRequest := int64(0) + lastLivestreamRequestMQTT := int64(0) + lastLivestreamRequestHTTP := int64(0) + + // HTTP transport (preferred when this agent is paired with a Kerberos + // Hub): ship preview frames to hub-api over HTTPS instead of pushing + // (large, base64) images through the MQTT broker. Viewers opt in per + // session via the "http" transport on their keepalive; the legacy MQTT + // push is kept for viewers (older frontends) that don't, and as a fallback. + region := "" + if config.S3 != nil { + region = config.S3.Region + } + var snapshotPublisher *livesnapshot.Publisher + if config.HubURI != "" && config.HubKey != "" { + snapshotPublisher = livesnapshot.NewPublisher(livesnapshot.PublisherConfig{ + HubURI: config.HubURI, + HubKey: config.HubKey, + HubPrivateKey: config.HubPrivateKey, + Region: region, + DeviceKey: deviceId, + }) + } var cursorError error var pkt packets.Packet @@ -695,20 +718,49 @@ func HandleLiveStreamSD(livestreamCursor *packets.QueueCursor, configuration *mo continue } now := time.Now().Unix() + // Drain both viewer keepalive channels (non-blocking): one for the + // HTTP transport, one for the legacy MQTT push. select { case <-communication.HandleLiveSD: - lastLivestreamRequest = now + lastLivestreamRequestMQTT = now default: } - if now-lastLivestreamRequest > 3 { + select { + case <-communication.HandleLiveSDHTTP: + lastLivestreamRequestHTTP = now + default: + } + + mqttViewerActive := now-lastLivestreamRequestMQTT <= 3 + httpViewerActive := now-lastLivestreamRequestHTTP <= 3 + if !mqttViewerActive && !httpViewerActive { continue } - log.Log.Info("cloud.HandleLiveStreamSD(): Sending base64 encoded images to MQTT.") - img, err := rtspClient.DecodePacket(pkt) - if err == nil { - imageResized, _ := utils.ResizeImage(&img, uint(config.Capture.IPCamera.BaseWidth), uint(config.Capture.IPCamera.BaseHeight)) - bytes, _ := utils.ImageToBytes(imageResized) + img, err := rtspClient.DecodePacket(pkt) + if err != nil { + continue + } + imageResized, _ := utils.ResizeImage(&img, uint(config.Capture.IPCamera.BaseWidth), uint(config.Capture.IPCamera.BaseHeight)) + bytes, _ := utils.ImageToBytes(imageResized) + + // Prefer HTTP for viewers that asked for it. Only if that did not + // deliver (Hub not configured, or the upload failed) do we also push + // over MQTT, so a new frontend can still fall back to its MQTT path. + httpPushed := false + if httpViewerActive && snapshotPublisher != nil { + ctx, cancel := context.WithTimeout(context.Background(), 4*time.Second) + if errPublish := snapshotPublisher.PublishSnapshot(ctx, bytes); errPublish != nil { + log.Log.Error("cloud.HandleLiveStreamSD(): failed to publish preview frame over HTTP: " + errPublish.Error()) + } else { + httpPushed = true + } + cancel() + } + + pushMQTT := mqttViewerActive || (httpViewerActive && !httpPushed) + if pushMQTT { + log.Log.Info("cloud.HandleLiveStreamSD(): Sending base64 encoded images to MQTT.") chunking := config.Capture.LiveviewChunking if chunking == "true" { diff --git a/machinery/src/cloud/livesnapshot/publisher.go b/machinery/src/cloud/livesnapshot/publisher.go new file mode 100644 index 0000000..2abe426 --- /dev/null +++ b/machinery/src/cloud/livesnapshot/publisher.go @@ -0,0 +1,151 @@ +// Package livesnapshot implements the agent-side producer for the live-view +// "preview" (SD) mode over HTTP. +// +// Historically the preview pipeline shipped each resized keyframe (a base64 +// JPEG, often chunked) to viewers over the MQTT broker. MQTT is a control plane +// for small messages, so pushing ~1 image/second of base64 image data per +// watched camera congests the broker and delays genuine control traffic. This +// package moves those frames off MQTT: the agent POSTs the latest resized JPEG +// straight to hub-api over plain HTTPS (outbound only), and viewers fetch it +// back with their session token. Only the tiny "a viewer is watching" keepalive +// stays on MQTT. +// +// The wire contract (agent -> hub-api) deliberately mirrors the live HLS ingest +// and the existing storage-upload convention (X-Kerberos-Storage-Device plus the +// Hub public/private key auth headers). hub-api authenticates the agent and +// stores the frame in an ephemeral, short-TTL per-device slot which it serves +// straight back to authorized viewers; the frame never enters the vault or the +// recordings collection. +// +// Like live HLS segments, a preview frame is worthless once stale: a frame that +// fails to upload is superseded by the next one a second later, so the publisher +// is fire-and-forget and drops on failure (logged) rather than retrying. +package livesnapshot + +import ( + "bytes" + "context" + "fmt" + "net/http" + "strings" + "time" + + "github.com/kerberos-io/agent/machinery/src/log" +) + +const ( + // snapshotIngestPath is the hub-api endpoint that accepts the latest preview + // frame and stores it in the device's ephemeral snapshot slot (mirrors the + // /storage/live live-HLS ingest convention). + snapshotIngestPath = "/storage/snapshot" + + contentTypeJPEG = "image/jpeg" + + // Header names for the snapshot ingest contract (shared with live HLS / storage). + headerHubPublicKey = "X-Kerberos-Hub-PublicKey" + headerHubPrivateKey = "X-Kerberos-Hub-PrivateKey" + headerHubRegion = "X-Kerberos-Hub-Region" + headerStorageDevice = "X-Kerberos-Storage-Device" + + // defaultPublishTimeout bounds a single snapshot upload. Preview frames are + // produced roughly once a second from a single goroutine, so an upload that + // cannot land in a few seconds is abandoned rather than allowed to back up the + // preview loop behind a slow request. + defaultPublishTimeout = 4 * time.Second +) + +// PublisherConfig carries the hub endpoint and credentials needed to ship +// preview frames. It is populated from the agent's models.Config (the same +// HubURI/HubKey/HubPrivateKey used by recordings and live HLS). +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 the latest preview frame to hub-api over plain HTTP POST. +// +// It is safe for sequential use from a single live-stream goroutine. PublishSnapshot +// is fire-and-forget: it returns 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/live-HLS upload clients. +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} +} + +// PublishSnapshot uploads a single resized preview frame (JPEG) as the device's +// latest snapshot. It overwrites whatever frame was there before, so viewers +// always fetch the most recent frame. +func (p *Publisher) PublishSnapshot(ctx context.Context, jpeg []byte) error { + if p.cfg.HubURI == "" { + return fmt.Errorf("livesnapshot: HubURI not configured") + } + if len(jpeg) == 0 { + return fmt.Errorf("livesnapshot: empty snapshot body") + } + + url := strings.TrimRight(p.cfg.HubURI, "/") + snapshotIngestPath + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(jpeg)) + if err != nil { + return fmt.Errorf("livesnapshot: build request: %w", err) + } + + req.Header.Set("Content-Type", contentTypeJPEG) + req.Header.Set(headerStorageDevice, p.cfg.DeviceKey) + 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("livesnapshot: upload snapshot: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return fmt.Errorf("livesnapshot: upload snapshot rejected: %s", resp.Status) + } + log.Log.Debug("livesnapshot.Publisher.PublishSnapshot(): shipped preview frame for device " + p.cfg.DeviceKey) + 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 +} diff --git a/machinery/src/components/kerberos.go b/machinery/src/components/kerberos.go index 65d5ee0..6839ddc 100644 --- a/machinery/src/components/kerberos.go +++ b/machinery/src/components/kerberos.go @@ -70,6 +70,7 @@ func Bootstrap(ctx context.Context, configDirectory string, configuration *model communication.HandleUpload = make(chan string, 1) communication.HandleHeartBeat = make(chan string, 1) communication.HandleLiveSD = make(chan int64, 1) + communication.HandleLiveSDHTTP = make(chan int64, 1) communication.HandleLiveHDKeepalive = make(chan string, 1) communication.HandleLiveHDPeers = make(chan string, 1) communication.HandleLiveHLS = make(chan int64, 1) diff --git a/machinery/src/models/communication.go b/machinery/src/models/communication.go index edb50db..1eafd4c 100644 --- a/machinery/src/models/communication.go +++ b/machinery/src/models/communication.go @@ -37,6 +37,7 @@ type Communication struct { HandleUpload chan string HandleHeartBeat chan string HandleLiveSD chan int64 + HandleLiveSDHTTP chan int64 HandleLiveHDKeepalive chan string HandleLiveHDHandshake chan LiveHDHandshake HandleLiveHDPeers chan string diff --git a/machinery/src/models/mqtt.go b/machinery/src/models/mqtt.go index fdb0bbc..f5ec05d 100644 --- a/machinery/src/models/mqtt.go +++ b/machinery/src/models/mqtt.go @@ -171,6 +171,12 @@ type UpdateConfigPayload struct { // We received a request SD stream request type RequestSDStreamPayload struct { Timestamp int64 `json:"timestamp"` // timestamp + // Transport selects how the agent should deliver the preview frames for this + // viewer. "http" asks the agent to POST frames to hub-api (keeping them off + // MQTT); empty/absent means the legacy MQTT image push. Older agents simply + // ignore this unknown field and keep doing MQTT, and older frontends never set + // it — so new/old agents and frontends interoperate in every combination. + Transport string `json:"transport,omitempty"` } // We received a live HLS stream request. Like SD it is a simple viewer diff --git a/machinery/src/routers/mqtt/main.go b/machinery/src/routers/mqtt/main.go index 43e2d63..6fd81b1 100644 --- a/machinery/src/routers/mqtt/main.go +++ b/machinery/src/routers/mqtt/main.go @@ -550,9 +550,20 @@ func HandleRequestSDStream(mqttClient mqtt.Client, hubKey string, payload models if requestSDStreamPayload.Timestamp != 0 { if communication.CameraConnected { - select { - case communication.HandleLiveSD <- time.Now().Unix(): - default: + // A viewer that opted into the HTTP transport is signalled on a separate + // channel so the producer ships its frames to hub-api over HTTP instead of + // publishing them over MQTT. Any other (or absent) transport keeps the + // legacy MQTT image push, so older frontends behave exactly as before. + if requestSDStreamPayload.Transport == "http" { + select { + case communication.HandleLiveSDHTTP <- time.Now().Unix(): + default: + } + } else { + select { + case communication.HandleLiveSD <- time.Now().Unix(): + default: + } } log.Log.Info("routers.mqtt.main.HandleRequestSDStream(): received request to livestream.") } else {