feat: allow unrestricted independent job concurrency
This commit is contained in:
+5
-2
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
|
|||||||
+13
-4
@@ -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
@@ -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
@@ -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 |
|
||||||
|
|
||||||
|
|||||||
@@ -43,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,
|
||||||
@@ -95,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,
|
||||||
@@ -560,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
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -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(
|
||||||
|
|||||||
@@ -22,6 +22,47 @@ 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):
|
async def test_silent_control_connection_ends_for_reconnect(self):
|
||||||
"""A lost server heartbeat must not leave durable commands stranded."""
|
"""A lost server heartbeat must not leave durable commands stranded."""
|
||||||
|
|
||||||
|
|||||||
@@ -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)
|
||||||
|
|||||||
Reference in New Issue
Block a user