#!/usr/bin/env python3 """Proxy TCP pour PAC Arkteos REG3.""" import asyncio import logging import sys KEEPALIVE_INTERVAL = 300 RECONNECT_DELAY = 10 logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)s %(message)s", stream=sys.stdout, ) logger = logging.getLogger(__name__) def parse_allow_client_writes(value: str | None) -> bool: """Retourne False tant que l'option n'est pas explicitement vraie.""" return value is not None and value.strip().lower() == "true" class ArkteosProxy: def __init__( self, pac_host: str, pac_port: int, proxy_port: int, allow_client_writes: bool = False, ) -> None: self.pac_host = pac_host self.pac_port = pac_port self.proxy_port = proxy_port self.allow_client_writes = allow_client_writes self.clients: set[asyncio.StreamWriter] = set() self.pac_writer: asyncio.StreamWriter | None = None self.pac_write_lock = asyncio.Lock() self.stop_event = asyncio.Event() async def write_to_pac(self, data: bytes) -> None: """Écrit un bloc complet vers la PAC sans l'altérer.""" async with self.pac_write_lock: if self.pac_writer is None: raise ConnectionError("PAC non connectée") self.pac_writer.write(data) await self.pac_writer.drain() async def send_keepalive(self) -> None: await self.write_to_pac(b"\x00") logger.info("Keepalive envoyé à la PAC") async def pac_reader(self, reader: asyncio.StreamReader) -> None: logger.info("Démarrage lecture PAC") try: while not self.stop_event.is_set(): data = await reader.read(4096) if not data: logger.info("PAC a fermé la connexion") return await self.broadcast_to_clients(data) except (ConnectionError, OSError) as error: logger.info("Erreur lecture PAC : %s", error) finally: logger.info("Arrêt lecture PAC") async def broadcast_to_clients(self, data: bytes) -> None: dead_clients: list[asyncio.StreamWriter] = [] for writer in tuple(self.clients): try: writer.write(data) await writer.drain() except (ConnectionError, OSError): dead_clients.append(writer) for writer in dead_clients: await self.close_client(writer) async def pac_keepalive(self) -> None: logger.info("Démarrage keepalive PAC") try: while not self.stop_event.is_set(): await asyncio.sleep(KEEPALIVE_INTERVAL) await self.send_keepalive() except (ConnectionError, OSError) as error: logger.info("Erreur keepalive PAC : %s", error) finally: logger.info("Arrêt keepalive PAC") async def handle_client( self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter, ) -> None: peername = writer.get_extra_info("peername") logger.info("Nouveau client : %s", peername) self.clients.add(writer) blocked_write_logged = False try: while not self.stop_event.is_set(): data = await reader.read(1024) if not data: return if not self.allow_client_writes: if not blocked_write_logged: logger.info("Tentative d’écriture client ignorée") blocked_write_logged = True continue try: await self.write_to_pac(data) except (ConnectionError, OSError) as error: logger.info("Erreur envoi PAC depuis %s : %s", peername, error) return except (ConnectionError, OSError) as error: logger.info("Erreur client %s : %s", peername, error) finally: await self.close_client(writer) async def close_client(self, writer: asyncio.StreamWriter) -> None: if writer not in self.clients: return self.clients.discard(writer) writer.close() try: await writer.wait_closed() except (ConnectionError, OSError): pass logger.info("Client déconnecté nettoyé") async def close_clients(self) -> None: for writer in tuple(self.clients): await self.close_client(writer) async def run_pac_connection(self) -> None: reader, writer = await asyncio.open_connection(self.pac_host, self.pac_port) self.pac_writer = writer logger.info("Connecté à la PAC %s:%s", self.pac_host, self.pac_port) reader_task = asyncio.create_task(self.pac_reader(reader)) keepalive_task = asyncio.create_task(self.pac_keepalive()) try: await reader_task finally: keepalive_task.cancel() await asyncio.gather(keepalive_task, return_exceptions=True) if self.pac_writer is writer: self.pac_writer = None writer.close() try: await writer.wait_closed() except (ConnectionError, OSError): pass await self.close_clients() async def serve(self) -> None: mode = ( "Mode bidirectionnel : écritures des clients autorisées" if self.allow_client_writes else "Mode lecture seule : écritures des clients bloquées" ) logger.info(mode) server = await asyncio.start_server(self.handle_client, "0.0.0.0", self.proxy_port) logger.info("Proxy en écoute sur 0.0.0.0:%s", self.proxy_port) try: while not self.stop_event.is_set(): try: await self.run_pac_connection() except (ConnectionError, OSError) as error: logger.info("Échec connexion PAC : %s", error) if not self.stop_event.is_set(): logger.info("Nouvelle tentative de connexion PAC dans 10s...") await asyncio.sleep(RECONNECT_DELAY) finally: self.stop_event.set() server.close() await server.wait_closed() await self.close_clients() logger.info("Proxy arrêté") def main() -> None: pac_host = sys.argv[1] if len(sys.argv) > 1 else "192.168.X.X" pac_port = int(sys.argv[2]) if len(sys.argv) > 2 else 9641 proxy_port = int(sys.argv[3]) if len(sys.argv) > 3 else 9641 allow_client_writes = parse_allow_client_writes(sys.argv[4] if len(sys.argv) > 4 else None) proxy = ArkteosProxy(pac_host, pac_port, proxy_port, allow_client_writes) try: asyncio.run(proxy.serve()) except KeyboardInterrupt: logger.info("Arrêt demandé") if __name__ == "__main__": main()