Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
515599f2d4 | ||
|
|
b67ca18403 | ||
|
|
6446846ee8 | ||
|
|
c717aea394 | ||
|
|
ea35ba2758 | ||
|
|
5e5e2f9095 | ||
|
|
0b11e2a3a2 |
+6
-2
@@ -11,8 +11,12 @@ CMD ["python", "-m", "unittest", "discover", "-s", "tests", "-v"]
|
|||||||
|
|
||||||
FROM python:3.11-slim
|
FROM python:3.11-slim
|
||||||
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
|
ENV PYTHONDONTWRITEBYTECODE=1 PYTHONUNBUFFERED=1
|
||||||
RUN groupadd --gid 1001 archive-control \
|
RUN if ! getent group 1001 >/dev/null; then \
|
||||||
&& useradd --uid 1001 --gid 1001 --no-create-home archive-control
|
groupadd --gid 1001 archive-control; \
|
||||||
|
fi \
|
||||||
|
&& if ! getent passwd 1001 >/dev/null; then \
|
||||||
|
useradd --uid 1001 --gid 1001 --no-create-home archive-control; \
|
||||||
|
fi
|
||||||
COPY --from=builder /wheels /wheels
|
COPY --from=builder /wheels /wheels
|
||||||
RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels
|
RUN pip install --no-cache-dir /wheels/*.whl && rm -rf /wheels
|
||||||
USER 1001:1001
|
USER 1001:1001
|
||||||
|
|||||||
@@ -259,6 +259,28 @@ Generated protobuf Python bindings are committed in each consumer with the
|
|||||||
exact `archive-control-proto` tag/commit recorded. Release order is proto,
|
exact `archive-control-proto` tag/commit recorded. Release order is proto,
|
||||||
control consumer, client consumer, E2E, then image publication.
|
control consumer, client consumer, E2E, then image publication.
|
||||||
|
|
||||||
|
### Reproducible multi-platform Buildx lifecycle
|
||||||
|
|
||||||
|
The named Buildx builders are local acceleration/cache only; they are not a
|
||||||
|
deployment dependency and can be removed after publication. To create a fresh
|
||||||
|
builder, verify its platforms, publish a release, and remove it afterwards:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
docker buildx create --name archive-control-release --driver docker-container --use
|
||||||
|
docker buildx inspect --bootstrap
|
||||||
|
docker buildx build --platform linux/amd64,linux/arm64 \
|
||||||
|
--tag sodium/archive-clients:vX.Y.Z --push .
|
||||||
|
docker buildx rm archive-control-release
|
||||||
|
```
|
||||||
|
|
||||||
|
`docker buildx inspect` must list both `linux/amd64` and `linux/arm64` before
|
||||||
|
publishing. If the host has no arm64 emulation, install/configure it according
|
||||||
|
to the host Docker distribution before the build; do not publish a partial
|
||||||
|
single-platform tag. Retain the pushed manifest digest in the release notes
|
||||||
|
and deploy the immutable tag or digest. The optional `archive-control-qemu`
|
||||||
|
builder follows the same lifecycle when it is used for an emulation smoke
|
||||||
|
build.
|
||||||
|
|
||||||
## Operator usage
|
## Operator usage
|
||||||
|
|
||||||
1. Prepare local qB, sync, state, backup, config, and secret mounts with the
|
1. Prepare local qB, sync, state, backup, config, and secret mounts with the
|
||||||
|
|||||||
+1
-1
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "archive-clients"
|
name = "archive-clients"
|
||||||
version = "0.1.14"
|
version = "0.1.19"
|
||||||
requires-python = ">=3.11"
|
requires-python = ">=3.11"
|
||||||
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
dependencies = ["protobuf==7.35.1", "websockets==16.0"]
|
||||||
|
|
||||||
|
|||||||
@@ -168,7 +168,7 @@ def _v1_files(info: dict[bytes, Any]) -> list[MetaFile]:
|
|||||||
components = [name] + [_component(part) for part in raw_path]
|
components = [name] + [_component(part) for part in raw_path]
|
||||||
result.append(MetaFile(
|
result.append(MetaFile(
|
||||||
"/".join(components), _length(item.get(b"length")),
|
"/".join(components), _length(item.get(b"length")),
|
||||||
_padding(item),
|
_padding(item) or _bitcomet_padding_name(components[-1]),
|
||||||
))
|
))
|
||||||
return result
|
return result
|
||||||
|
|
||||||
@@ -223,3 +223,13 @@ def _length(value: Any) -> int:
|
|||||||
def _padding(value: dict[bytes, Any]) -> bool:
|
def _padding(value: dict[bytes, Any]) -> bool:
|
||||||
attributes = value.get(b"attr", b"")
|
attributes = value.get(b"attr", b"")
|
||||||
return isinstance(attributes, bytes) and b"p" in attributes
|
return isinstance(attributes, bytes) and b"p" in attributes
|
||||||
|
|
||||||
|
|
||||||
|
def _bitcomet_padding_name(component: str) -> bool:
|
||||||
|
"""Recognize BitComet's legacy padding-file convention.
|
||||||
|
|
||||||
|
Such v1 torrents often omit the standard ``attr=p`` flag, but
|
||||||
|
qBittorrent/libtorrent still hides these synthetic entries from its file
|
||||||
|
API. The exact reserved prefix is the interoperable marker.
|
||||||
|
"""
|
||||||
|
return component.startswith("_____padding_file_")
|
||||||
|
|||||||
@@ -671,6 +671,15 @@ class ClientJobExecutor:
|
|||||||
# Syncthing can expose the per-job directory before the
|
# Syncthing can expose the per-job directory before the
|
||||||
# ready marker and manifest have arrived atomically as a set.
|
# ready marker and manifest have arrived atomically as a set.
|
||||||
pass
|
pass
|
||||||
|
except RouteSetupError as exc:
|
||||||
|
# A network interruption can temporarily make the local
|
||||||
|
# Syncthing REST API unavailable. The staged payload is
|
||||||
|
# durable and the job must remain recoverable; keep polling
|
||||||
|
# instead of turning a transient outage into a failed job.
|
||||||
|
logger.warning("syncthing_status_deferred", extra={
|
||||||
|
"job_id": definition.job_id,
|
||||||
|
"error_type": type(exc).__name__,
|
||||||
|
})
|
||||||
time.sleep(self.poll_interval)
|
time.sleep(self.poll_interval)
|
||||||
|
|
||||||
def _target_materialize(
|
def _target_materialize(
|
||||||
|
|||||||
@@ -3,6 +3,7 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
import logging
|
||||||
import threading
|
import threading
|
||||||
import time
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
@@ -16,6 +17,8 @@ from urllib import error, parse, request
|
|||||||
from archive_clients.config import ServiceConfig
|
from archive_clients.config import ServiceConfig
|
||||||
from archive_clients.resources import NormalizedResource, normalize_resource
|
from archive_clients.resources import NormalizedResource, normalize_resource
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class QBittorrentError(RuntimeError):
|
class QBittorrentError(RuntimeError):
|
||||||
pass
|
pass
|
||||||
@@ -72,7 +75,17 @@ class QBittorrentReader:
|
|||||||
torrent_hash = torrent.get("hash")
|
torrent_hash = torrent.get("hash")
|
||||||
if not isinstance(torrent_hash, str):
|
if not isinstance(torrent_hash, str):
|
||||||
raise QBittorrentError("qBittorrent torrent hash is invalid")
|
raise QBittorrentError("qBittorrent torrent hash is invalid")
|
||||||
result.append(self._normalize(torrent, torrent_hash))
|
try:
|
||||||
|
result.append(self._normalize(torrent, torrent_hash))
|
||||||
|
except ValueError as exc:
|
||||||
|
# A stale or malformed qBittorrent entry must not hide every
|
||||||
|
# otherwise valid resource from Archive/Evict inventory.
|
||||||
|
# Keep it ineligible and leave an operator-visible diagnosis.
|
||||||
|
logger.warning(
|
||||||
|
"Skipping qBittorrent resource with inconsistent "
|
||||||
|
"metadata: hash=%s name=%r reason=%s",
|
||||||
|
torrent_hash, torrent.get("name"), exc,
|
||||||
|
)
|
||||||
return result
|
return result
|
||||||
|
|
||||||
def get_resource(self, torrent_hash: str) -> NormalizedResource | None:
|
def get_resource(self, torrent_hash: str) -> NormalizedResource | None:
|
||||||
|
|||||||
@@ -91,12 +91,21 @@ def normalize_resource(
|
|||||||
observed_at: datetime | None = None,
|
observed_at: datetime | None = None,
|
||||||
) -> NormalizedResource:
|
) -> NormalizedResource:
|
||||||
metainfo = decode_metainfo(metainfo_bytes)
|
metainfo = decode_metainfo(metainfo_bytes)
|
||||||
if len(raw_files) != len(metainfo.files):
|
# qBittorrent/libtorrent may omit torrent padding files from
|
||||||
|
# /torrents/files while keeping them in the exported metainfo.
|
||||||
|
visible_metainfo_files = tuple(
|
||||||
|
item for item in metainfo.files if not item.padding
|
||||||
|
)
|
||||||
|
if len(raw_files) == len(metainfo.files):
|
||||||
|
matched_metainfo_files = metainfo.files
|
||||||
|
elif len(raw_files) == len(visible_metainfo_files):
|
||||||
|
matched_metainfo_files = visible_metainfo_files
|
||||||
|
else:
|
||||||
raise ResourceError("qBittorrent and metainfo file counts differ")
|
raise ResourceError("qBittorrent and metainfo file counts differ")
|
||||||
files = []
|
files = []
|
||||||
canonical = True
|
canonical = True
|
||||||
for expected_index, (raw, meta_file) in enumerate(
|
for expected_index, (raw, meta_file) in enumerate(
|
||||||
zip(raw_files, metainfo.files, strict=True)
|
zip(raw_files, matched_metainfo_files, strict=True)
|
||||||
):
|
):
|
||||||
index = _integer(raw.get("index"), "file index")
|
index = _integer(raw.get("index"), "file index")
|
||||||
if index != expected_index:
|
if index != expected_index:
|
||||||
|
|||||||
@@ -7,6 +7,8 @@ import json
|
|||||||
import os
|
import os
|
||||||
import sqlite3
|
import sqlite3
|
||||||
import stat
|
import stat
|
||||||
|
import threading
|
||||||
|
from contextlib import contextmanager
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
@@ -14,6 +16,12 @@ from pathlib import Path
|
|||||||
SCHEMA_VERSION = 3
|
SCHEMA_VERSION = 3
|
||||||
|
|
||||||
|
|
||||||
|
# A client daemon executes unrelated jobs concurrently. SQLite still permits
|
||||||
|
# only one writer, so serialize this process's short state transactions rather
|
||||||
|
# than allowing an otherwise healthy job to fail after its busy timeout.
|
||||||
|
_DATABASE_LOCK = threading.RLock()
|
||||||
|
|
||||||
|
|
||||||
class CommandConflict(RuntimeError):
|
class CommandConflict(RuntimeError):
|
||||||
pass
|
pass
|
||||||
|
|
||||||
@@ -557,14 +565,30 @@ class ClientStore:
|
|||||||
).fetchall()
|
).fetchall()
|
||||||
return [dict(row) for row in rows]
|
return [dict(row) for row in rows]
|
||||||
|
|
||||||
def _connect(self) -> sqlite3.Connection:
|
@contextmanager
|
||||||
connection = sqlite3.connect(self.database, isolation_level=None, timeout=5)
|
def _connect(self):
|
||||||
connection.row_factory = sqlite3.Row
|
"""Yield one connection while serializing local SQLite writers.
|
||||||
connection.execute("PRAGMA foreign_keys = ON")
|
|
||||||
connection.execute("PRAGMA journal_mode = WAL")
|
The longer SQLite timeout also covers a short lock held by a separate
|
||||||
connection.execute("PRAGMA synchronous = FULL")
|
maintenance process such as a backup. The lock is deliberately held
|
||||||
connection.execute("PRAGMA busy_timeout = 5000")
|
for the full transaction, including ``BEGIN IMMEDIATE``.
|
||||||
return connection
|
"""
|
||||||
|
with _DATABASE_LOCK:
|
||||||
|
connection = sqlite3.connect(
|
||||||
|
self.database, isolation_level=None, timeout=30
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
connection.row_factory = sqlite3.Row
|
||||||
|
connection.execute("PRAGMA foreign_keys = ON")
|
||||||
|
connection.execute("PRAGMA journal_mode = WAL")
|
||||||
|
connection.execute("PRAGMA synchronous = FULL")
|
||||||
|
connection.execute("PRAGMA busy_timeout = 30000")
|
||||||
|
# Preserve the original ``with connection`` commit/rollback
|
||||||
|
# behavior used by every store operation.
|
||||||
|
with connection:
|
||||||
|
yield connection
|
||||||
|
finally:
|
||||||
|
connection.close()
|
||||||
|
|
||||||
|
|
||||||
def _canonical(value: object) -> str:
|
def _canonical(value: object) -> str:
|
||||||
|
|||||||
+39
-1
@@ -12,7 +12,7 @@ from archive_clients.jobs import (
|
|||||||
JobExecutionError,
|
JobExecutionError,
|
||||||
_resource_fingerprint,
|
_resource_fingerprint,
|
||||||
)
|
)
|
||||||
from archive_clients.syncthing import RouteSetupError
|
from archive_clients.syncthing import RouteSetupError, SyncthingTransferStatus
|
||||||
from archive_clients.resources import normalize_resource
|
from archive_clients.resources import normalize_resource
|
||||||
from archive_clients.state import ClientStore
|
from archive_clients.state import ClientStore
|
||||||
from archive_control.v1 import control_pb2, job_pb2
|
from archive_control.v1 import control_pb2, job_pb2
|
||||||
@@ -39,6 +39,44 @@ class SlowRescanSyncthing(CompleteSyncthing):
|
|||||||
|
|
||||||
|
|
||||||
class ClientJobHappyPathTests(unittest.TestCase):
|
class ClientJobHappyPathTests(unittest.TestCase):
|
||||||
|
def test_syncthing_api_outage_during_transfer_is_retried(self):
|
||||||
|
"""A transient local REST outage must not terminally fail the job."""
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
store = ClientStore(root / "client.db")
|
||||||
|
store.initialize()
|
||||||
|
definition = job_pb2.JobDefinition(
|
||||||
|
job_id=str(uuid4()),
|
||||||
|
transfer={
|
||||||
|
"source_client_id": "cache-1",
|
||||||
|
"target_client_id": "archive-1",
|
||||||
|
"route_id": "route-1",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
observer = Mock()
|
||||||
|
observer.status.side_effect = [
|
||||||
|
RouteSetupError("Syncthing API is unavailable"),
|
||||||
|
SyncthingTransferStatus(1, 42, 42, True, 0),
|
||||||
|
]
|
||||||
|
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._observer = Mock(return_value=observer)
|
||||||
|
progress = Mock()
|
||||||
|
|
||||||
|
executor._wait_for_syncthing(definition, progress)
|
||||||
|
|
||||||
|
self.assertEqual(observer.status.call_count, 2)
|
||||||
|
progress.assert_called_once_with(1, 42, 42, "0 Syncthing items still needed")
|
||||||
|
|
||||||
def test_reconnect_replay_never_duplicates_any_transfer_step(self):
|
def test_reconnect_replay_never_duplicates_any_transfer_step(self):
|
||||||
for step in (
|
for step in (
|
||||||
job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
|
job_pb2.JOB_STEP_KIND_SOURCE_STAGE,
|
||||||
|
|||||||
@@ -205,6 +205,54 @@ class QBittorrentReaderTests(unittest.TestCase):
|
|||||||
self.assertEqual(resources, [])
|
self.assertEqual(resources, [])
|
||||||
self.assertEqual(len(opener.calls), 2)
|
self.assertEqual(len(opener.calls), 2)
|
||||||
|
|
||||||
|
def test_listing_skips_malformed_torrent_and_keeps_valid_resources(self):
|
||||||
|
def torrent(name, content):
|
||||||
|
info = {
|
||||||
|
b"length": len(content), b"name": name.encode(),
|
||||||
|
b"piece length": 16384, b"pieces": b"x" * 20,
|
||||||
|
}
|
||||||
|
return encode({b"info": info}), hashlib.sha1(encode(info)).hexdigest()
|
||||||
|
|
||||||
|
first_bytes, first_hash = torrent("first.txt", b"one")
|
||||||
|
malformed_bytes, malformed_hash = torrent("broken.txt", b"bad")
|
||||||
|
second_bytes, second_hash = torrent("second.txt", b"two")
|
||||||
|
torrents = [
|
||||||
|
{"hash": first_hash, "name": "first.txt", "state": "uploading"},
|
||||||
|
{"hash": malformed_hash, "name": "broken.txt", "state": "stalledUP"},
|
||||||
|
{"hash": second_hash, "name": "second.txt", "state": "uploading"},
|
||||||
|
]
|
||||||
|
responses = [
|
||||||
|
b"Ok.", json.dumps(torrents).encode(),
|
||||||
|
b'[{"index":0,"name":"first.txt","size":3,"progress":1,"priority":1}]',
|
||||||
|
first_bytes,
|
||||||
|
b'[{"index":0,"name":"broken.txt","size":3,"progress":1,"priority":1},'
|
||||||
|
b'{"index":1,"name":"unexpected.txt","size":1,"progress":1,"priority":1}]',
|
||||||
|
malformed_bytes,
|
||||||
|
b'[{"index":0,"name":"second.txt","size":3,"progress":1,"priority":1}]',
|
||||||
|
second_bytes,
|
||||||
|
]
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
root = Path(directory)
|
||||||
|
password = root / "password"
|
||||||
|
password.write_text("secret", encoding="utf-8")
|
||||||
|
os.chmod(password, 0o600)
|
||||||
|
config = ServiceConfig(
|
||||||
|
"http://qb", PurePosixPath("/downloads"), root,
|
||||||
|
username="admin", password_file=password,
|
||||||
|
)
|
||||||
|
opener = _Opener(responses)
|
||||||
|
with patch(
|
||||||
|
"archive_clients.qbittorrent.request.build_opener",
|
||||||
|
return_value=opener,
|
||||||
|
), self.assertLogs("archive_clients.qbittorrent", "WARNING") as logs:
|
||||||
|
resources = QBittorrentReader(config).list_resources()
|
||||||
|
|
||||||
|
self.assertEqual(
|
||||||
|
[resource.summary.display_name for resource in resources],
|
||||||
|
["first.txt", "second.txt"],
|
||||||
|
)
|
||||||
|
self.assertIn(malformed_hash, "\n".join(logs.output))
|
||||||
|
|
||||||
def test_stopped_add_selection_recheck_and_entry_only_delete(self):
|
def test_stopped_add_selection_recheck_and_entry_only_delete(self):
|
||||||
torrent_hash = "a" * 40
|
torrent_hash = "a" * 40
|
||||||
responses = [
|
responses = [
|
||||||
|
|||||||
@@ -118,6 +118,47 @@ class ResourceTests(unittest.TestCase):
|
|||||||
self.assertEqual(root.available_file_count, 1)
|
self.assertEqual(root.available_file_count, 1)
|
||||||
self.assertEqual(root.available_logical_bytes, 3)
|
self.assertEqual(root.available_logical_bytes, 3)
|
||||||
|
|
||||||
|
def test_qbittorrent_hidden_padding_files_are_normalized(self):
|
||||||
|
info = {
|
||||||
|
b"files": [
|
||||||
|
{b"length": 3, b"path": [b"first.bin"]},
|
||||||
|
{
|
||||||
|
b"length": 5,
|
||||||
|
b"path": [b"_____padding_file_5"],
|
||||||
|
},
|
||||||
|
{b"length": 7, b"path": [b"last.bin"]},
|
||||||
|
],
|
||||||
|
b"name": b"with-padding",
|
||||||
|
b"piece length": 16384,
|
||||||
|
b"pieces": b"x" * 20,
|
||||||
|
}
|
||||||
|
metainfo = encode({b"info": info})
|
||||||
|
torrent_hash = hashlib.sha1(encode(info)).hexdigest()
|
||||||
|
normalized = normalize_resource(
|
||||||
|
{
|
||||||
|
"hash": torrent_hash,
|
||||||
|
"name": "with-padding",
|
||||||
|
"state": "stalledUP",
|
||||||
|
},
|
||||||
|
[
|
||||||
|
{
|
||||||
|
"index": 0, "name": "with-padding/first.bin",
|
||||||
|
"size": 3, "completed": 3, "priority": 1,
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"index": 1, "name": "with-padding/last.bin",
|
||||||
|
"size": 7, "completed": 7, "priority": 1,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
metainfo,
|
||||||
|
)
|
||||||
|
self.assertEqual(
|
||||||
|
[item.canonical_path for item in normalized.files],
|
||||||
|
["with-padding/first.bin", "with-padding/last.bin"],
|
||||||
|
)
|
||||||
|
self.assertEqual(normalized.summary.total_file_count, 2)
|
||||||
|
self.assertEqual(normalized.summary.selected_complete_bytes, 10)
|
||||||
|
|
||||||
def test_noncanonical_and_unsafe_paths_are_distinct(self):
|
def test_noncanonical_and_unsafe_paths_are_distinct(self):
|
||||||
info = {
|
info = {
|
||||||
b"length": 3,
|
b"length": 3,
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import tempfile
|
import tempfile
|
||||||
import unittest
|
import unittest
|
||||||
import os
|
import os
|
||||||
|
import threading
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from uuid import uuid4
|
from uuid import uuid4
|
||||||
|
|
||||||
@@ -13,6 +14,40 @@ from archive_clients.state import (
|
|||||||
|
|
||||||
|
|
||||||
class ClientStoreTests(unittest.TestCase):
|
class ClientStoreTests(unittest.TestCase):
|
||||||
|
def test_database_connections_serialize_local_writers(self):
|
||||||
|
"""Concurrent jobs share one daemon DB without SQLite lock failures."""
|
||||||
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
|
database = Path(directory) / "state.db"
|
||||||
|
first = ClientStore(database)
|
||||||
|
second = ClientStore(database)
|
||||||
|
first.initialize()
|
||||||
|
barrier = threading.Barrier(2)
|
||||||
|
failures: list[Exception] = []
|
||||||
|
|
||||||
|
def write(store, command_id):
|
||||||
|
try:
|
||||||
|
barrier.wait()
|
||||||
|
store.accept_command(command_id, '{"kind":"job"}', '{"ok":true}')
|
||||||
|
except Exception as exc: # pragma: no cover - assertion below
|
||||||
|
failures.append(exc)
|
||||||
|
|
||||||
|
left = threading.Thread(target=write, args=(first, "command-1"))
|
||||||
|
right = threading.Thread(target=write, args=(second, "command-2"))
|
||||||
|
left.start()
|
||||||
|
right.start()
|
||||||
|
left.join(1)
|
||||||
|
right.join(1)
|
||||||
|
|
||||||
|
self.assertFalse(left.is_alive())
|
||||||
|
self.assertFalse(right.is_alive())
|
||||||
|
self.assertEqual(failures, [])
|
||||||
|
self.assertEqual(len(first.list_accepted_commands()), 2)
|
||||||
|
with first._connect() as connection:
|
||||||
|
self.assertEqual(
|
||||||
|
connection.execute("PRAGMA busy_timeout").fetchone()[0],
|
||||||
|
30000,
|
||||||
|
)
|
||||||
|
|
||||||
def test_command_acceptance_is_durable_and_content_addressed(self):
|
def test_command_acceptance_is_durable_and_content_addressed(self):
|
||||||
with tempfile.TemporaryDirectory() as directory:
|
with tempfile.TemporaryDirectory() as directory:
|
||||||
database = Path(directory) / "state.db"
|
database = Path(directory) / "state.db"
|
||||||
|
|||||||
Reference in New Issue
Block a user