Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions AzuriteConfig
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
{"instaceID":"7d61017f-125f-488f-b15d-143f1a2fc570"}
1 change: 1 addition & 0 deletions __azurite_db_table__.json
Original file line numberDiff line numberDiff line change
@@ -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"}
12 changes: 9 additions & 3 deletions ccbt/discovery/xet_cas.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -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.

Expand All@@ -63,13 +64,15 @@ 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
self.tracker = tracker_client
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
Expand DownExpand Up@@ -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")

Expand Down
14 changes: 7 additions & 7 deletions ccbt/i18n/__init__.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -6,7 +6,6 @@
from __future__ import annotations

import gettext
import locale
import logging
import os
from pathlib import Path
Expand DownExpand Up@@ -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


Expand Down
6 changes: 5 additions & 1 deletion ccbt/session/session.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -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)
Expand DownExpand Up@@ -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)
Expand Down
5 changes: 2 additions & 3 deletions ccbt/storage/folder_watcher.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -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:
Expand Down
4 changes: 3 additions & 1 deletion ccbt/utils/events.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -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
Expand DownExpand Up@@ -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(
Expand Down
58 changes: 30 additions & 28 deletions tests/conftest.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -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
Expand DownExpand Up@@ -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:
Expand Down
95 changes: 51 additions & 44 deletions tests/integration/test_xet_integration.py
Original file line numberDiff line numberDiff line change
Expand Up@@ -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"
Expand DownExpand Up@@ -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)

Expand DownExpand Up@@ -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()

Expand All@@ -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)

Expand DownExpand Up@@ -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)

Expand DownExpand Up@@ -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)
Expand Down
Loading