Compare commits

...
4 Commits
8 changed files with 248 additions and 48 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.8 image: sodium/archive-clients:v0.1.12
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.8 image: sodium/archive-clients:v0.1.12
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.8 image: sodium/archive-clients:v0.1.12
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.8 image: sodium/archive-clients:v0.1.12
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.8" version = "0.1.12"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+54 -27
View File
@@ -120,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,
@@ -132,32 +136,51 @@ 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 # A newly-started step must publish its first observation. In
speed_sample = [time.monotonic(), 0] # 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( def progress(
fraction: float, fraction: float,
@@ -167,8 +190,12 @@ class ClientJobExecutor:
) -> None: ) -> None:
nonlocal cursor nonlocal cursor
now = time.monotonic() now = time.monotonic()
elapsed = now - float(speed_sample[0]) previous_time = speed_sample[0]
if fraction < 1 and elapsed < 1: 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 return
speed = ( speed = (
max(0, bytes_complete - int(speed_sample[1])) / elapsed max(0, bytes_complete - int(speed_sample[1])) / elapsed
+80 -16
View File
@@ -139,19 +139,14 @@ class SyncthingTransferObserver:
def status(self) -> SyncthingTransferStatus: def status(self) -> SyncthingTransferStatus:
published = load_published_transfer(self.local_job_directory) published = load_published_transfer(self.local_job_directory)
total = 0 declared_files = [
for entry in published.manifest.files: (entry.payload_relative_path, entry.logical_bytes)
total += _verified_job_file_size( for entry in published.manifest.files
self.local_job_directory, ] + [
entry.payload_relative_path, (artifact.payload_relative_path, artifact.logical_bytes)
entry.logical_bytes, for artifact in published.manifest.artifacts
) ]
for artifact in published.manifest.artifacts: total = sum(size for _, size in declared_files)
total += _verified_job_file_size(
self.local_job_directory,
artifact.payload_relative_path,
artifact.logical_bytes,
)
completion = self.transport.get_json( completion = self.transport.get_json(
"/rest/db/completion?" "/rest/db/completion?"
@@ -178,11 +173,44 @@ class SyncthingTransferObserver:
name for name in needed_names name for name in needed_names
if name == self.job_relative_path or name.startswith(prefix) if name == self.job_relative_path or name.startswith(prefix)
} }
fraction = float(raw_fraction) / 100 complete = not relevant
complete = fraction == 1 and 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
# ``/db/completion`` is folder-wide. It can be near zero when a
# route contains a freshly-created item even though the tracked
# job's temporary payload has already received many blocks. For
# an incomplete temporary file, its allocated blocks are the only
# job-specific signal, so do not cap them with that unrelated
# folder aggregate. Once every declared file is atomically
# present, retain the completion value as a conservative guard
# until the need queue has caught up.
fraction = (
observed_fraction
if observed_fraction < 1.0
else min(1.0, float(raw_fraction) / 100)
)
completed_bytes = int(total * fraction)
if complete:
fraction = 1.0
return SyncthingTransferStatus( return SyncthingTransferStatus(
fraction, fraction,
total if complete else int(total * fraction), completed_bytes,
total, total,
complete, complete,
len(relevant), len(relevant),
@@ -472,6 +500,42 @@ def _verified_job_file_size(
return metadata.st_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]: def _needed_names(value: dict[str, Any]) -> set[str]:
result: set[str] = set() result: set[str] = set()
for key in ("progress", "queued", "rest"): for key in ("progress", "queued", "rest"):
+109
View File
@@ -111,6 +111,115 @@ 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 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( def _run_transfer(
self, self,
root: Path, root: Path,