diff --git a/AzuriteConfig b/AzuriteConfig new file mode 100644 index 0000000..296fcc2 --- /dev/null +++ b/AzuriteConfig @@ -0,0 +1 @@ +{"instaceID":"7d61017f-125f-488f-b15d-143f1a2fc570"} \ No newline at end of file diff --git a/__azurite_db_table__.json b/__azurite_db_table__.json new file mode 100644 index 0000000..d0a1963 --- /dev/null +++ b/__azurite_db_table__.json @@ -0,0 +1 @@ +{"filename":"c:\\Users\\MeMyself\\bittorrentclient\\__azurite_db_table__.json","collections":[{"name":"$TABLES_COLLECTION$","data":[],"idIndex":null,"binaryIndices":{"account":{"name":"account","dirty":false,"values":[]},"table":{"name":"table","dirty":false,"values":[]}},"constraints":null,"uniqueNames":[],"transforms":{},"objType":"$TABLES_COLLECTION$","dirty":false,"cachedIndex":null,"cachedBinaryIndex":null,"cachedData":null,"adaptiveBinaryIndices":true,"transactional":false,"cloneObjects":false,"cloneMethod":"parse-stringify","asyncListeners":false,"disableMeta":false,"disableChangesApi":true,"disableDeltaChangesApi":true,"autoupdate":false,"serializableIndices":true,"disableFreeze":true,"ttl":null,"maxId":0,"DynamicViews":[],"events":{"insert":[],"update":[],"pre-insert":[],"pre-update":[],"close":[],"flushbuffer":[],"error":[],"delete":[null],"warning":[null]},"changes":[],"dirtyIds":[]},{"name":"$SERVICES_COLLECTION$","data":[],"idIndex":null,"binaryIndices":{},"constraints":null,"uniqueNames":["accountName"],"transforms":{},"objType":"$SERVICES_COLLECTION$","dirty":false,"cachedIndex":null,"cachedBinaryIndex":null,"cachedData":null,"adaptiveBinaryIndices":true,"transactional":false,"cloneObjects":false,"cloneMethod":"parse-stringify","asyncListeners":false,"disableMeta":false,"disableChangesApi":true,"disableDeltaChangesApi":true,"autoupdate":false,"serializableIndices":true,"disableFreeze":true,"ttl":null,"maxId":0,"DynamicViews":[],"events":{"insert":[],"update":[],"pre-insert":[],"pre-update":[],"close":[],"flushbuffer":[],"error":[],"delete":[null],"warning":[null]},"changes":[],"dirtyIds":[]}],"databaseVersion":1.5,"engineVersion":1.5,"autosave":true,"autosaveInterval":5000,"autosaveHandle":null,"throttledSaves":true,"options":{"persistenceMethod":"fs","autosave":true,"autosaveInterval":5000,"serializationMethod":"normal","destructureDelimiter":"$<\n"},"persistenceMethod":"fs","persistenceAdapter":null,"verbose":false,"events":{"init":[null],"loaded":[],"flushChanges":[],"close":[],"changes":[],"warning":[]},"ENV":"NODEJS"} \ No newline at end of file diff --git a/ccbt/discovery/xet_cas.py b/ccbt/discovery/xet_cas.py index 7c16fef..559b25b 100644 --- a/ccbt/discovery/xet_cas.py +++ b/ccbt/discovery/xet_cas.py @@ -54,6 +54,7 @@ def __init__( key_manager: Any = None, # Ed25519KeyManager bloom_filter: Optional[Any] = None, # XetChunkBloomFilter catalog: Optional[Any] = None, # XetChunkCatalog + extension_manager: Optional[Any] = None, # ExtensionManager ): """Initialize P2P CAS with DHT and tracker clients. @@ -63,6 +64,7 @@ def __init__( key_manager: Optional Ed25519KeyManager for signing chunks bloom_filter: Optional bloom filter for chunk availability catalog: Optional chunk catalog for bulk queries + extension_manager: Optional ExtensionManager (preferred over get_extension_manager()) """ self.dht = dht_client @@ -70,6 +72,7 @@ def __init__( self.key_manager = key_manager self.bloom_filter = bloom_filter self.catalog = catalog + self.extension_manager = extension_manager self.pex_manager: Optional[Any] = None self.peer_authorizer: Optional[Any] = None self.discovery_backend_success_notifier: Optional[Any] = None @@ -592,10 +595,13 @@ async def download_chunk( msg = f"Chunk hash must be 32 bytes, got {len(chunk_hash)}" raise ValueError(msg) - # Get extension manager and Xet extension - from ccbt.extensions.manager import get_extension_manager + # Get extension manager and Xet extension (prefer injected manager) + if self.extension_manager is None: + from ccbt.extensions.manager import get_extension_manager - extension_manager = get_extension_manager() + extension_manager = get_extension_manager() + else: + extension_manager = self.extension_manager extension_protocol = extension_manager.get_extension("protocol") xet_ext = extension_manager.get_extension("xet") diff --git a/ccbt/i18n/__init__.py b/ccbt/i18n/__init__.py index f8df228..6e9fda6 100644 --- a/ccbt/i18n/__init__.py +++ b/ccbt/i18n/__init__.py @@ -6,7 +6,6 @@ from __future__ import annotations import gettext -import locale import logging import os from pathlib import Path @@ -75,16 +74,17 @@ def get_locale() -> str: locale_code, ) - # Fall back to system locale - try: - system_locale, _ = locale.getdefaultlocale() + # Fall back to system locale (same env order as getdefaultlocale; avoid deprecated getdefaultlocale()) + for env_var in ("LC_ALL", "LC_CTYPE", "LANG", "LANGUAGE"): + raw = os.environ.get(env_var) + if not raw: + continue + # Take first language if LANGUAGE is a colon-separated list + system_locale = raw.split(":")[0].split(".")[0].strip() if system_locale: locale_code = system_locale.split("_")[0].lower() if _is_valid_locale(locale_code): return locale_code - except Exception: - pass - return DEFAULT_LOCALE diff --git a/ccbt/session/session.py b/ccbt/session/session.py index 9cd8fbc..94ec41e 100644 --- a/ccbt/session/session.py +++ b/ccbt/session/session.py @@ -3937,6 +3937,7 @@ def _ensure_xet_discovery_graph(self) -> None: key_manager=self.key_manager, bloom_filter=self.xet_bloom_filter, catalog=self.xet_catalog, + extension_manager=self.extension_manager, ) if hasattr(self, "pex_manager") and self.pex_manager is not None: self.xet_cas_client.register_pex_manager(self.pex_manager) @@ -4144,8 +4145,11 @@ async def start_dht_client() -> None: try: from ccbt.extensions.manager import get_extension_manager - self._ensure_xet_discovery_graph() + # Set extension_manager before _ensure_xet_discovery_graph so P2PCASClient + # receives it via injection (avoids deprecated get_extension_manager() in + # download_chunk) and uses the same lifecycle-bound instance. self.extension_manager = get_extension_manager() + self._ensure_xet_discovery_graph() xet_ext = self.extension_manager.extensions.get("xet") if xet_ext is not None: metadata_exchange = XetMetadataExchange(xet_ext) diff --git a/ccbt/storage/folder_watcher.py b/ccbt/storage/folder_watcher.py index 3d09382..2fa55c0 100644 --- a/ccbt/storage/folder_watcher.py +++ b/ccbt/storage/folder_watcher.py @@ -234,9 +234,8 @@ async def _emit_folder_changed() -> None: self._loop.call_soon_threadsafe( lambda: asyncio.create_task(_emit_folder_changed()) ) - else: - with contextlib.suppress(RuntimeError): - asyncio.create_task(_emit_folder_changed()) # noqa: RUF006 + # When loop is not running (e.g. during teardown), skip emit to avoid + # "coroutine was never awaited" and to avoid scheduling on a dead loop. # Call all callbacks for callback in self.change_callbacks: diff --git a/ccbt/utils/events.py b/ccbt/utils/events.py index aee57dd..f551fe2 100644 --- a/ccbt/utils/events.py +++ b/ccbt/utils/events.py @@ -485,6 +485,7 @@ def __init__( "dht_node_added": 0.1, "monitoring_heartbeat": 1.0, "global_metrics_update": 0.5, + "folder_changed": 0.5, } # Statistics @@ -854,12 +855,13 @@ def get_event_bus() -> EventBus: config = get_config() obs_config = config.observability - # Build throttle intervals from config + # Build throttle intervals from config (folder_changed not in config, use default) throttle_intervals = { "dht_node_found": obs_config.event_bus_throttle_dht_node_found, "dht_node_added": obs_config.event_bus_throttle_dht_node_added, "monitoring_heartbeat": obs_config.event_bus_throttle_monitoring_heartbeat, "global_metrics_update": obs_config.event_bus_throttle_global_metrics_update, + "folder_changed": 0.5, } _event_bus = EventBus( diff --git a/tests/conftest.py b/tests/conftest.py index 5f7ba8b..7b2b5fb 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -660,19 +660,18 @@ def cleanup_network_ports(): This fixture provides best-effort cleanup by waiting for ports to be released. Actual port cleanup happens in component stop() methods. - CRITICAL FIX: Increased wait time from 0.1s to 2.0s to ensure ports are released - before next test starts. This prevents "Address already in use" errors. + Wait time is 0.2s by default so integration runs don't appear to hang (2s per + test was causing 762 tests to take 25+ minutes in teardown alone). Set + CCBT_TEST_PORT_RELEASE_DELAY=2.0 in CI if "Address already in use" appears. Also releases ports from port pool manager to prevent pool exhaustion. """ yield import time - # CRITICAL FIX: Increased from 0.1s to 2.0s to ensure ports are fully released - # Ports can take time to be released by the OS, especially on CI/CD systems - # Note: Actual port cleanup happens in component stop() methods - # This fixture just ensures we wait for cleanup to complete - time.sleep(2.0) + delay = float(os.environ.get("CCBT_TEST_PORT_RELEASE_DELAY", "0.2")) + if delay > 0: + time.sleep(delay) # Release all ports from port pool after each test # This ensures the pool doesn't get exhausted over many tests @@ -771,27 +770,30 @@ def verify_test_isolation(): except Exception: pass # Thread enumeration may fail, ignore - # Check for open file handles (if psutil available) - try: - import psutil - import os as os_module - process = psutil.Process(os_module.getpid()) - open_files = process.open_files() - # Filter out known system files and pytest files - suspicious_files = [ - f.path for f in open_files - if not any( - skip in f.path.lower() - for skip in ["/dev/", "/proc/", "pytest", ".pyc", "__pycache__", ".cursor"] - ) - ] - if suspicious_files and len(suspicious_files) > 5: # Allow some files, warn on many - warnings.append(f"Many open file handles detected: {len(suspicious_files)} files") - except ImportError: - # psutil not available, skip file handle check - pass - except Exception: - pass # File handle check may fail, ignore + # Check for open file handles (if psutil available). + # Skip on Windows: psutil.Process.open_files() is very slow there (~2s per call) + # and causes integration runs to appear to hang during teardown. + if sys.platform != "win32": + try: + import psutil + import os as os_module + process = psutil.Process(os_module.getpid()) + open_files = process.open_files() + # Filter out known system files and pytest files + suspicious_files = [ + f.path for f in open_files + if not any( + skip in f.path.lower() + for skip in ["/dev/", "/proc/", "pytest", ".pyc", "__pycache__", ".cursor"] + ) + ] + if suspicious_files and len(suspicious_files) > 5: # Allow some files, warn on many + warnings.append(f"Many open file handles detected: {len(suspicious_files)} files") + except ImportError: + # psutil not available, skip file handle check + pass + except Exception: + pass # File handle check may fail, ignore # Log warnings if any found if warnings: diff --git a/tests/integration/test_xet_integration.py b/tests/integration/test_xet_integration.py index 11db7ee..65190f3 100644 --- a/tests/integration/test_xet_integration.py +++ b/tests/integration/test_xet_integration.py @@ -288,20 +288,21 @@ async def test_xet_extension_message_handling(self): async def test_xet_cas_download_chunk_with_existing_connection(self): """Test downloading chunk with existing connection from connection manager.""" from ccbt.discovery.xet_cas import P2PCASClient - from ccbt.extensions.manager import get_extension_manager + from ccbt.extensions.manager import ExtensionManager from ccbt.extensions.xet import XetExtension from ccbt.models import PeerInfo from ccbt.peer.async_peer_connection import AsyncPeerConnection, ConnectionState from ccbt.storage.xet_hashing import XetHasher - # Create CAS client - cas_client = P2PCASClient(dht_client=None, tracker_client=None) - - # Xet extension is already registered - extension_manager = get_extension_manager() + extension_manager = ExtensionManager() xet_ext = extension_manager.get_extension("xet") assert xet_ext is not None + # Create CAS client with extension manager (avoids deprecated get_extension_manager()) + cas_client = P2PCASClient( + dht_client=None, tracker_client=None, extension_manager=extension_manager + ) + # Create mock connection chunk_hash = b"DOWNLOAD" * 4 # 32 bytes chunk_data = b"Test chunk data for download" @@ -372,43 +373,41 @@ async def test_xet_cas_download_chunk_extension_not_available(self): from ccbt.discovery.xet_cas import P2PCASClient from ccbt.models import PeerInfo - cas_client = P2PCASClient(dht_client=None, tracker_client=None) + mock_manager = MagicMock() + mock_manager.get_extension = MagicMock( + side_effect=lambda name: None if name == "protocol" else MagicMock() + ) + cas_client = P2PCASClient( + dht_client=None, + tracker_client=None, + extension_manager=mock_manager, + ) chunk_hash = b"EXTENSION" * 4 # 32 bytes peer = PeerInfo(ip="192.168.1.1", port=6881) - # Mock extension manager to return None for extension protocol - with patch("ccbt.extensions.manager.get_extension_manager") as mock_get_manager: - mock_manager = MagicMock() - # Make get_extension return None for "protocol" - def mock_get_ext(name): - if name == "protocol": - return None - return MagicMock() - mock_manager.get_extension = MagicMock(side_effect=mock_get_ext) - mock_get_manager.return_value = mock_manager - - with pytest.raises((NotImplementedError, ValueError)): - await cas_client.download_chunk( - chunk_hash, peer, {"info_hash": b"X" * 20}, None - ) + with pytest.raises((NotImplementedError, ValueError)): + await cas_client.download_chunk( + chunk_hash, peer, {"info_hash": b"X" * 20}, None + ) @pytest.mark.asyncio @pytest.mark.slow async def test_xet_cas_download_chunk_peer_not_support_xet(self): """Test downloading chunk when peer doesn't support Xet extension.""" from ccbt.discovery.xet_cas import P2PCASClient - from ccbt.extensions.manager import get_extension_manager + from ccbt.extensions.manager import ExtensionManager from ccbt.extensions.xet import XetExtension from ccbt.models import PeerInfo from ccbt.peer.async_peer_connection import AsyncPeerConnection, ConnectionState - cas_client = P2PCASClient(dht_client=None, tracker_client=None) - - # Xet extension is already registered - extension_manager = get_extension_manager() + extension_manager = ExtensionManager() xet_ext = extension_manager.get_extension("xet") assert xet_ext is not None + cas_client = P2PCASClient( + dht_client=None, tracker_client=None, extension_manager=extension_manager + ) + chunk_hash = b"0" * 32 # 32 bytes peer = PeerInfo(ip="192.168.1.1", port=6881) @@ -510,9 +509,14 @@ async def test_xet_cas_establish_peer_connection_handshake_error(self): peer = PeerInfo(ip="127.0.0.1", port=6881) torrent_data = {"info_hash": b"X" * 20} - # Mock asyncio.open_connection + # Mock asyncio.open_connection. StreamWriter.write() is sync (buffers only); + # drain() is async. Use MagicMock base so .write is not AsyncMock (avoids + # "coroutine was never awaited" when production code correctly does not + # await writer.write()). mock_reader = AsyncMock() - mock_writer = AsyncMock() + mock_writer = MagicMock() + mock_writer.write = MagicMock() + mock_writer.drain = AsyncMock() mock_writer.close = MagicMock() mock_writer.wait_closed = AsyncMock() @@ -533,18 +537,19 @@ async def test_xet_cas_establish_peer_connection_handshake_error(self): async def test_xet_cas_download_chunk_chunk_not_found(self): """Test downloading chunk when peer responds with CHUNK_NOT_FOUND.""" from ccbt.discovery.xet_cas import P2PCASClient - from ccbt.extensions.manager import get_extension_manager + from ccbt.extensions.manager import ExtensionManager from ccbt.extensions.xet import XetExtension, XetMessageType from ccbt.models import PeerInfo from ccbt.peer.async_peer_connection import AsyncPeerConnection, ConnectionState - cas_client = P2PCASClient(dht_client=None, tracker_client=None) - - # Xet extension is already registered - extension_manager = get_extension_manager() + extension_manager = ExtensionManager() xet_ext = extension_manager.get_extension("xet") assert xet_ext is not None + cas_client = P2PCASClient( + dht_client=None, tracker_client=None, extension_manager=extension_manager + ) + chunk_hash = b"NOTFOUND" * 4 # 32 bytes peer = PeerInfo(ip="192.168.1.1", port=6881) @@ -595,18 +600,19 @@ async def mock_receive_ext_msg(conn, timeout): async def test_xet_cas_download_chunk_chunk_error(self): """Test downloading chunk when peer responds with CHUNK_ERROR.""" from ccbt.discovery.xet_cas import P2PCASClient - from ccbt.extensions.manager import get_extension_manager + from ccbt.extensions.manager import ExtensionManager from ccbt.extensions.xet import XetExtension, XetMessageType from ccbt.models import PeerInfo from ccbt.peer.async_peer_connection import AsyncPeerConnection, ConnectionState - cas_client = P2PCASClient(dht_client=None, tracker_client=None) - - # Xet extension is already registered - extension_manager = get_extension_manager() + extension_manager = ExtensionManager() xet_ext = extension_manager.get_extension("xet") assert xet_ext is not None + cas_client = P2PCASClient( + dht_client=None, tracker_client=None, extension_manager=extension_manager + ) + chunk_hash = b"CHUNKERR" * 4 # 32 bytes peer = PeerInfo(ip="192.168.1.1", port=6881) @@ -655,18 +661,19 @@ async def mock_receive_ext_msg(conn, timeout): async def test_xet_cas_download_chunk_hash_mismatch(self): """Test downloading chunk when hash doesn't match.""" from ccbt.discovery.xet_cas import P2PCASClient - from ccbt.extensions.manager import get_extension_manager + from ccbt.extensions.manager import ExtensionManager from ccbt.extensions.xet import XetExtension from ccbt.models import PeerInfo from ccbt.peer.async_peer_connection import AsyncPeerConnection, ConnectionState - cas_client = P2PCASClient(dht_client=None, tracker_client=None) - - # Xet extension is already registered - extension_manager = get_extension_manager() + extension_manager = ExtensionManager() xet_ext = extension_manager.get_extension("xet") assert xet_ext is not None + cas_client = P2PCASClient( + dht_client=None, tracker_client=None, extension_manager=extension_manager + ) + chunk_hash = b"HASHMISMATCH" * 2 + b"XXXXXXXX" # 32 bytes (12*2=24, +8=32) wrong_chunk_data = b"Wrong chunk data" peer = PeerInfo(ip="192.168.1.1", port=6881) diff --git a/tests/integration/test_xet_sync_workflow.py b/tests/integration/test_xet_sync_workflow.py index fea9a0b..63feaa9 100644 --- a/tests/integration/test_xet_sync_workflow.py +++ b/tests/integration/test_xet_sync_workflow.py @@ -93,29 +93,49 @@ async def test_sync_mode_changes(self, folder_path): @pytest.mark.slow async def test_folder_change_detection(self, folder_path): """Test folder change detection and queuing.""" + from ccbt.utils.events import get_event_bus from ccbt.storage.xet_folder_manager import XetFolder - folder = XetFolder( - folder_path=folder_path, - sync_mode="best_effort", - check_interval=0.5, - session_manager=_build_session_manager_stub(), - ) + # Start event bus so FOLDER_CHANGED events are consumed; otherwise the queue + # fills (watchdog can emit many events on Windows), emit blocks time out + # repeatedly, and the test can hit the global timeout. + event_bus = get_event_bus() + bus_was_running = event_bus.running + if not bus_was_running: + await event_bus.start() + try: + folder = XetFolder( + folder_path=folder_path, + sync_mode="best_effort", + check_interval=0.5, + session_manager=_build_session_manager_stub(), + ) - await folder.start() + await folder.start() - # Create new file - new_file = folder_path / "new_file.txt" - new_file.write_text("new content") + # Create new file + new_file = folder_path / "new_file.txt" + new_file.write_text("new content") - # Wait for change detection - await asyncio.sleep(1.0) + # Wait for change detection (realtime sync runs every check_interval=0.5s). + await asyncio.sleep(1.5) - # Trigger sync - synced = await folder.sync() - assert synced is True + # Trigger sync; retry if realtime sync has already started one. + success = False + for _ in range(6): + ok, _ = await folder.sync() + if ok: + success = True + break + await asyncio.sleep(0.5) + assert success, "folder.sync() did not succeed within retries" - await folder.stop() + await folder.stop() + # Give Windows time to release db file handle before temp_dir teardown. + await asyncio.sleep(0.25) + finally: + if not bus_was_running and event_bus.running: + await event_bus.stop() @pytest.mark.asyncio @pytest.mark.slow