import json import tempfile import time import unittest from pathlib import Path, PurePosixPath from uuid import uuid4 from archive_clients.config import ServiceConfig from archive_clients.state import ClientStore from archive_clients.syncthing import ( RoutePathConflict, RouteSetupTimeout, SyncthingRouteManager, SyncthingTransferObserver, ) from archive_clients.transfer import stage_transfer from archive_control.v1 import route_pb2, transfer_pb2 class FakeTransport: def __init__(self): self.status = {"myID": "LOCAL"} self.config = {"devices": [], "folders": []} self.puts = [] self.posts = [] def get_json(self, path): if path == "/rest/system/status": return self.status if path == "/rest/config": return self.config raise AssertionError(path) def put_json(self, path, payload): self.puts.append((path, payload)) if path.startswith("/rest/config/devices/"): self.config["devices"].append(payload) elif path.startswith("/rest/config/folders/"): self.config["folders"].append(payload) else: raise AssertionError(path) def post(self, path): self.posts.append(path) class FakeTransferTransport: def __init__(self, completion, need): self.completion = completion self.need = need self.posts = [] def get_json(self, path): if path.startswith("/rest/db/completion?"): return self.completion if path.startswith("/rest/db/need?"): return self.need raise AssertionError(path) def put_json(self, path, payload): raise AssertionError((path, payload)) def post(self, path): self.posts.append(path) class SyncthingRouteManagerTests(unittest.TestCase): def setUp(self): self.temp_dir = tempfile.TemporaryDirectory() self.root = Path(self.temp_dir.name) / "sync" self.root.mkdir() self.transport = FakeTransport() self.manager = SyncthingRouteManager( ServiceConfig( "http://syncthing", PurePosixPath("/sync"), self.root, ), sparse_supported=True, transport=self.transport, poll_interval=0, ) self.spec = route_pb2.EnsureRouteSpec( route_id="route-1", peer_client_id="archive-1", peer_syncthing_device_id="PEER", peer_addresses=["tcp://archive:22000"], local_relative_path="routes/route-1", setup_timeout_seconds=1800, ) def tearDown(self): self.temp_dir.cleanup() def test_configure_adds_only_peer_and_pairwise_folder(self): configured = self.manager.configure( self.spec, time.monotonic() + 1 ) self.assertEqual(configured.local_device_id, "LOCAL") self.assertTrue(configured.local_path.is_dir()) self.assertEqual(len(self.transport.puts), 2) folder = self.transport.config["folders"][0] self.assertEqual(folder["path"], "/sync/routes/route-1") self.assertEqual(folder["type"], "sendreceive") self.assertEqual( {item["deviceID"] for item in folder["devices"]}, {"LOCAL", "PEER"}, ) self.assertTrue(configured.local_route.archive_control_created) repeated = self.manager.configure(self.spec, time.monotonic() + 1) self.assertEqual(len(self.transport.puts), 2) self.assertFalse(repeated.local_route.archive_control_created) def test_existing_folder_conflicts_are_never_overwritten(self): self.transport.config["folders"].append( { "id": "route-1", "path": "/somewhere-else", "type": "sendreceive", "devices": [{"deviceID": "LOCAL"}, {"deviceID": "PEER"}], } ) with self.assertRaisesRegex(RoutePathConflict, "different path"): self.manager.configure(self.spec, time.monotonic() + 1) self.assertEqual(self.transport.puts, []) def test_existing_home_relative_folder_path_is_accepted(self): self.transport.config["devices"].append({"deviceID": "PEER"}) self.transport.config["folders"].append( { "id": "route-1", "path": "~/routes/route-1", "type": "sendreceive", "devices": [{"deviceID": "LOCAL"}, {"deviceID": "PEER"}], } ) configured = self.manager.configure( self.spec, time.monotonic() + 1 ) self.assertEqual(self.transport.puts, []) self.assertFalse(configured.local_route.archive_control_created) def test_bidirectional_nonce_and_ack_are_required(self): configured = self.manager.configure(self.spec, time.monotonic() + 1) peer_nonce = configured.local_path / ( ".archive-control-route-nonce.archive-1" ) peer_nonce.write_text( json.dumps({ "route_id": "route-1", "client_id": "archive-1", "nonce": "peer-nonce", }), encoding="utf-8", ) acknowledgement = configured.local_path / ( ".archive-control-route-ack.cache-1.archive-1" ) acknowledgement.write_text( json.dumps({"route_id": "route-1", "nonce": "local-nonce"}), encoding="utf-8", ) verified = self.manager.verify( configured, self.spec, "cache-1", "local-nonce", time.monotonic() + 1, ) self.assertEqual(verified, (True, True)) peer_ack = configured.local_path / ( ".archive-control-route-ack.archive-1.cache-1" ) self.assertEqual(json.loads(peer_ack.read_text())["nonce"], "peer-nonce") local_nonce = configured.local_path / ( ".archive-control-route-nonce.cache-1" ) before = (local_nonce.stat(), peer_ack.stat()) repeated = self.manager.verify( configured, self.spec, "cache-1", "local-nonce", time.monotonic() + 1, ) after = (local_nonce.stat(), peer_ack.stat()) self.assertEqual(repeated, (True, True)) self.assertEqual( [(item.st_ino, item.st_mtime_ns) for item in before], [(item.st_ino, item.st_mtime_ns) for item in after], ) def test_verification_times_out_without_peer_evidence(self): configured = self.manager.configure(self.spec, time.monotonic() + 1) with self.assertRaises(RouteSetupTimeout): self.manager.verify( configured, self.spec, "cache-1", "local-nonce", time.monotonic() + 0.01, ) self.assertTrue(self.transport.posts) def test_job_prefix_rescan_progress_and_completion(self): source = Path(self.temp_dir.name) / "source" source.mkdir() (source / "payload.bin").write_bytes(b"x" * 8192) metainfo = Path(self.temp_dir.name) / "source.torrent" metainfo.write_bytes(b"torrent") store = ClientStore(Path(self.temp_dir.name) / "state.db") store.initialize() manifest = transfer_pb2.TransferManifest( manifest_version=1, job_id=str(uuid4()), source_client_id="cache-1", target_client_id="archive-1", route_id="route-1", ) manifest.resource_id.info_hash_v1_hex = "a" * 40 manifest.created_at.seconds = 1_700_000_000 manifest.files.add( file_index=0, payload_relative_path="payload/payload.bin", target_canonical_path="payload.bin", logical_bytes=8192, ) manifest.artifacts.add( kind=transfer_pb2.ARTIFACT_KIND_TORRENT_FILE, payload_relative_path="metainfo/source.torrent", logical_bytes=metainfo.stat().st_size, ) published = stage_transfer( manifest, source_root=source, sync_root=self.root, store=store, artifact_sources={"metainfo/source.torrent": metainfo}, ) relative = ( f".archive-control/jobs/{manifest.job_id}" ) in_progress_transport = FakeTransferTransport( {"completion": 50}, { "progress": [{"name": f"{relative}/payload/payload.bin"}], "queued": [], "rest": [], }, ) observer = SyncthingTransferObserver( in_progress_transport, "route-1", relative, published.job_directory, ) observer.rescan() in_progress = observer.status() self.assertEqual(in_progress.fraction_complete, 0.5) self.assertFalse(in_progress.complete) self.assertEqual(in_progress.needed_items, 1) self.assertIn("folder=route-1", in_progress_transport.posts[0]) self.assertIn("sub=.archive-control%2Fjobs%2F", in_progress_transport.posts[0]) complete = SyncthingTransferObserver( FakeTransferTransport( {"completion": 100}, {"progress": [], "queued": [], "rest": []}, ), "route-1", relative, published.job_directory, ).status() self.assertTrue(complete.complete) self.assertEqual(complete.bytes_complete, complete.bytes_total) if __name__ == "__main__": unittest.main()