mirror of
https://github.com/calibrain/shelfmark.git
synced 2026-10-05 16:31:14 +01:00
Fix: IRC connection threading + fs logging (#445)
This commit is contained in:
@@ -62,7 +62,8 @@ def is_onboarding_complete() -> bool:
|
||||
return True
|
||||
|
||||
return False
|
||||
except (json.JSONDecodeError, OSError):
|
||||
except (json.JSONDecodeError, OSError) as e:
|
||||
logger.warning(f"Could not read onboarding status from settings.json: {e}")
|
||||
return False
|
||||
|
||||
|
||||
@@ -438,8 +439,8 @@ def save_onboarding_settings(values: Dict[str, Any]) -> Dict[str, Any]:
|
||||
try:
|
||||
from shelfmark.core.config import config
|
||||
config.refresh()
|
||||
except ImportError:
|
||||
pass
|
||||
except ImportError as e:
|
||||
logger.debug(f"Could not refresh config after onboarding: {e}")
|
||||
|
||||
return {"success": True, "message": "Onboarding complete!"}
|
||||
|
||||
|
||||
@@ -171,8 +171,9 @@ def atomic_move(source_path: Path, dest_path: Path, max_attempts: int = 100) ->
|
||||
logger.info(f"File collision resolved (fallback): {try_path.name}")
|
||||
return try_path
|
||||
except Exception as fallback_error:
|
||||
# Fallback failed, raise original error
|
||||
raise e
|
||||
# Fallback failed, chain exceptions for better debugging
|
||||
logger.error(f"NFS fallback also failed: {fallback_error}")
|
||||
raise e from fallback_error
|
||||
raise
|
||||
|
||||
raise RuntimeError(f"Could not move file after {max_attempts} attempts: {dest_path}")
|
||||
@@ -249,7 +250,8 @@ def atomic_copy(source_path: Path, dest_path: Path, max_attempts: int = 100) ->
|
||||
try:
|
||||
_perform_nfs_fallback(source_path, temp_path, is_move=False)
|
||||
except Exception as fallback_error:
|
||||
raise e
|
||||
logger.error(f"NFS fallback also failed: {fallback_error}")
|
||||
raise e from fallback_error
|
||||
else:
|
||||
raise
|
||||
|
||||
|
||||
@@ -42,7 +42,8 @@ def should_bypass_proxy(url: str) -> bool:
|
||||
try:
|
||||
parsed = urllib.parse.urlparse(url)
|
||||
hostname = (parsed.hostname or "").lower()
|
||||
except Exception:
|
||||
except Exception as e:
|
||||
logger.debug(f"Failed to parse URL for proxy bypass check: {url} - {e}")
|
||||
return False
|
||||
|
||||
if not hostname:
|
||||
|
||||
@@ -44,6 +44,7 @@ class IRCConnectionManager:
|
||||
self._connections: dict[str, IRCClient] = {}
|
||||
self._last_used: dict[str, float] = {}
|
||||
self._channels: dict[str, str] = {} # connection_key -> joined channel
|
||||
self._connecting: dict[str, bool] = {} # Track keys currently being connected
|
||||
self._conn_lock = threading.Lock()
|
||||
self._cleanup_thread: Optional[threading.Thread] = None
|
||||
self._running = True
|
||||
@@ -112,6 +113,8 @@ class IRCConnectionManager:
|
||||
Connected IRCClient instance that has joined the channel
|
||||
"""
|
||||
key = self._connection_key(server, port, nick)
|
||||
need_new_connection = False
|
||||
dead_client = None
|
||||
|
||||
with self._conn_lock:
|
||||
# Check for existing connection
|
||||
@@ -130,29 +133,56 @@ class IRCConnectionManager:
|
||||
|
||||
return existing
|
||||
|
||||
# Clean up dead connection if it exists
|
||||
if existing:
|
||||
logger.debug(f"Removing dead connection: {key}")
|
||||
self._connections.pop(key, None)
|
||||
self._last_used.pop(key, None)
|
||||
self._channels.pop(key, None)
|
||||
try:
|
||||
existing.disconnect()
|
||||
except Exception:
|
||||
pass
|
||||
# Check if another thread is already connecting
|
||||
if self._connecting.get(key):
|
||||
logger.debug(f"Another thread is connecting to {key}, waiting...")
|
||||
# Release lock and wait, then retry
|
||||
pass # Fall through to retry logic below
|
||||
else:
|
||||
# Clean up dead connection if it exists
|
||||
if existing:
|
||||
logger.debug(f"Removing dead connection: {key}")
|
||||
self._connections.pop(key, None)
|
||||
self._last_used.pop(key, None)
|
||||
self._channels.pop(key, None)
|
||||
dead_client = existing
|
||||
|
||||
# Create new connection
|
||||
# Mark that we're connecting (prevents duplicate attempts)
|
||||
self._connecting[key] = True
|
||||
need_new_connection = True
|
||||
|
||||
# Clean up dead client outside lock
|
||||
if dead_client:
|
||||
try:
|
||||
dead_client.disconnect()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# If another thread is connecting, wait and retry
|
||||
if not need_new_connection:
|
||||
time.sleep(0.5)
|
||||
return self.get_connection(server, port, nick, use_tls, channel)
|
||||
|
||||
# Create new connection OUTSIDE the lock to avoid blocking other threads
|
||||
try:
|
||||
logger.info(f"Creating new IRC connection to {server}:{port}")
|
||||
client = IRCClient(nick, server, port, use_tls=use_tls)
|
||||
client.connect()
|
||||
client.join_channel(channel)
|
||||
|
||||
# Store connection
|
||||
self._connections[key] = client
|
||||
self._last_used[key] = time.time()
|
||||
self._channels[key] = channel
|
||||
# Store connection (re-acquire lock)
|
||||
with self._conn_lock:
|
||||
self._connections[key] = client
|
||||
self._last_used[key] = time.time()
|
||||
self._channels[key] = channel
|
||||
self._connecting.pop(key, None)
|
||||
|
||||
return client
|
||||
except Exception:
|
||||
# Clear connecting flag on failure
|
||||
with self._conn_lock:
|
||||
self._connecting.pop(key, None)
|
||||
raise
|
||||
|
||||
def release_connection(self, client: IRCClient) -> None:
|
||||
"""Mark a connection as available for reuse.
|
||||
|
||||
Reference in New Issue
Block a user