Compare commits

...
1 Commits
Author SHA1 Message Date
cabbage 94a3233ce2 fix: clean interrupted staging temporaries 2026-07-24 17:00:10 +00:00
8 changed files with 69 additions and 5 deletions
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.7 image: sodium/archive-clients:v0.1.8
user: "1000:1000" user: "1000:1000"
restart: unless-stopped restart: unless-stopped
command: ["--config", "/etc/archive-control/client.toml"] command: ["--config", "/etc/archive-control/client.toml"]
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.7 image: sodium/archive-clients:v0.1.8
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.7 image: sodium/archive-clients:v0.1.8
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
+1 -1
View File
@@ -218,7 +218,7 @@ cache/archive routes according to policy.
```yaml ```yaml
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.7 image: sodium/archive-clients:v0.1.8
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
command: ["archive-client", "--config", "/etc/archive-control/client.toml"] command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "archive-clients" name = "archive-clients"
version = "0.1.7" version = "0.1.8"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+6
View File
@@ -32,6 +32,7 @@ from archive_clients.syncthing import RouteSetupError, SyncthingTransferObserver
from archive_clients.transfer import ( from archive_clients.transfer import (
TransferError, TransferError,
TransferIntegrityError, TransferIntegrityError,
cleanup_partial_transfer,
cleanup_transfer, cleanup_transfer,
load_published_transfer, load_published_transfer,
materialize_transfer, materialize_transfer,
@@ -833,6 +834,11 @@ class ClientJobExecutor:
qb_root=job_directory, qb_root=job_directory,
store=self.store, store=self.store,
) )
cleanup_partial_transfer(
job_directory,
job_id=definition.job_id,
store=self.store,
)
for name in ("ready.json", "manifest.json"): for name in ("ready.json", "manifest.json"):
candidate = job_directory / name candidate = job_directory / name
if candidate.is_file() and not candidate.is_symlink(): if candidate.is_file() and not candidate.is_symlink():
+30
View File
@@ -456,6 +456,36 @@ def cleanup_transfer(job_directory: Path) -> bool:
return True return True
def cleanup_partial_transfer(
job_directory: Path,
*,
job_id: str,
store: ClientStore,
) -> None:
"""Remove copy/reflink temporaries left before a transfer is published.
A source-stage failure before ``manifest.json``/``ready.json`` exists
cannot use :func:`cleanup_transfer`. The file-operation journal is the
authoritative list of owned destinations, so derive each temporary path
from it rather than recursively removing arbitrary content from a shared
Syncthing folder.
"""
for row in store.file_operation_rows(job_id):
intent = json.loads(str(row["intent_json"]))
destination = job_directory / _relative_path(str(intent["destination"]))
temporary = _temporary_path(destination, str(row["operation_id"]))
try:
metadata = temporary.lstat()
except FileNotFoundError:
continue
if not stat.S_ISREG(metadata.st_mode):
raise TransferIntegrityError(
"job-owned temporary cleanup path is not a regular file"
)
temporary.unlink()
def canonical_message_json(message: object) -> bytes: def canonical_message_json(message: object) -> bytes:
value = json_format.MessageToDict( value = json_format.MessageToDict(
message, message,
+28
View File
@@ -11,6 +11,7 @@ from archive_clients.transfer import (
FileMaterializer, FileMaterializer,
TransferIntegrityError, TransferIntegrityError,
canonical_message_json, canonical_message_json,
cleanup_partial_transfer,
load_published_transfer, load_published_transfer,
materialize_transfer, materialize_transfer,
stage_transfer, stage_transfer,
@@ -206,6 +207,33 @@ class TransferHappyPathTests(unittest.TestCase):
os.stat(copied).st_blocks * 512, os.stat(copied).st_size os.stat(copied).st_blocks * 512, os.stat(copied).st_size
) )
def test_partial_cleanup_removes_only_journalled_temporary(self):
job_id = str(uuid4())
job_root = self.sync / ".archive-control/jobs" / job_id
destination = job_root / "payload/album/one.bin"
destination.parent.mkdir(parents=True)
operation_id = "interrupted-copy"
self.store.begin_file_operation(
operation_id,
job_id,
'{"destination":"payload/album/one.bin"}',
)
temporary = destination.with_name(
".one.bin.archive-control-"
+ hashlib.sha256(operation_id.encode()).hexdigest()[:16]
+ ".tmp"
)
temporary.write_bytes(b"partial")
unrelated = destination.parent / "keep-me"
unrelated.write_bytes(b"unrelated")
cleanup_partial_transfer(
job_root, job_id=job_id, store=self.store
)
self.assertFalse(temporary.exists())
self.assertTrue(unrelated.is_file())
def test_canonical_manifest_json_is_stable(self): def test_canonical_manifest_json_is_stable(self):
manifest = self._manifest() manifest = self._manifest()
first = canonical_message_json(manifest) first = canonical_message_json(manifest)