Compare commits

..
5 Commits
16 changed files with 298 additions and 26 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 `archive_control_token`, `qb_password`, and `syncthing_api_key`, each a
regular non-empty file with mode `0600`. regular non-empty file with mode `0600`.
The Syncthing mounts intentionally reproduce each instance's `/var/syncthing` ## Hardlink-safe bind-mount topology
layout, including nested data binds. This lets route discovery and route
provisioning use one safe API-to-local path mapping without altering an For an archive source, the qB content path and every Syncthing route used for
existing Syncthing configuration. 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: Before starting a stack, validate it with:
+1 -1
View File
@@ -2,7 +2,7 @@ name: archive-control-archive
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.12 image: sodium/archive-clients:v0.1.14
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.12 image: sodium/archive-clients:v0.1.14
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host 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_key_file = "/run/secrets/syncthing_api_key"
api_root = "/var/syncthing" api_root = "/var/syncthing"
local_root = "/data/sync" 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"] advertised_addresses = ["dynamic"]
+1 -2
View File
@@ -2,7 +2,7 @@ name: archive-control-cache
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.12 image: sodium/archive-clients:v0.1.14
user: "1001:1001" user: "1001:1001"
restart: unless-stopped restart: unless-stopped
network_mode: host network_mode: host
@@ -16,4 +16,3 @@ services:
- ./backups:/var/backups/archive-control - ./backups:/var/backups/archive-control
- /home/ubuntu/Downloads:/data/qb - /home/ubuntu/Downloads:/data/qb
- /home/ubuntu/compose/syncthing/st_home:/data/sync - /home/ubuntu/compose/syncthing/st_home:/data/sync
- /home/ubuntu/Downloads/Sync:/data/sync/Downloads/Sync
+5 -2
View File
@@ -35,8 +35,11 @@ expand these rules but must not contradict them.
manifests, and observed external state. manifests, and observed external state.
- SQLite backups use the online backup API. Defaults are every six hours, 12 - SQLite backups use the online backup API. Defaults are every six hours, 12
recent, 14 daily, and eight weekly copies, plus pre/post-migration backups. recent, 14 daily, and eight weekly copies, plus pre/post-migration backups.
- Default concurrency is one active data-moving job per client and per route. - Default concurrency is unrestricted: every eligible queued job is admitted
Disjoint node pairs may run concurrently. Queueing is durable FIFO. in one scheduler pass and independent jobs execute concurrently on a client.
Physical qBittorrent, Syncthing, disk, and network capacity are therefore
the natural limit. `enforce_concurrency_limits=true` restores the optional
legacy per-client/per-route gates. Commands for one job remain serialized.
- A queued job owns a per-resource reservation. Cancelling it removes only the - A queued job owns a per-resource reservation. Cancelling it removes only the
queued record/reservation and never sends cleanup commands. queued record/reservation and never sends cleanup commands.
- Connectivity or transfer stalls wait indefinitely. A configurable 30-minute - Connectivity or transfer stalls wait indefinitely. A configurable 30-minute
+10 -2
View File
@@ -44,6 +44,9 @@ command_max_attempts = 3 # Total sends, including the initial attempt.
stall_after = "30m" # Warning state only; jobs continue waiting. stall_after = "30m" # Warning state only; jobs continue waiting.
route_policy = "on_demand" # Alternative: eager_mesh. route_policy = "on_demand" # Alternative: eager_mesh.
route_setup_timeout = "30m" route_setup_timeout = "30m"
# Default: allow all eligible jobs; disk/network/service capacity is the limit.
enforce_concurrency_limits = false
# Used only when enforce_concurrency_limits is true.
max_active_per_client = 1 max_active_per_client = 1
max_active_per_route = 1 max_active_per_route = 1
max_envelope_bytes = 1048576 max_envelope_bytes = 1048576
@@ -128,7 +131,12 @@ advertised_addresses = ["dynamic"]
When an existing Syncthing folder is physically nested in the qB data root, 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 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 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 The remaining node examples omit optional `[connection]`, `[jobs]`, and
`[backup]` tables and therefore use these same defaults; deployments may `[backup]` tables and therefore use these same defaults; deployments may
@@ -218,7 +226,7 @@ cache/archive routes according to policy.
```yaml ```yaml
services: services:
archive-client: archive-client:
image: sodium/archive-clients:v0.1.12 image: sodium/archive-clients:v0.1.14
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"]
+5 -3
View File
@@ -89,9 +89,11 @@ left unchanged.
Syncthing may serialize a folder path relative to its home as `~/...`. The Syncthing may serialize a folder path relative to its home as `~/...`. The
client normalizes that notation beneath the configured API-visible sync root client normalizes that notation beneath the configured API-visible sync root
before applying the API-to-local root mapping. Deployments must mirror before applying the API-to-local root mapping. When that folder is nested
Syncthing's nested bind mounts into the client so the normalized API path and below qB's content root, an exact `local_path_overrides` entry must map it
the client filesystem path refer to the same bytes. 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 Provisioning uses idempotent device and folder configuration updates. Each
client receives its peer device ID and optional advertised addresses client receives its peer device ID and optional advertised addresses
+13 -4
View File
@@ -133,10 +133,13 @@ a resource may be active or queued at a time. This prevents incompatible
baselines even when jobs would use different nodes or observe different sides baselines even when jobs would use different nodes or observe different sides
of a hybrid identity. of a hybrid identity.
Defaults allow one active data-moving job per client and one per route. A job By default every eligible queued job is claimed in the same scheduler pass and
must acquire its source client, target client, route, and resource reservation independent jobs run concurrently on each client; physical service and storage
atomically. Disjoint node pairs may run concurrently. Route setup is a capacity are the limit. Set `enforce_concurrency_limits=true` on control to
preflight activity and does not permit a data step to bypass these leases. enable the optional one-per-client/one-per-route gates. A job always retains
its resource reservation, and commands for the same job remain serialized.
Route setup is a preflight activity and does not permit a data step to bypass
the resource reservation.
Offline nodes do not prevent unrelated jobs from being listed or run. A job Offline nodes do not prevent unrelated jobs from being listed or run. A job
requiring an offline node remains waiting indefinitely; it does not consume an requiring an offline node remains waiting indefinitely; it does not consume an
@@ -154,6 +157,12 @@ advances to the next step. On reconnect:
4. Control observes relevant qBittorrent, Syncthing, staging, and manifest 4. Control observes relevant qBittorrent, Syncthing, staging, and manifest
state before selecting retry, resume, compensation, cleanup, or manual state before selecting retry, resume, compensation, cleanup, or manual
intervention. intervention.
If a connection disappears while a synchronous job operation is in progress,
the client journals events locally, reconnects indefinitely, and serializes a
replayed command behind the in-flight operation for that job. The replay reads
the journal rather than repeating the data operation. This applies to every
transfer step, including post-commit staging cleanup.
5. A command is reissued with its original ID when the acceptance result is 5. A command is reissued with its original ID when the acceptance result is
uncertain. uncertain.
+4 -2
View File
@@ -25,8 +25,10 @@ What do you want to do?
[ Evict Cache ] [ Job Status ] [ Evict Cache ] [ Job Status ]
``` ```
All subsequent pages edit this message. `Cancel` closes the active selection All subsequent pages edit this message. Leaving a resource/tree selection
flow. `Back` returns one level while retaining validated filters/selections. returns to the Archive Control operation menu without closing the shared
conversation. Confirmation `Cancel` returns to its immediate prior selection;
the main menu alone can close Archive Control.
Callback payloads contain opaque session/action IDs, not resource names or Callback payloads contain opaque session/action IDs, not resource names or
paths, and are validated against persisted session revision and expiry. paths, and are validated against persisted session revision and expiry.
+1 -1
View File
@@ -225,7 +225,7 @@ Fixtures include small deterministic v1, v2, and hybrid torrents with:
| Recovery | restart every step on source/target/control, lost DB with proof, ambiguous loss fails closed | | Recovery | restart every step on source/target/control, lost DB with proof, ambiguous loss fails closed |
| Messaging | duplicate command/event, lost ack, sequence gap, stale revision, duplicate client ID | | Messaging | duplicate command/event, lost ack, sequence gap, stale revision, duplicate client ID |
| Safety | attempted download, corrupt same-size file, symlink/path escape, special file, partfile mismatch | | Safety | attempted download, corrupt same-size file, symlink/path escape, special file, partfile mismatch |
| Operations | stall retains percent, offline wait, cancellation each phase, postcommit cleanup retry | | Operations | stall retains percent, offline wait, cancellation each phase, postcommit cleanup retry, disconnect/reconnect during every transfer step with exactly-once replay |
| Database | online backup under load, pre/post migration, retention, corrupt backup rejection, offline restore | | Database | online backup under load, pre/post migration, retention, corrupt backup rejection, offline restore |
| UI | pagination/filter chain, stale/double callback, persisted session, clear vs evict separation | | UI | pagination/filter chain, stale/double callback, persisted session, clear vs evict separation |
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project] [project]
name = "archive-clients" name = "archive-clients"
version = "0.1.12" version = "0.1.14"
requires-python = ">=3.11" requires-python = ">=3.11"
dependencies = ["protobuf==7.35.1", "websockets==16.0"] dependencies = ["protobuf==7.35.1", "websockets==16.0"]
+38 -3
View File
@@ -11,6 +11,7 @@ from pathlib import Path, PurePosixPath
from typing import Any from typing import Any
from websockets.asyncio.client import connect from websockets.asyncio.client import connect
from websockets.exceptions import ConnectionClosedOK
from archive_clients.backup import SQLiteBackupManager from archive_clients.backup import SQLiteBackupManager
from archive_clients.config import ClientConfig from archive_clients.config import ClientConfig
@@ -42,6 +43,19 @@ from archive_control.v1 import (
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
def _job_id_for_command(command: control_pb2.Command) -> str:
"""Return the durable job key for commands whose execution is serialized."""
payload = command.WhichOneof("payload")
if payload == "assign_job":
return command.assign_job.job.job_id
if payload == "execute_step":
return command.execute_step.job_id
if payload == "cancel_job":
return command.cancel_job.job_id
raise ValueError(f"command {command.command_id} does not execute a job")
class ArchiveClientDaemon: class ArchiveClientDaemon:
def __init__( def __init__(
self, self,
@@ -94,7 +108,10 @@ class ArchiveClientDaemon:
self._lease = DatabaseLease(config.state_db) self._lease = DatabaseLease(config.state_db)
self._active_route_commands: set[str] = set() self._active_route_commands: set[str] = set()
self._active_job_commands: set[str] = set() self._active_job_commands: set[str] = set()
self._job_execution_lock = asyncio.Lock() # Commands for one job remain ordered locally, while unrelated jobs
# may use the node's available qB/Syncthing/filesystem capacity in
# parallel. The control daemon owns admission policy.
self._job_execution_locks: dict[str, asyncio.Lock] = {}
self.jobs = ( self.jobs = (
ClientJobExecutor( ClientJobExecutor(
client_id=config.client_id, client_id=config.client_id,
@@ -205,7 +222,23 @@ class ArchiveClientDaemon:
command_tasks: set[asyncio.Task[None]] = set() command_tasks: set[asyncio.Task[None]] = set()
try: try:
await self._resume_commands(outbound, command_tasks) 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) await self._handle(decode(frame), outbound, command_tasks)
finally: finally:
writer.cancel() writer.cancel()
@@ -543,7 +576,9 @@ class ArchiveClientDaemon:
correlation_id: str, correlation_id: str,
outbound: asyncio.Queue[str], outbound: asyncio.Queue[str],
) -> None: ) -> None:
async with self._job_execution_lock: job_id = _job_id_for_command(command)
lock = self._job_execution_locks.setdefault(job_id, asyncio.Lock())
async with lock:
await self._execute_job_command_locked( await self._execute_job_command_locked(
command, correlation_id, outbound command, correlation_id, outbound
) )
+18
View File
@@ -86,6 +86,8 @@ class ClientJobExecutor:
self.verification_timeout = verification_timeout self.verification_timeout = verification_timeout
self.free_space_reserve_bytes = free_space_reserve_bytes self.free_space_reserve_bytes = free_space_reserve_bytes
self._cancel_events: dict[str, threading.Event] = {} self._cancel_events: dict[str, threading.Event] = {}
self._execution_locks: dict[str, threading.Lock] = {}
self._execution_locks_guard = threading.Lock()
def request_cancel(self, job_id: str) -> None: def request_cancel(self, job_id: str) -> None:
self._cancel_events.setdefault(job_id, threading.Event()).set() self._cancel_events.setdefault(job_id, threading.Event()).set()
@@ -115,6 +117,22 @@ class ClientJobExecutor:
self, self,
command: control_pb2.ExecuteStepCommand, command: control_pb2.ExecuteStepCommand,
event_callback: Callable[[control_pb2.JobEvent], None] | None = None, event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
) -> list[control_pb2.JobEvent]:
# asyncio cancellation of a connection-bound task cannot stop the
# synchronous filesystem/qB operation already running in its worker
# thread. A replay after reconnect therefore waits for that operation
# and then reads its durable event journal instead of executing twice.
with self._execution_lock(command.job_id):
return self._execute_locked(command, event_callback)
def _execution_lock(self, job_id: str) -> threading.Lock:
with self._execution_locks_guard:
return self._execution_locks.setdefault(job_id, threading.Lock())
def _execute_locked(
self,
command: control_pb2.ExecuteStepCommand,
event_callback: Callable[[control_pb2.JobEvent], None] | None = None,
) -> list[control_pb2.JobEvent]: ) -> list[control_pb2.JobEvent]:
definition = self._definition(command.job_id) definition = self._definition(command.job_id)
replay = self._replay( replay = self._replay(
+88
View File
@@ -22,6 +22,94 @@ from archive_control.v1 import (
class DaemonTransportTests(unittest.IsolatedAsyncioTestCase): class DaemonTransportTests(unittest.IsolatedAsyncioTestCase):
async def test_unrelated_job_commands_execute_concurrently(self):
with tempfile.TemporaryDirectory() as directory:
root = Path(directory)
token = root / "token"
token.write_text("shared-secret", encoding="utf-8")
os.chmod(token, 0o600)
service = ServiceConfig("http://local", PurePosixPath("/api"), root)
config = ClientConfig(
"cache-1", "Cache 1", "cache", "ws://control", token,
root / "state.db", root / "backups", service, service,
)
probe = FilesystemProbe(root, True, True, True, True, True)
daemon = ArchiveClientDaemon(config, [probe, probe], [])
started: set[str] = set()
release = asyncio.Event()
async def execute(command, correlation_id, outbound):
started.add(command.assign_job.job.job_id)
await release.wait()
daemon._execute_job_command_locked = execute
commands = []
for _ in range(2):
command = control_pb2.Command(command_id=str(uuid4()))
command.assign_job.job.job_id = str(uuid4())
commands.append(command)
outbound = asyncio.Queue()
tasks = [
asyncio.create_task(
daemon._execute_job_command(command, "", outbound)
)
for command in commands
]
for _ in range(100):
if len(started) == 2:
break
await asyncio.sleep(0.01)
self.assertEqual(len(started), 2)
release.set()
await asyncio.gather(*tasks)
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): async def test_eviction_assignment_and_steps_are_admitted(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
root = Path(directory) root = Path(directory)
+77
View File
@@ -1,5 +1,6 @@
import hashlib import hashlib
import tempfile import tempfile
import threading
import unittest import unittest
from pathlib import Path from pathlib import Path
from unittest.mock import Mock, patch from unittest.mock import Mock, patch
@@ -38,6 +39,82 @@ class SlowRescanSyncthing(CompleteSyncthing):
class ClientJobHappyPathTests(unittest.TestCase): class ClientJobHappyPathTests(unittest.TestCase):
def test_reconnect_replay_never_duplicates_any_transfer_step(self):
for step in (
job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
job_pb2.JOB_STEP_KIND_SYNCTHING_TRANSFER,
job_pb2.JOB_STEP_KIND_TARGET_MATERIALIZE,
job_pb2.JOB_STEP_KIND_QB_VERIFY,
job_pb2.JOB_STEP_KIND_STAGING_CLEANUP,
):
with self.subTest(step=step), 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="reconnect 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="cache-1",
qbittorrent=Mock(),
store=store,
qb_root=root,
qb_api_root=Path("/downloads"),
route_path=lambda _: root,
syncthing_transport=Mock(),
sparse_supported=True,
)
executor.assign(control_pb2.AssignJobCommand(
job=definition,
expected_job_revision=1,
expected_last_event_sequence=0,
))
entered = threading.Event()
release = threading.Event()
def execute_step(*_args):
entered.set()
release.wait(1)
return None
executor._execute_step = Mock(side_effect=execute_step)
command = control_pb2.ExecuteStepCommand(
job_id=definition.job_id,
expected_job_revision=1,
expected_last_event_sequence=1,
step=step,
attempt=1,
)
results: list[list[control_pb2.JobEvent]] = []
first = threading.Thread(
target=lambda: results.append(executor.execute(command))
)
second = threading.Thread(
target=lambda: results.append(executor.execute(command))
)
first.start()
self.assertTrue(entered.wait(1))
second.start()
release.set()
first.join(1)
second.join(1)
self.assertFalse(first.is_alive())
self.assertFalse(second.is_alive())
self.assertEqual(executor._execute_step.call_count, 1)
self.assertEqual(len(results), 2)
self.assertEqual(results[0][-1].event_id, results[1][-1].event_id)
def test_capacity_guard_fails_before_data_movement(self): def test_capacity_guard_fails_before_data_movement(self):
with tempfile.TemporaryDirectory() as directory: with tempfile.TemporaryDirectory() as directory:
root = Path(directory) root = Path(directory)