Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9df2e6991e | ||
|
|
d2fa69a1d6 | ||
|
|
0f94388f49 | ||
|
|
94a3233ce2 |
@@ -2,7 +2,7 @@ name: archive-control-archive
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.7
|
||||
image: sodium/archive-clients:v0.1.11
|
||||
user: "1000:1000"
|
||||
restart: unless-stopped
|
||||
command: ["--config", "/etc/archive-control/client.toml"]
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.7
|
||||
image: sodium/archive-clients:v0.1.11
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
network_mode: host
|
||||
|
||||
@@ -2,7 +2,7 @@ name: archive-control-cache
|
||||
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.7
|
||||
image: sodium/archive-clients:v0.1.11
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
network_mode: host
|
||||
|
||||
@@ -218,7 +218,7 @@ cache/archive routes according to policy.
|
||||
```yaml
|
||||
services:
|
||||
archive-client:
|
||||
image: sodium/archive-clients:v0.1.7
|
||||
image: sodium/archive-clients:v0.1.11
|
||||
user: "1001:1001"
|
||||
restart: unless-stopped
|
||||
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
|
||||
|
||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
||||
|
||||
[project]
|
||||
name = "archive-clients"
|
||||
version = "0.1.7"
|
||||
version = "0.1.11"
|
||||
requires-python = ">=3.11"
|
||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||
|
||||
|
||||
+60
-27
@@ -32,6 +32,7 @@ from archive_clients.syncthing import RouteSetupError, SyncthingTransferObserver
|
||||
from archive_clients.transfer import (
|
||||
TransferError,
|
||||
TransferIntegrityError,
|
||||
cleanup_partial_transfer,
|
||||
cleanup_transfer,
|
||||
load_published_transfer,
|
||||
materialize_transfer,
|
||||
@@ -119,7 +120,11 @@ class ClientJobExecutor:
|
||||
replay = self._replay(
|
||||
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_COMMITTED,
|
||||
control_pb2.JOB_EVENT_TYPE_SUCCEEDED,
|
||||
@@ -131,32 +136,51 @@ class ClientJobExecutor:
|
||||
for event in replay:
|
||||
event_callback(event)
|
||||
return replay
|
||||
started = replay[0] if replay else self._event(
|
||||
definition,
|
||||
sequence=command.expected_last_event_sequence + 1,
|
||||
revision=command.expected_job_revision + 1,
|
||||
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
|
||||
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,
|
||||
)
|
||||
if not replay:
|
||||
if step_replay:
|
||||
started = step_replay[0]
|
||||
cursor = step_replay[-1]
|
||||
emitted = list(step_replay)
|
||||
else:
|
||||
previous = replay[-1] if replay else None
|
||||
started = self._event(
|
||||
definition,
|
||||
sequence=(
|
||||
previous.sequence + 1
|
||||
if previous is not None
|
||||
else command.expected_last_event_sequence + 1
|
||||
),
|
||||
revision=(
|
||||
previous.job_revision + 1
|
||||
if previous is not None
|
||||
else command.expected_job_revision + 1
|
||||
),
|
||||
event_type=control_pb2.JOB_EVENT_TYPE_STEP_STARTED,
|
||||
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)
|
||||
cursor = started
|
||||
emitted = [started]
|
||||
if event_callback is not None:
|
||||
event_callback(started)
|
||||
emitted = [started]
|
||||
cursor = started
|
||||
speed_sample = [time.monotonic(), 0]
|
||||
for event in emitted:
|
||||
event_callback(event)
|
||||
# A newly-started step must publish its first observation. In
|
||||
# particular, a sparse Syncthing temporary file may retain the same
|
||||
# allocated-block count for a long time while data is still needed;
|
||||
# suppressing that first observation made a healthy transfer look
|
||||
# permanently stalled to the control daemon.
|
||||
speed_sample: list[float | int | None] = [None, 0]
|
||||
|
||||
def progress(
|
||||
fraction: float,
|
||||
@@ -166,8 +190,12 @@ class ClientJobExecutor:
|
||||
) -> None:
|
||||
nonlocal cursor
|
||||
now = time.monotonic()
|
||||
elapsed = now - float(speed_sample[0])
|
||||
if fraction < 1 and elapsed < 1:
|
||||
previous_time = speed_sample[0]
|
||||
elapsed = (
|
||||
now - float(previous_time)
|
||||
if previous_time is not None else 0.0
|
||||
)
|
||||
if previous_time is not None and fraction < 1 and elapsed < 1:
|
||||
return
|
||||
speed = (
|
||||
max(0, bytes_complete - int(speed_sample[1])) / elapsed
|
||||
@@ -833,6 +861,11 @@ class ClientJobExecutor:
|
||||
qb_root=job_directory,
|
||||
store=self.store,
|
||||
)
|
||||
cleanup_partial_transfer(
|
||||
job_directory,
|
||||
job_id=definition.job_id,
|
||||
store=self.store,
|
||||
)
|
||||
for name in ("ready.json", "manifest.json"):
|
||||
candidate = job_directory / name
|
||||
if candidate.is_file() and not candidate.is_symlink():
|
||||
|
||||
@@ -139,19 +139,14 @@ class SyncthingTransferObserver:
|
||||
|
||||
def status(self) -> SyncthingTransferStatus:
|
||||
published = load_published_transfer(self.local_job_directory)
|
||||
total = 0
|
||||
for entry in published.manifest.files:
|
||||
total += _verified_job_file_size(
|
||||
self.local_job_directory,
|
||||
entry.payload_relative_path,
|
||||
entry.logical_bytes,
|
||||
)
|
||||
for artifact in published.manifest.artifacts:
|
||||
total += _verified_job_file_size(
|
||||
self.local_job_directory,
|
||||
artifact.payload_relative_path,
|
||||
artifact.logical_bytes,
|
||||
)
|
||||
declared_files = [
|
||||
(entry.payload_relative_path, entry.logical_bytes)
|
||||
for entry in published.manifest.files
|
||||
] + [
|
||||
(artifact.payload_relative_path, artifact.logical_bytes)
|
||||
for artifact in published.manifest.artifacts
|
||||
]
|
||||
total = sum(size for _, size in declared_files)
|
||||
|
||||
completion = self.transport.get_json(
|
||||
"/rest/db/completion?"
|
||||
@@ -178,11 +173,36 @@ class SyncthingTransferObserver:
|
||||
name for name in needed_names
|
||||
if name == self.job_relative_path or name.startswith(prefix)
|
||||
}
|
||||
fraction = float(raw_fraction) / 100
|
||||
complete = fraction == 1 and not relevant
|
||||
complete = not relevant
|
||||
if complete:
|
||||
for relative_path, expected_bytes in declared_files:
|
||||
_verified_job_file_size(
|
||||
self.local_job_directory,
|
||||
relative_path,
|
||||
expected_bytes,
|
||||
)
|
||||
completed_bytes = total
|
||||
else:
|
||||
observed_bytes = sum(
|
||||
_received_job_file_bytes(
|
||||
self.local_job_directory,
|
||||
relative_path,
|
||||
expected_bytes,
|
||||
)
|
||||
for relative_path, expected_bytes in declared_files
|
||||
)
|
||||
observed_fraction = observed_bytes / total if total else 1.0
|
||||
# Folder completion remains useful when Syncthing has already
|
||||
# atomically published a file but still reports it in its need
|
||||
# queue; the allocated-byte estimate is needed for sparse temp
|
||||
# files that are pre-sized before their blocks arrive.
|
||||
fraction = min(observed_fraction, float(raw_fraction) / 100)
|
||||
completed_bytes = int(total * fraction)
|
||||
if complete:
|
||||
fraction = 1.0
|
||||
return SyncthingTransferStatus(
|
||||
fraction,
|
||||
total if complete else int(total * fraction),
|
||||
completed_bytes,
|
||||
total,
|
||||
complete,
|
||||
len(relevant),
|
||||
@@ -472,6 +492,42 @@ def _verified_job_file_size(
|
||||
return metadata.st_size
|
||||
|
||||
|
||||
def _received_job_file_bytes(
|
||||
job_directory: Path,
|
||||
relative_path: str,
|
||||
expected_bytes: int,
|
||||
) -> int:
|
||||
"""Return a conservative receive estimate for one job-owned file.
|
||||
|
||||
Syncthing writes incomplete files as ``.syncthing.<name>.tmp`` and may
|
||||
pre-size that sparse temporary to its final logical length. Allocated
|
||||
blocks, rather than ``st_size``, therefore provide the useful progress
|
||||
signal until the final atomic rename occurs.
|
||||
"""
|
||||
|
||||
relative = PurePosixPath(relative_path)
|
||||
current = job_directory
|
||||
for component in relative.parts[:-1]:
|
||||
current = current / component
|
||||
final = current / relative.name
|
||||
try:
|
||||
metadata = final.lstat()
|
||||
except FileNotFoundError:
|
||||
metadata = None
|
||||
if metadata is not None:
|
||||
if stat.S_ISREG(metadata.st_mode) and metadata.st_size == expected_bytes:
|
||||
return expected_bytes
|
||||
return 0
|
||||
temporary = current / f".syncthing.{relative.name}.tmp"
|
||||
try:
|
||||
temporary_metadata = temporary.lstat()
|
||||
except FileNotFoundError:
|
||||
return 0
|
||||
if not stat.S_ISREG(temporary_metadata.st_mode):
|
||||
return 0
|
||||
return min(expected_bytes, temporary_metadata.st_blocks * 512)
|
||||
|
||||
|
||||
def _needed_names(value: dict[str, Any]) -> set[str]:
|
||||
result: set[str] = set()
|
||||
for key in ("progress", "queued", "rest"):
|
||||
|
||||
@@ -456,6 +456,36 @@ def cleanup_transfer(job_directory: Path) -> bool:
|
||||
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:
|
||||
value = json_format.MessageToDict(
|
||||
message,
|
||||
|
||||
@@ -111,6 +111,115 @@ class ClientJobHappyPathTests(unittest.TestCase):
|
||||
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 test_first_partial_progress_is_durable(self):
|
||||
"""A sparse transfer must not lose its only initial observation."""
|
||||
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,
|
||||
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=0,
|
||||
))
|
||||
|
||||
def partial_progress(_definition, _step, progress):
|
||||
progress(0.001, 10, 10_000, "still receiving")
|
||||
|
||||
with patch.object(executor, "_execute_step", partial_progress):
|
||||
events = executor.execute(control_pb2.ExecuteStepCommand(
|
||||
job_id=definition.job_id, expected_job_revision=1,
|
||||
expected_last_event_sequence=1,
|
||||
step=job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER, attempt=1,
|
||||
))
|
||||
|
||||
self.assertEqual(len(events), 3)
|
||||
self.assertEqual(events[1].type, control_pb2.JOB_EVENT_TYPE_PROGRESS)
|
||||
self.assertEqual(events[1].progress.bytes_complete, 10)
|
||||
|
||||
def _run_transfer(
|
||||
self,
|
||||
root: Path,
|
||||
|
||||
@@ -11,6 +11,7 @@ from archive_clients.transfer import (
|
||||
FileMaterializer,
|
||||
TransferIntegrityError,
|
||||
canonical_message_json,
|
||||
cleanup_partial_transfer,
|
||||
load_published_transfer,
|
||||
materialize_transfer,
|
||||
stage_transfer,
|
||||
@@ -206,6 +207,33 @@ class TransferHappyPathTests(unittest.TestCase):
|
||||
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):
|
||||
manifest = self._manifest()
|
||||
first = canonical_message_json(manifest)
|
||||
|
||||
Reference in New Issue
Block a user