Skip to content

Commit 958cdec

Browse files
authored
Merge pull request #609 from opensensor/fix/spaces-scoped-archive
fix(storage): support bucket-scoped Spaces archives
2 parents 5ad7905 + 24405ec commit 958cdec

8 files changed

Lines changed: 451 additions & 9 deletions

File tree

‎.github/workflows/integration-test.yml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -117,7 +117,7 @@ jobs:
117117
run: |
118118
cd build
119119
ctest --output-on-failure -V \
120-
-R "test_storage_pressure$|test_storage_pressure_extended|test_detection_result_structures|test_storage_retention_sqlite|test_storage_manager_retention|test_storage_target_pressure_cleanup|test_storage_migration|test_storage_archive_s3|test_authorization|test_api_handlers_recordings_playback|test_db_storage_targets|test_db_storage_policies|test_api_handlers_storage_targets|test_api_handlers_storage_policies|test_config|test_logger|test_strings|test_memory|test_request_response|test_shutdown_coordinator|test_detection_config|test_detection_model_motion|test_db_streams|test_db_recordings_extended|test_db_detections|test_db_zones|test_db_events|test_db_auth|test_db_transactions|test_db_maintenance|test_db_query_builder|test_logger_json|test_batch_delete_progress|test_db_recordings_sync|test_httpd_utils|test_zone_filter|test_onvif_soap_fault|test_stream_manager|test_stream_state|test_packet_buffer|test_timestamp_manager|test_motion_trigger_parse|test_url_utils"
120+
-R "test_storage_pressure$|test_storage_pressure_extended|test_detection_result_structures|test_storage_retention_sqlite|test_storage_manager_retention|test_storage_target_pressure_cleanup|test_storage_migration|test_storage_archive_s3|test_storage_archive_operator|test_authorization|test_api_handlers_recordings_playback|test_db_storage_targets|test_db_storage_policies|test_api_handlers_storage_targets|test_api_handlers_storage_policies|test_config|test_logger|test_strings|test_memory|test_request_response|test_shutdown_coordinator|test_detection_config|test_detection_model_motion|test_db_streams|test_db_recordings_extended|test_db_detections|test_db_zones|test_db_events|test_db_auth|test_db_transactions|test_db_maintenance|test_db_query_builder|test_logger_json|test_batch_delete_progress|test_db_recordings_sync|test_httpd_utils|test_zone_filter|test_onvif_soap_fault|test_stream_manager|test_stream_state|test_packet_buffer|test_timestamp_manager|test_motion_trigger_parse|test_url_utils"
121121
122122
- name: Generate coverage report
123123
if: always()
Lines changed: 128 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,128 @@
1+
#!/usr/bin/env python3
2+
"""Publish a short-lived lifecycle inspection for bucket-scoped NVR credentials.
3+
4+
Run in the trusted control plane with boto3 and operator AWS credentials. Never
5+
mount those operator credentials into the tenant. Refresh this file every 15–30
6+
minutes and deliver it atomically beside the NVR's ordinary credentials file.
7+
"""
8+
9+
import argparse
10+
import json
11+
import os
12+
from pathlib import Path
13+
import re
14+
import tempfile
15+
import time
16+
from urllib.parse import urlsplit
17+
import xml.etree.ElementTree as ET
18+
19+
20+
def lifecycle_xml(configuration):
21+
"""Reject every lifecycle action except bounded incomplete-upload cleanup."""
22+
root = ET.Element("LifecycleConfiguration")
23+
if set(configuration) - {"Rules", "ResponseMetadata"}:
24+
raise ValueError("Unrecognized lifecycle configuration")
25+
rules = configuration.get("Rules", [])
26+
if not isinstance(rules, list):
27+
raise ValueError("Lifecycle Rules must be a list")
28+
for rule in rules:
29+
if not isinstance(rule, dict) or set(rule) - {
30+
"ID", "Status", "Filter", "Prefix", "AbortIncompleteMultipartUpload"
31+
}:
32+
raise ValueError("Object expiry, transitions and unknown lifecycle actions are unsupported")
33+
abort = rule.get("AbortIncompleteMultipartUpload", {})
34+
days = abort.get("DaysAfterInitiation") if isinstance(abort, dict) else None
35+
if (not isinstance(abort, dict) or set(abort) != {"DaysAfterInitiation"}
36+
or type(days) is not int or not 1 <= days <= 7):
37+
raise ValueError("Each rule must only abort incomplete uploads after 1–7 days")
38+
# Filters and status only narrow this abort action; they cannot expire
39+
# completed objects. Preserve them in the inspection for operator review.
40+
element = ET.SubElement(root, "Rule")
41+
for key, value in rule.items():
42+
append_xml(element, key, value)
43+
return ET.tostring(root, encoding="unicode")
44+
45+
46+
def append_xml(parent, key, value):
47+
element = ET.SubElement(parent, key)
48+
if isinstance(value, dict):
49+
for child, content in value.items():
50+
append_xml(element, child, content)
51+
elif isinstance(value, list):
52+
for item in value:
53+
append_xml(element, "Tag" if key == "Tags" else "Item", item)
54+
else:
55+
element.text = str(value)
56+
57+
58+
def inspect_bucket(client, endpoint, region, bucket):
59+
from botocore.exceptions import ClientError
60+
61+
checked_at = int(time.time())
62+
versioning = client.get_bucket_versioning(Bucket=bucket)
63+
if set(versioning) - {"ResponseMetadata"}:
64+
raise ValueError("Use a never-versioned bucket; suspended versioning is unsupported")
65+
try:
66+
configuration = client.get_bucket_lifecycle_configuration(Bucket=bucket)
67+
except ClientError as error:
68+
if (error.response.get("Error", {}).get("Code") != "NoSuchLifecycleConfiguration"
69+
or error.response.get("ResponseMetadata", {}).get("HTTPStatusCode") != 404):
70+
raise
71+
configuration = {}
72+
else:
73+
if "Rules" not in configuration:
74+
raise ValueError("Successful lifecycle inspection did not contain Rules")
75+
return {"endpoint": endpoint, "region": region, "bucket": bucket,
76+
"checked_at": checked_at, "lifecycle_configuration_xml": lifecycle_xml(configuration)}
77+
78+
79+
def write_inspection(path, inspection):
80+
encoded = json.dumps(inspection).encode()
81+
if len(encoded) > 8192:
82+
raise ValueError("Lifecycle inspection exceeds the NVR's 8192-byte limit")
83+
path = Path(path)
84+
fd, temporary = tempfile.mkstemp(prefix=f".{path.name}-", dir=path.parent)
85+
try:
86+
with os.fdopen(fd, "wb") as output:
87+
output.write(encoded)
88+
output.flush()
89+
os.fsync(output.fileno())
90+
os.replace(temporary, path)
91+
finally:
92+
if os.path.exists(temporary):
93+
os.unlink(temporary)
94+
95+
96+
def main():
97+
import boto3
98+
from botocore.config import Config
99+
from botocore.exceptions import BotoCoreError, ClientError
100+
101+
parser = argparse.ArgumentParser(description=__doc__)
102+
parser.add_argument("--endpoint", required=True)
103+
parser.add_argument("--region", required=True)
104+
parser.add_argument("--bucket", required=True)
105+
parser.add_argument("--output", required=True, help="e.g. /out/instance.lifecycle.json")
106+
args = parser.parse_args()
107+
url = urlsplit(args.endpoint)
108+
if (url.scheme != "https" or not url.hostname or url.username or url.password
109+
or url.path not in ("", "/") or url.query or url.fragment):
110+
parser.error("Endpoint must be an HTTPS origin")
111+
if not re.fullmatch(r"[a-zA-Z0-9._-]+", args.bucket):
112+
parser.error("Invalid bucket name")
113+
client = boto3.client("s3", endpoint_url=args.endpoint, region_name=args.region,
114+
config=Config(signature_version="s3v4", s3={"addressing_style": "path"},
115+
connect_timeout=10, read_timeout=30,
116+
retries={"max_attempts": 2}))
117+
try:
118+
inspection = inspect_bucket(client, args.endpoint, args.region, args.bucket)
119+
write_inspection(args.output, inspection)
120+
except (BotoCoreError, ClientError, ValueError, OSError) as error:
121+
# Do not print provider bodies, request headers, or credentials. Retain
122+
# the prior file: it will expire instead of publishing an unsafe update.
123+
parser.exit(1, f"Bucket inspection failed ({type(error).__name__}); no verification published\n")
124+
print(f"Lifecycle verified for {args.bucket}; valid for at most one hour")
125+
126+
127+
if __name__ == "__main__":
128+
main()

‎docs/STORAGE_ARCHIVE.md‎

Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -61,6 +61,51 @@ certification: run the connection test and outage/restart pilot on your actual
6161
provider before moving production footage. Bucket policy changes outside LightNVR
6262
remain the operator's responsibility; do not modify managed object keys manually.
6363

64+
### DigitalOcean Spaces and scoped credentials
65+
66+
Spaces bucket-scoped Read/Write/Delete keys can inspect versioning and manage
67+
objects, but lifecycle inspection requires an account-wide key. See the provider's
68+
[permissions table](https://docs.digitalocean.com/reference/api/spaces/).
69+
Keep the account-wide key in the trusted provisioning/control plane. Do not mount
70+
it into a tenant NVR.
71+
72+
The control plane can use
73+
[verify_bucket_lifecycle.py](../deployment/archive/verify_bucket_lifecycle.py)
74+
with Python 3, `boto3`, and its operator AWS credentials to inspect the dedicated
75+
bucket. It rejects versioned buckets and unsafe lifecycle actions, and atomically
76+
writes a mode-0600 verification file:
77+
78+
```sh
79+
python3 deployment/archive/verify_bucket_lifecycle.py \
80+
--endpoint https://nyc3.digitaloceanspaces.com --region nyc3 \
81+
--bucket YOUR_BUCKET --output /out/instance.lifecycle.json
82+
```
83+
84+
Deliver that file beside the tenant's ordinary credential file, using the same
85+
reference plus `.lifecycle.json`. For credential reference `instance`, the files
86+
are `instance` and `instance.lifecycle.json`. Both must be ordinary files owned by
87+
root or the NVR UID without group/other access. A Kubernetes projected Secret must
88+
be copied into ordinary files, as with credentials. An init container alone is
89+
insufficient for continuous refresh: a trusted sidecar must atomically copy updates
90+
into the mounted directory. The tenant must not have permission to update the
91+
source Secret or access operator credentials.
92+
93+
Refresh and deliver the inspection every 15–30 minutes. It is valid for less than
94+
one hour from inspection, binds the exact endpoint, region and bucket, and carries
95+
the inspected lifecycle rules. Only a denied lifecycle GET (HTTP 403) may use it.
96+
Readable unsafe rules, failed versioning checks, public object access, provider
97+
outages, stale timestamps, and malformed files still fail the connection test.
98+
No verification file is needed when the tenant's key can inspect lifecycle rules.
99+
100+
Treat this file as a short-lived operator assertion, not a provider-signed
101+
certificate. The control plane must own bucket configuration and deliver new
102+
inspections securely; changes to lifecycle rules can take up to the remaining
103+
one-hour validity plus the target probe interval to be noticed. If refresh fails,
104+
new archival stops after verification expires and the target is reprobed; existing
105+
objects remain readable with the scoped object credentials. Monitor refresh age
106+
and the archive backlog. Do not synthesize a successful verification after a failed
107+
inspection. This helper does not provision a bucket or schedule refresh for you.
108+
64109
For NFS/SAN, mount the filesystem outside LightNVR and add a filesystem target
65110
with its mount guard enabled. The same archive policy and verified migration
66111
worker apply. S3 targets cannot be capture defaults, capture pool members selected

‎src/storage/storage_s3.c‎

Lines changed: 54 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -121,30 +121,38 @@ bool storage_s3_validate(const storage_target_t *target, char *error, size_t len
121121
return message == NULL;
122122
}
123123

124-
static int credentials_read(const storage_target_t *target, s3_credentials_t *credentials) {
124+
/* Credentials and operator attestations share the same protected file boundary. */
125+
static cJSON *private_json_read(const storage_target_t *target, const char *suffix) {
125126
const char *directory = getenv("LIGHTNVR_ARCHIVE_CREDENTIALS_DIR");
126127
if (!directory) directory = "/etc/lightnvr/archive-credentials";
127-
if (!safe_identifier(target->credential_ref, sizeof(target->credential_ref), false)) return -1;
128+
if (!safe_identifier(target->credential_ref, sizeof(target->credential_ref), false)) return NULL;
129+
char filename[sizeof(target->credential_ref) + 32];
130+
snprintf(filename, sizeof(filename), "%s%s", target->credential_ref, suffix);
128131
int dir = open(directory, O_RDONLY | O_DIRECTORY | O_CLOEXEC | O_NOFOLLOW);
129-
if (dir < 0) return -1;
130-
int descriptor = openat(dir, target->credential_ref, O_RDONLY | O_CLOEXEC | O_NOFOLLOW);
132+
if (dir < 0) return NULL;
133+
int descriptor = openat(dir, filename, O_RDONLY | O_CLOEXEC | O_NOFOLLOW | O_NONBLOCK);
131134
close(dir);
132-
if (descriptor < 0) return -1;
135+
if (descriptor < 0) return NULL;
133136
struct stat status;
134137
char buffer[S3_SECRET_LIMIT + 1];
135138
bool valid = fstat(descriptor, &status) == 0 && S_ISREG(status.st_mode) &&
136139
!(status.st_mode & 0077) && (status.st_uid == 0 || status.st_uid == geteuid()) &&
137140
status.st_size > 0 && status.st_size <= S3_SECRET_LIMIT;
138141
ssize_t count = valid ? read(descriptor, buffer, sizeof(buffer) - 1) : -1;
139142
close(descriptor);
140-
if (count <= 0 || count != status.st_size) return -1;
143+
if (count <= 0 || count != status.st_size) return NULL;
141144
buffer[count] = 0;
142145
cJSON *json = cJSON_Parse(buffer);
143146
wipe(buffer, sizeof(buffer));
147+
return json;
148+
}
149+
150+
static int credentials_read(const storage_target_t *target, s3_credentials_t *credentials) {
151+
cJSON *json = private_json_read(target, "");
144152
const cJSON *access = cJSON_GetObjectItemCaseSensitive(json, "access_key_id");
145153
const cJSON *secret = cJSON_GetObjectItemCaseSensitive(json, "secret_access_key");
146154
const cJSON *token = cJSON_GetObjectItemCaseSensitive(json, "session_token");
147-
valid = cJSON_IsString(access) && cJSON_IsString(secret) &&
155+
bool valid = cJSON_IsString(access) && cJSON_IsString(secret) &&
148156
access->valuestring[0] && secret->valuestring[0] &&
149157
strlen(access->valuestring) < sizeof(credentials->access) &&
150158
strlen(secret->valuestring) < sizeof(credentials->secret) &&
@@ -235,6 +243,19 @@ static int request(const storage_target_t *target, const char *key, const char *
235243
snprintf(raw_key, sizeof(raw_key), "%s%s%s", key ? prefix : "",
236244
key && *prefix && prefix[strlen(prefix) - 1] != '/' ? "/" : "", key ? key : "");
237245
char *encoded = curl ? curl_easy_escape(curl, raw_key, 0) : NULL;
246+
/* S3 SigV4 preserves object-key path separators. Escaping a slash as %2F
247+
* signs a different canonical URI on providers such as Spaces. Literal
248+
* percent signs stay escaped, so a key containing "%2F" is not a slash. */
249+
if (encoded) {
250+
char *read = encoded, *write = encoded;
251+
while (*read) {
252+
if (!strncmp(read, "%2F", 3)) {
253+
*write++ = '/';
254+
read += 3;
255+
} else *write++ = *read++;
256+
}
257+
*write = 0;
258+
}
238259
char url[4 * MAX_PATH_LENGTH + 512];
239260
char endpoint[MAX_PATH_LENGTH];
240261
safe_strcpy(endpoint, target->endpoint, sizeof(endpoint), 0);
@@ -681,6 +702,28 @@ static bool safe_lifecycle(char *body) {
681702
return safe;
682703
}
683704

705+
/* Some providers deny lifecycle GET to bucket-scoped object credentials.
706+
* A trusted operator may supply a fresh, bucket-bound inspection beside the
707+
* credentials. Never infer safety from a 403 or bypass a readable unsafe rule.
708+
* Refresh at least every 30 minutes; stale/malformed files fail closed. */
709+
static bool lifecycle_attested(const storage_target_t *target) {
710+
cJSON *json = private_json_read(target, ".lifecycle.json");
711+
const cJSON *endpoint = cJSON_GetObjectItemCaseSensitive(json, "endpoint");
712+
const cJSON *region = cJSON_GetObjectItemCaseSensitive(json, "region");
713+
const cJSON *bucket = cJSON_GetObjectItemCaseSensitive(json, "bucket");
714+
const cJSON *checked = cJSON_GetObjectItemCaseSensitive(json, "checked_at");
715+
const cJSON *xml = cJSON_GetObjectItemCaseSensitive(json, "lifecycle_configuration_xml");
716+
double now = (double)time(NULL);
717+
bool valid = cJSON_IsString(endpoint) && !strcmp(endpoint->valuestring, target->endpoint) &&
718+
cJSON_IsString(region) && !strcmp(region->valuestring, target->region) &&
719+
cJSON_IsString(bucket) && !strcmp(bucket->valuestring, target->bucket) &&
720+
cJSON_IsNumber(checked) && checked->valuedouble > 0 &&
721+
checked->valuedouble <= now && now - checked->valuedouble < 3600 &&
722+
cJSON_IsString(xml) && safe_lifecycle(xml->valuestring);
723+
cJSON_Delete(json);
724+
return valid;
725+
}
726+
684727
int storage_s3_probe(storage_target_t *target, bool write_test) {
685728
char error[STORAGE_TARGET_ERROR_MAX] = {0};
686729
char *body = calloc(1, S3_REPLY_LIMIT);
@@ -696,7 +739,10 @@ int storage_s3_probe(storage_target_t *target, bool write_test) {
696739
io.length = 0; io.bytes = 0;
697740
result = request(target, NULL, "lifecycle", "GET", NULL, 0, NULL, &io, NULL, error);
698741
if (result == STORAGE_S3_MISSING) result = 0;
699-
else if (result == 0 && !safe_lifecycle(body)) {
742+
else if (result != 0 && io.status == 403) {
743+
if (lifecycle_attested(target)) result = 0;
744+
else error_set(error, "Lifecycle inspection denied; a fresh operator-verified .lifecycle.json file is required with bucket-scoped credentials");
745+
} else if (result == 0 && !safe_lifecycle(body)) {
700746
error_set(error, "Bucket lifecycle may only abort incomplete uploads after 1-7 days; expiry and transitions are unsupported");
701747
result = STORAGE_S3_CONFLICT;
702748
}

‎tests/unit/CMakeLists.txt‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -51,6 +51,8 @@ macro(add_layer2_test TEST_NAME)
5151
find_package(Python3 COMPONENTS Interpreter REQUIRED)
5252
add_test(NAME test_storage_archive_s3 COMMAND ${Python3_EXECUTABLE}
5353
${CMAKE_CURRENT_SOURCE_DIR}/s3_fixture.py $<TARGET_FILE:test_storage_archive>)
54+
add_test(NAME test_storage_archive_operator COMMAND ${Python3_EXECUTABLE}
55+
${CMAKE_CURRENT_SOURCE_DIR}/test_archive_operator.py)
5456
else()
5557
add_test(NAME ${TEST_NAME} COMMAND ${TEST_NAME})
5658
endif()

‎tests/unit/s3_fixture.py‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,10 @@ def handle_request(self):
6868
return self.respond(413)
6969
body = self.rfile.read(length) if length else b''
7070
url = urlsplit(self.path)
71+
# S3 canonical paths preserve key separators. Spaces rejects signatures
72+
# using encoded separators even though the decoded object key matches.
73+
if '%2F' in url.path.upper():
74+
return self.respond(403, b'<Error><Code>SignatureDoesNotMatch</Code></Error>')
7175
path = unquote(url.path).strip('/')
7276
bucket, _, key = path.partition('/')
7377
anonymous_public = bucket == 'public' and self.command == 'GET' and not self.headers.get('Authorization')
@@ -78,6 +82,10 @@ def handle_request(self):
7882
value = b'<Status>Enabled</Status>' if bucket == 'versioned' else b''
7983
return self.respond(200, b'<VersioningConfiguration>' + value + b'</VersioningConfiguration>')
8084
if url.query == 'lifecycle':
85+
if bucket == 'scoped-lifecycle':
86+
return self.respond(403)
87+
if bucket == 'lifecycle-unavailable':
88+
return self.respond(503)
8189
if bucket == 'expiry':
8290
return self.respond(200, b'<LifecycleConfiguration><Rule><Status>Enabled</Status><Expiration><Days>1</Days></Expiration></Rule></LifecycleConfiguration>')
8391
if bucket == 'abort-only':

0 commit comments

Comments
 (0)