Compare commits

...

7 Commits

Author SHA1 Message Date
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
22 changed files with 284 additions and 31 deletions

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

@@ -93,17 +93,26 @@ func logTusUploadProgress(label string, offset, size int64, loggedBucket *int64)
log.Log.Infof("%s: resumable upload progress %d%% (%d/%d bytes)", label, percent, offset, size)
}
// uploadVaultResumable uploads a recording to a Kerberos Vault using the tus
// resumable upload protocol.
// 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 vault.
// - responded: the vault returned a definitive HTTP response (used by the
// - 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 vault exposes a tus endpoint. When false, the caller should
// fall back to the legacy single-POST upload (older vault deployments).
// - 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 uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName, label, slot string) (uploaded bool, responded bool, supported bool, body string, err error) {
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)
@@ -124,17 +133,20 @@ func uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName,
}
size := info.Size()
baseURL := strings.TrimRight(vault.URI, "/") + tusUploadPath
client := newVaultHTTPClient(0)
metadata := encodeTusMetadata(map[string]string{
"filename": fileName,
"device": deviceKey,
"directory": vault.Directory,
"provider": vault.Provider,
"capture": "IPCamera",
"cloudkey": publicKey,
})
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)
@@ -145,7 +157,7 @@ func uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName,
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, vault, publicKey, deviceKey, fileName)
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.
@@ -160,7 +172,7 @@ func uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName,
}
// (2) Query the current server-side offset.
offset, status, herr := tusHead(client, uploadURL, vault, publicKey, deviceKey)
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.
@@ -180,7 +192,7 @@ func uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName,
if restartedAfterComplete {
return false, true, true, "resumable finalize did not complete", errors.New(label + ": resumable finalize did not complete")
}
tusTerminate(client, uploadURL, vault, publicKey, deviceKey)
tusTerminate(client, uploadURL, setHeaders)
removeTusResumeState(sidecar)
uploadURL = ""
restartedAfterComplete = true
@@ -207,7 +219,7 @@ func uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName,
if chunkSize > 0 && chunkSize < patchLen {
patchLen = chunkSize
}
newOffset, status, respBody, perr := tusPatch(client, uploadURL, offset, patchLen, file, vault, publicKey, deviceKey)
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).
@@ -249,9 +261,47 @@ func uploadVaultResumable(vault models.KStorage, publicKey, deviceKey, fileName,
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, vault models.KStorage, publicKey, deviceKey, fileName string) (string, int, error) {
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
@@ -261,7 +311,7 @@ func tusCreate(client *http.Client, baseURL string, size int64, metadata string,
if metadata != "" {
req.Header.Set("Upload-Metadata", metadata)
}
setVaultTusHeaders(req.Header, vault, publicKey, deviceKey, fileName)
setHeaders(req.Header, fileName)
resp, err := client.Do(req)
if resp != nil {
@@ -284,13 +334,13 @@ func tusCreate(client *http.Client, baseURL string, size int64, metadata string,
// tusHead performs the tus "offset" request (HEAD) and returns the current
// server-side upload offset.
func tusHead(client *http.Client, uploadURL string, vault models.KStorage, publicKey, deviceKey string) (int64, int, error) {
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)
setVaultTusHeaders(req.Header, vault, publicKey, deviceKey, "")
setHeaders(req.Header, "")
resp, err := client.Do(req)
if resp != nil {
@@ -315,7 +365,7 @@ func tusHead(client *http.Client, uploadURL string, vault models.KStorage, publi
// 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, vault models.KStorage, publicKey, deviceKey string) (int64, int, string, error) {
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
@@ -324,7 +374,7 @@ func tusPatch(client *http.Client, uploadURL string, offset, length int64, file
req.Header.Set("Tus-Resumable", tusResumableVersion)
req.Header.Set("Content-Type", "application/offset+octet-stream")
req.Header.Set("Upload-Offset", strconv.FormatInt(offset, 10))
setVaultTusHeaders(req.Header, vault, publicKey, deviceKey, "")
setHeaders(req.Header, "")
resp, err := client.Do(req)
if resp != nil {
@@ -349,13 +399,13 @@ func tusPatch(client *http.Client, uploadURL string, offset, length int64, file
}
// tusTerminate best-effort deletes an upload server-side (DELETE).
func tusTerminate(client *http.Client, uploadURL string, vault models.KStorage, publicKey, deviceKey string) {
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)
setVaultTusHeaders(req.Header, vault, publicKey, deviceKey, "")
setHeaders(req.Header, "")
resp, derr := client.Do(req)
if resp != nil {
@@ -383,6 +433,22 @@ func setVaultTusHeaders(h http.Header, vault models.KStorage, publicKey, deviceK
}
}
// 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.

View File

@@ -2,6 +2,7 @@ package cloud
import (
"bytes"
"encoding/base64"
"fmt"
"io"
"net/http"
@@ -22,6 +23,13 @@ type fakeUpload struct {
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 {
@@ -38,6 +46,10 @@ type fakeTus struct {
// 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 {
@@ -84,10 +96,27 @@ func (s *fakeTus) createCount() int {
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 {
@@ -378,6 +407,137 @@ func TestUploadVaultResumable_ResumeFromSidecar(t *testing.T) {
}
}
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",

View File

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