Compare commits

...
7 Commits
13 changed files with 359 additions and 58 deletions
+31 -4
View File
@@ -13,10 +13,37 @@ initial x1/x2/lithium topology.
`archive_control_token`, `qb_password`, and `syncthing_api_key`, each a
regular non-empty file with mode `0600`.
The Syncthing mounts intentionally reproduce each instance's `/var/syncthing`
layout, including nested data binds. This lets route discovery and route
provisioning use one safe API-to-local path mapping without altering an
existing Syncthing configuration.
## Hardlink-safe bind-mount topology
For an archive source, the qB content path and every Syncthing route used for
staging must resolve through the **same container mount**. Matching host
filesystem device IDs alone is insufficient: two separate Docker bind mounts
have different mount IDs and `link(2)` may return `EXDEV` across them. The
client deliberately treats that case as copy-only and performs a full payload
free-space check.
When a Syncthing route is physically nested below the qB root, mount the qB
root once and map the exact Syncthing API folder through it:
```yaml
volumes:
- /srv/downloads:/data/qb
- /srv/syncthing-config:/data/sync
```
```toml
[syncthing]
api_root = "/var/syncthing"
local_root = "/data/sync"
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" }
```
Do **not** additionally mount `/srv/downloads/Sync` at a path beneath
`/data/sync`. The override is the authoritative mapping for that folder and
keeps qB source files and staging destinations in one mount namespace. Use
the folder's normalized API-visible **path** as the override key (for example,
Syncthing `~/Downloads/Sync` becomes `/var/syncthing/Downloads/Sync`); do not
use the folder ID.
Before starting a stack, validate it with:
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services:
archive-client:
image: sodium/archive-clients:v0.1.8
image: sodium/archive-clients:v0.1.13
user: "1000:1000"
restart: unless-stopped
command: ["--config", "/etc/archive-control/client.toml"]
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services:
archive-client:
image: sodium/archive-clients:v0.1.8
image: sodium/archive-clients:v0.1.13
user: "1001:1001"
restart: unless-stopped
network_mode: host
+4
View File
@@ -39,4 +39,8 @@ endpoint = "http://127.0.0.1:8384"
api_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing"
local_root = "/data/sync"
# `DownloadsSync-X2` is ~/Downloads/Sync on the host, nested below the qB
# root. Resolve it through /data/qb rather than a second nested bind mount so
# source staging can hardlink it.
local_path_overrides = { "/var/syncthing/Downloads/Sync" = "/data/qb/Sync" }
advertised_addresses = ["dynamic"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services:
archive-client:
image: sodium/archive-clients:v0.1.8
image: sodium/archive-clients:v0.1.13
user: "1001:1001"
restart: unless-stopped
network_mode: host
@@ -16,4 +16,3 @@ services:
- ./backups:/var/backups/archive-control
- /home/ubuntu/Downloads:/data/qb
- /home/ubuntu/compose/syncthing/st_home:/data/sync
- /home/ubuntu/Downloads/Sync:/data/sync/Downloads/Sync
+7 -2
View File
@@ -128,7 +128,12 @@ advertised_addresses = ["dynamic"]
When an existing Syncthing folder is physically nested in the qB data root,
use a `local_path_overrides` entry to map that exact Syncthing API path through
the same client bind mount. This enables hardlinks without creating two Docker
mount boundaries for the same host files.
mount boundaries for the same host files. Do not add a second bind mount for
the nested folder: Linux treats it as a distinct mount even when it has the
same `st_dev`, and the client correctly falls back to copy-only capacity
accounting. The override key is the normalized Syncthing folder path beneath
`api_root`, not its folder ID. See the production deployment README for the
required compose and override pattern.
The remaining node examples omit optional `[connection]`, `[jobs]`, and
`[backup]` tables and therefore use these same defaults; deployments may
@@ -218,7 +223,7 @@ cache/archive routes according to policy.
```yaml
services:
archive-client:
image: sodium/archive-clients:v0.1.8
image: sodium/archive-clients:v0.1.13
user: "1001:1001"
restart: unless-stopped
command: ["archive-client", "--config", "/etc/archive-control/client.toml"]
+5 -3
View File
@@ -89,9 +89,11 @@ left unchanged.
Syncthing may serialize a folder path relative to its home as `~/...`. The
client normalizes that notation beneath the configured API-visible sync root
before applying the API-to-local root mapping. Deployments must mirror
Syncthing's nested bind mounts into the client so the normalized API path and
the client filesystem path refer to the same bytes.
before applying the API-to-local root mapping. When that folder is nested
below qB's content root, an exact `local_path_overrides` entry must map it
through the qB bind mount. Do not mirror it as a second nested client bind
mount: it denotes the same host bytes but a distinct mount namespace boundary,
which prevents hardlink staging.
Provisioning uses idempotent device and folder configuration updates. Each
client receives its peer device ID and optional advertised addresses
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "archive-clients"
version = "0.1.8"
version = "0.1.13"
requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+18 -1
View File
@@ -11,6 +11,7 @@ from pathlib import Path, PurePosixPath
from typing import Any
from websockets.asyncio.client import connect
from websockets.exceptions import ConnectionClosedOK
from archive_clients.backup import SQLiteBackupManager
from archive_clients.config import ClientConfig
@@ -205,7 +206,23 @@ class ArchiveClientDaemon:
command_tasks: set[asyncio.Task[None]] = set()
try:
await self._resume_commands(outbound, command_tasks)
async for frame in websocket:
# A proxy can leave the TCP/WebSocket socket apparently open
# after the control server has discarded its session. The
# server then cannot deliver durable commands and its pending
# outbox remains stranded unless the client independently
# detects the missing application heartbeats and reconnects.
while True:
try:
frame = await asyncio.wait_for(
websocket.recv(),
self.config.connection.offline_timeout,
)
except ConnectionClosedOK:
return
except asyncio.TimeoutError as exc:
raise RuntimeError(
"control heartbeat timed out"
) from exc
await self._handle(decode(frame), outbound, command_tasks)
finally:
writer.cancel()
+54 -27
View File
@@ -120,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,
@@ -132,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,
@@ -167,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
+80 -16
View File
@@ -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,44 @@ 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
# ``/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(
fraction,
total if complete else int(total * fraction),
completed_bytes,
total,
complete,
len(relevant),
@@ -472,6 +500,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"):
+47
View File
@@ -22,6 +22,53 @@ from archive_control.v1 import (
class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
async def test_silent_control_connection_ends_for_reconnect(self):
"""A lost server heartbeat must not leave durable commands stranded."""
async def control(websocket):
registration = decode(await websocket.recv())
self.assertEqual(registration.WhichOneof("payload"), "register_request")
response = new_envelope()
response.correlation_id = registration.message_id
response.register_response.status = (
client_pb2.REGISTRATION_STATUS_ACCEPTED
)
response.register_response.negotiated_version.major = 1
await websocket.send(encode(response))
# Deliberately keep TCP/WebSocket open but send no application
# heartbeats. This models a stale proxy/server-side session.
await websocket.wait_closed()
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
token = root / "token"
token.write_text("shared-secret", encoding="utf-8")
os.chmod(token, 0o600)
async with serve(control, "127.0.0.1", 0, ping_interval=None) as server:
port = server.sockets[0].getsockname()[1]
service = ServiceConfig(
"http://local", PurePosixPath("/api"), root,
)
config = ClientConfig(
"cache-1", "Cache 1", "cache",
f"ws://127.0.0.1:{port}", token,
root / "state.db", root / "backups", service, service,
ConnectionConfig(
registration_timeout=1,
offline_timeout=0.05,
reconnect_initial=0.01,
reconnect_max=0.01,
reconnect_jitter=False,
),
)
probe = FilesystemProbe(root, True, True, True, True, True)
daemon = ArchiveClientDaemon(config, [probe, probe], [])
await asyncio.to_thread(daemon.store.initialize)
with self.assertRaisesRegex(
RuntimeError, "control heartbeat timed out"
):
await asyncio.wait_for(daemon._connection(), 1)
async def test_eviction_assignment_and_steps_are_admitted(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
+109
View File
@@ -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,