release: prepare arkteos proxy addon 1.0.0
This commit is contained in:
+165
-152
@@ -1,180 +1,193 @@
|
||||
#!/usr/bin/env python3
|
||||
# Proxy TCP pour PAC Arkteos REG3
|
||||
# Permet les connexions simultanées de Node-RED et de l'app mobile Arkteos
|
||||
# Usage : arkteos_proxy.py <pac_host> <pac_port> <proxy_port>
|
||||
"""Proxy TCP pour PAC Arkteos REG3."""
|
||||
|
||||
import socket
|
||||
import threading
|
||||
import sys
|
||||
import time
|
||||
import asyncio
|
||||
import logging
|
||||
import sys
|
||||
|
||||
# Configuration depuis les arguments (injectés par run.sh depuis les options HA)
|
||||
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
|
||||
|
||||
KEEPALIVE_INTERVAL = 300
|
||||
RECONNECT_DELAY = 10
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
format='%(asctime)s %(levelname)s %(message)s',
|
||||
stream=sys.stdout
|
||||
format="%(asctime)s %(levelname)s %(message)s",
|
||||
stream=sys.stdout,
|
||||
)
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
stop_event = threading.Event()
|
||||
clients = []
|
||||
clients_lock = threading.Lock()
|
||||
|
||||
def log(msg):
|
||||
logger.info(msg)
|
||||
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"
|
||||
|
||||
def connect_to_pac():
|
||||
while not stop_event.is_set():
|
||||
|
||||
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:
|
||||
s = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
s.connect((PAC_HOST, PAC_PORT))
|
||||
log(f"Connecté à la PAC {PAC_HOST}:{PAC_PORT}")
|
||||
return s
|
||||
except Exception as e:
|
||||
log(f"Échec connexion PAC : {e}. Nouvelle tentative dans 10s...")
|
||||
time.sleep(10)
|
||||
return None
|
||||
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")
|
||||
|
||||
def pac_reader(pac_socket):
|
||||
log("Démarrage thread lecture PAC")
|
||||
try:
|
||||
while not stop_event.is_set():
|
||||
async def broadcast_to_clients(self, data: bytes) -> None:
|
||||
dead_clients: list[asyncio.StreamWriter] = []
|
||||
for writer in tuple(self.clients):
|
||||
try:
|
||||
data = pac_socket.recv(4096)
|
||||
except Exception as e:
|
||||
log(f"Erreur lecture PAC : {e}")
|
||||
break
|
||||
if not data:
|
||||
log("PAC a fermé la connexion")
|
||||
break
|
||||
writer.write(data)
|
||||
await writer.drain()
|
||||
except (ConnectionError, OSError):
|
||||
dead_clients.append(writer)
|
||||
for writer in dead_clients:
|
||||
await self.close_client(writer)
|
||||
|
||||
dead_clients = []
|
||||
with clients_lock:
|
||||
for c in clients:
|
||||
try:
|
||||
c.sendall(data)
|
||||
except Exception:
|
||||
dead_clients.append(c)
|
||||
|
||||
if dead_clients:
|
||||
with clients_lock:
|
||||
for c in dead_clients:
|
||||
try:
|
||||
c.close()
|
||||
except Exception:
|
||||
pass
|
||||
if c in clients:
|
||||
clients.remove(c)
|
||||
log("Client déconnecté nettoyé")
|
||||
finally:
|
||||
log("Arrêt thread lecture PAC")
|
||||
pac_socket.close()
|
||||
|
||||
def pac_keepalive(pac_socket):
|
||||
log("Démarrage thread keepalive PAC")
|
||||
while not stop_event.is_set():
|
||||
time.sleep(300)
|
||||
async def pac_keepalive(self) -> None:
|
||||
logger.info("Démarrage keepalive PAC")
|
||||
try:
|
||||
pac_socket.sendall(b'\x00')
|
||||
log("Keepalive envoyé à la PAC")
|
||||
except Exception as e:
|
||||
log(f"Erreur keepalive PAC : {e}")
|
||||
break
|
||||
log("Arrêt thread keepalive PAC")
|
||||
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")
|
||||
|
||||
def handle_client(client_socket, client_addr, pac_socket):
|
||||
log(f"Nouveau client : {client_addr}")
|
||||
with clients_lock:
|
||||
clients.append(client_socket)
|
||||
try:
|
||||
while not stop_event.is_set():
|
||||
try:
|
||||
data = client_socket.recv(1024)
|
||||
except Exception:
|
||||
break
|
||||
if not data:
|
||||
break
|
||||
try:
|
||||
pac_socket.sendall(data)
|
||||
except Exception as e:
|
||||
log(f"Erreur envoi PAC depuis {client_addr} : {e}")
|
||||
break
|
||||
finally:
|
||||
with clients_lock:
|
||||
if client_socket in clients:
|
||||
clients.remove(client_socket)
|
||||
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:
|
||||
client_socket.close()
|
||||
except Exception:
|
||||
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
|
||||
log(f"Client déconnecté : {client_addr}")
|
||||
logger.info("Client déconnecté nettoyé")
|
||||
|
||||
def start_proxy():
|
||||
server_socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
server_socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
server_socket.bind(('0.0.0.0', PROXY_PORT))
|
||||
server_socket.listen(5)
|
||||
log(f"Proxy en écoute sur 0.0.0.0:{PROXY_PORT}")
|
||||
async def close_clients(self) -> None:
|
||||
for writer in tuple(self.clients):
|
||||
await self.close_client(writer)
|
||||
|
||||
pac_socket = connect_to_pac()
|
||||
if pac_socket is None:
|
||||
log("Impossible de se connecter à la PAC. Arrêt.")
|
||||
return
|
||||
|
||||
reader_thread = threading.Thread(target=pac_reader, args=(pac_socket,), daemon=True)
|
||||
reader_thread.start()
|
||||
|
||||
keepalive_thread = threading.Thread(target=pac_keepalive, args=(pac_socket,), daemon=True)
|
||||
keepalive_thread.start()
|
||||
|
||||
try:
|
||||
while not stop_event.is_set():
|
||||
server_socket.settimeout(1)
|
||||
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:
|
||||
client_socket, addr = server_socket.accept()
|
||||
t = threading.Thread(target=handle_client, args=(client_socket, addr, pac_socket), daemon=True)
|
||||
t.start()
|
||||
except socket.timeout:
|
||||
await writer.wait_closed()
|
||||
except (ConnectionError, OSError):
|
||||
pass
|
||||
await self.close_clients()
|
||||
|
||||
# Reconnexion si la PAC s'est déconnectée
|
||||
if not reader_thread.is_alive():
|
||||
log("Thread lecture PAC mort, reconnexion...")
|
||||
with clients_lock:
|
||||
for c in clients:
|
||||
try:
|
||||
c.close()
|
||||
except Exception:
|
||||
pass
|
||||
clients.clear()
|
||||
pac_socket = connect_to_pac()
|
||||
if pac_socket is None:
|
||||
break
|
||||
reader_thread = threading.Thread(target=pac_reader, args=(pac_socket,), daemon=True)
|
||||
reader_thread.start()
|
||||
keepalive_thread = threading.Thread(target=pac_keepalive, args=(pac_socket,), daemon=True)
|
||||
keepalive_thread.start()
|
||||
except Exception as e:
|
||||
log(f"Erreur serveur proxy : {e}")
|
||||
finally:
|
||||
stop_event.set()
|
||||
server_socket.close()
|
||||
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:
|
||||
pac_socket.close()
|
||||
except Exception:
|
||||
pass
|
||||
log("Proxy arrêté.")
|
||||
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__":
|
||||
try:
|
||||
start_proxy()
|
||||
except KeyboardInterrupt:
|
||||
log("Arrêt demandé")
|
||||
stop_event.set()
|
||||
sys.exit(0)
|
||||
main()
|
||||
|
||||
Reference in New Issue
Block a user