Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0f94388f49 | ||
|
|
94a3233ce2 |
@@ -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.9
|
||||||
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"]
|
||||||
|
|||||||
@@ -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.9
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
network_mode: host
|
network_mode: host
|
||||||
|
|||||||
@@ -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.9
|
||||||
user: "1001:1001"
|
user: "1001:1001"
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
network_mode: host
|
network_mode: host
|
||||||
|
|||||||
@@ -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.9
|
||||||
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
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "archive-clients"
|
name = "archive-clients"
|
||||||
version = "0.1.7"
|
version = "0.1.9"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||||
|
|
||||||
|
|||||||
+48
-24
@@ -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,
|
||||||
@@ -119,7 +120,11 @@ class ClientJobExecutor:
|
|||||||
replay = self._replay(
|
replay = self._replay(
|
||||||
command.job_id, command.expected_last_event_sequence
|
command.job_id, command.expected_last_event_sequence
|
||||||
)
|
)
|
||||||
if replay and replay[-1].type in {
|
step_replay = [
|
||||||
|
event for event in replay
|
||||||
|
if event.progress.step == command.step
|
||||||
|
]
|
||||||
|
if step_replay and step_replay[-1].type in {
|
||||||
control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
|
control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
|
||||||
control_pb2.JOB_EVENT_TYPE_COMMITTED,
|
control_pb2.JOB_EVENT_TYPE_COMMITTED,
|
||||||
control_pb2.JOB_EVENT_TYPE_SUCCEEDED,
|
control_pb2.JOB_EVENT_TYPE_SUCCEEDED,
|
||||||
@@ -131,31 +136,45 @@ class ClientJobExecutor:
|
|||||||
for event in replay:
|
for event in replay:
|
||||||
event_callback(event)
|
event_callback(event)
|
||||||
return replay
|
return replay
|
||||||
started = replay[0] if replay else self._event(
|
if step_replay:
|
||||||
definition,
|
started = step_replay[0]
|
||||||
sequence=command.expected_last_event_sequence + 1,
|
cursor = step_replay[-1]
|
||||||
revision=command.expected_job_revision + 1,
|
emitted = list(step_replay)
|
||||||
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
|
else:
|
||||||
state=job_pb2.JOB_STATE_RUNNING,
|
previous = replay[-1] if replay else None
|
||||||
committed=(
|
started = self._event(
|
||||||
self._committed(command.job_id)
|
definition,
|
||||||
or command.step == job_pb2.JOB_STEP_KIND_STAGING_CLEANUP
|
sequence=(
|
||||||
or (
|
previous.sequence + 1
|
||||||
command.step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK
|
if previous is not None
|
||||||
and self.store.get_job_artifact(
|
else command.expected_last_event_sequence + 1
|
||||||
command.job_id, "qb-entry-removed"
|
),
|
||||||
) is not None
|
revision=(
|
||||||
)
|
previous.job_revision + 1
|
||||||
),
|
if previous is not None
|
||||||
step=command.step,
|
else command.expected_job_revision + 1
|
||||||
step_state=job_pb2.STEP_STATE_RUNNING,
|
),
|
||||||
)
|
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
|
||||||
if not replay:
|
state=job_pb2.JOB_STATE_RUNNING,
|
||||||
|
committed=(
|
||||||
|
self._committed(command.job_id)
|
||||||
|
or command.step == job_pb2.JOB_STEP_KIND_STAGING_CLEANUP
|
||||||
|
or (
|
||||||
|
command.step == job_pb2.JOB_STEP_KIND_SAFE_FILE_UNLINK
|
||||||
|
and self.store.get_job_artifact(
|
||||||
|
command.job_id, "qb-entry-removed"
|
||||||
|
) is not None
|
||||||
|
)
|
||||||
|
),
|
||||||
|
step=command.step,
|
||||||
|
step_state=job_pb2.STEP_STATE_RUNNING,
|
||||||
|
)
|
||||||
self._record(definition, started)
|
self._record(definition, started)
|
||||||
|
cursor = started
|
||||||
|
emitted = [started]
|
||||||
if event_callback is not None:
|
if event_callback is not None:
|
||||||
event_callback(started)
|
for event in emitted:
|
||||||
emitted = [started]
|
event_callback(event)
|
||||||
cursor = started
|
|
||||||
speed_sample = [time.monotonic(), 0]
|
speed_sample = [time.monotonic(), 0]
|
||||||
|
|
||||||
def progress(
|
def progress(
|
||||||
@@ -833,6 +852,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():
|
||||||
|
|||||||
@@ -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,
|
||||||
|
|||||||
@@ -111,6 +111,72 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
|||||||
syncthing=SlowRescanSyncthing(),
|
syncthing=SlowRescanSyncthing(),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
def test_target_step_advances_past_replayed_source_completion(self):
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
store = ClientStore(root / "client.db")
|
||||||
|
store.initialize()
|
||||||
|
definition = job_pb2.JobDefinition(
|
||||||
|
job_id=str(uuid4()),
|
||||||
|
idempotency_key=str(uuid4()),
|
||||||
|
operation=job_pb2.JOB_OPERATION_ARCHIVE,
|
||||||
|
resource_display_name="fixture",
|
||||||
|
transfer={
|
||||||
|
"source_client_id": "cache-1",
|
||||||
|
"target_client_id": "archive-1",
|
||||||
|
"route_id": "route-1",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
definition.resource_id.info_hash_v1_hex = "a" * 40
|
||||||
|
definition.created_at.GetCurrentTime()
|
||||||
|
executor = ClientJobExecutor(
|
||||||
|
client_id="archive-1",
|
||||||
|
qbittorrent=Mock(),
|
||||||
|
store=store,
|
||||||
|
qb_root=root,
|
||||||
|
qb_api_root=Path("/downloads"),
|
||||||
|
route_path=lambda _: root,
|
||||||
|
syncthing_transport=Mock(),
|
||||||
|
sparse_supported=True,
|
||||||
|
poll_interval=0,
|
||||||
|
)
|
||||||
|
executor.assign(control_pb2.AssignJobCommand(
|
||||||
|
job=definition,
|
||||||
|
expected_job_revision=1,
|
||||||
|
expected_last_event_sequence=1,
|
||||||
|
))
|
||||||
|
source_complete = executor._event(
|
||||||
|
definition,
|
||||||
|
sequence=5,
|
||||||
|
revision=4,
|
||||||
|
event_type=control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
|
||||||
|
state=job_pb2.JOB_STATE_RUNNING,
|
||||||
|
committed=False,
|
||||||
|
step=job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
|
||||||
|
step_state=job_pb2.STEP_STATE_SUCCEEDED,
|
||||||
|
)
|
||||||
|
executor._record(definition, source_complete)
|
||||||
|
|
||||||
|
with patch.object(executor, "_wait_for_syncthing"):
|
||||||
|
events = executor.execute(control_pb2.ExecuteStepCommand(
|
||||||
|
job_id=definition.job_id,
|
||||||
|
expected_job_revision=4,
|
||||||
|
expected_last_event_sequence=5,
|
||||||
|
step=job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
|
||||||
|
attempt=1,
|
||||||
|
))
|
||||||
|
|
||||||
|
self.assertEqual(events[0].sequence, 6)
|
||||||
|
self.assertEqual(
|
||||||
|
events[0].progress.step,
|
||||||
|
job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
|
||||||
|
)
|
||||||
|
self.assertEqual(events[-1].sequence, 7)
|
||||||
|
self.assertEqual(
|
||||||
|
events[-1].type,
|
||||||
|
control_pb2.JOB_EVENT_TYPE_STEP_SUCCEEDED,
|
||||||
|
)
|
||||||
|
|
||||||
def _run_transfer(
|
def _run_transfer(
|
||||||
self,
|
self,
|
||||||
root: Path,
|
root: Path,
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
Reference in New Issue
Block a user