From 8ed8e75447710aa0b79176424ab30da54ef27b91 Mon Sep 17 00:00:00 2001 From: Cabbagec Date: Mon, 27 Jul 2026 03:19:47 +0000 Subject: [PATCH] feat: allow unrestricted independent job concurrency --- docs/decisions.md | 7 +++- docs/deployment-and-usage.md | 3 ++ docs/system-design.md | 17 ++++++-- docs/telegram-ux.md | 6 ++- docs/testing.md | 2 +- src/archive_clients/daemon.py | 22 +++++++++- src/archive_clients/jobs.py | 18 ++++++++ tests/test_daemon.py | 41 +++++++++++++++++++ tests/test_jobs.py | 77 +++++++++++++++++++++++++++++++++++ 9 files changed, 182 insertions(+), 11 deletions(-) diff --git a/docs/decisions.md b/docs/decisions.md index e79ff11..d0f8afa 100644 --- a/docs/decisions.md +++ b/docs/decisions.md @@ -35,8 +35,11 @@ expand these rules but must not contradict them. manifests, and observed external state. - 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. -- Default concurrency is one active data-moving job per client and per route. - Disjoint node pairs may run concurrently. Queueing is durable FIFO. +- Default concurrency is unrestricted: every eligible queued job is admitted + 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 queued record/reservation and never sends cleanup commands. - Connectivity or transfer stalls wait indefinitely. A configurable 30-minute diff --git a/docs/deployment-and-usage.md b/docs/deployment-and-usage.md index bf24a7b..f88a53a 100644 --- a/docs/deployment-and-usage.md +++ b/docs/deployment-and-usage.md @@ -44,6 +44,9 @@ command_max_attempts = 3 # Total sends, including the initial attempt. stall_after = "30m" # Warning state only; jobs continue waiting. route_policy = "on_demand" # Alternative: eager_mesh. 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_route = 1 max_envelope_bytes = 1048576 diff --git a/docs/system-design.md b/docs/system-design.md index 79a52bb..ce54ae3 100644 --- a/docs/system-design.md +++ b/docs/system-design.md @@ -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 of a hybrid identity. -Defaults allow one active data-moving job per client and one per route. A job -must acquire its source client, target client, route, and resource reservation -atomically. Disjoint node pairs may run concurrently. Route setup is a -preflight activity and does not permit a data step to bypass these leases. +By default every eligible queued job is claimed in the same scheduler pass and +independent jobs run concurrently on each client; physical service and storage +capacity are the limit. Set `enforce_concurrency_limits=true` on control to +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 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 state before selecting retry, resume, compensation, cleanup, or manual 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 uncertain. diff --git a/docs/telegram-ux.md b/docs/telegram-ux.md index 2fc87ff..2adf821 100644 --- a/docs/telegram-ux.md +++ b/docs/telegram-ux.md @@ -25,8 +25,10 @@ What do you want to do? [ Evict Cache ] [ Job Status ] ``` -All subsequent pages edit this message. `Cancel` closes the active selection -flow. `Back` returns one level while retaining validated filters/selections. +All subsequent pages edit this message. Leaving a resource/tree selection +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 paths, and are validated against persisted session revision and expiry. diff --git a/docs/testing.md b/docs/testing.md index 704171c..fdf0ff2 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -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 | | 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 | -| 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 | | UI | pagination/filter chain, stale/double callback, persisted session, clear vs evict separation | diff --git a/src/archive_clients/daemon.py b/src/archive_clients/daemon.py index 758f47c..0a6f2fd 100644 --- a/src/archive_clients/daemon.py +++ b/src/archive_clients/daemon.py @@ -43,6 +43,19 @@ from archive_control.v1 import ( 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: def __init__( self, @@ -95,7 +108,10 @@ class ArchiveClientDaemon: self._lease = DatabaseLease(config.state_db) self._active_route_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 = ( ClientJobExecutor( client_id=config.client_id, @@ -560,7 +576,9 @@ class ArchiveClientDaemon: correlation_id: str, outbound: asyncio.Queue[str], ) -> 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( command, correlation_id, outbound ) diff --git a/src/archive_clients/jobs.py b/src/archive_clients/jobs.py index 82ef3a6..8c59779 100644 --- a/src/archive_clients/jobs.py +++ b/src/archive_clients/jobs.py @@ -86,6 +86,8 @@ class ClientJobExecutor: self.verification_timeout = verification_timeout self.free_space_reserve_bytes = free_space_reserve_bytes 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: self._cancel_events.setdefault(job_id, threading.Event()).set() @@ -115,6 +117,22 @@ class ClientJobExecutor: self, command: control_pb2.ExecuteStepCommand, 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]: definition = self._definition(command.job_id) replay = self._replay( diff --git a/tests/test_daemon.py b/tests/test_daemon.py index ed64fbd..bd00917 100644 --- a/tests/test_daemon.py +++ b/tests/test_daemon.py @@ -22,6 +22,47 @@ from archive_control.v1 import ( 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.""" diff --git a/tests/test_jobs.py b/tests/test_jobs.py index 19f2722..eef895f 100644 --- a/tests/test_jobs.py +++ b/tests/test_jobs.py @@ -1,5 +1,6 @@ import hashlib import tempfile +import threading import unittest from pathlib import Path from unittest.mock import Mock, patch @@ -38,6 +39,82 @@ class SlowRescanSyncthing(CompleteSyncthing): 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): with tempfile.TemporaryDirectory() as directory: root = Path(directory)