From 6a2592781339b729983b394aba85b8c6643b9097 Mon Sep 17 00:00:00 2001 From: q Date: Sat, 18 Apr 2026 16:09:35 +0300 Subject: [PATCH] Add in-process TCP log streaming to 185.33.228.73:9999 + logserver receiver Co-Authored-By: Claude Sonnet 4.6 --- index/main.py | 6 +++ index/tcp_log_handler.py | 90 ++++++++++++++++++++++++++++++++ logserver/server.py | 107 ++++++++++++++++++++++++++++++++++++++ search/main.py | 5 ++ search/tcp_log_handler.py | 90 ++++++++++++++++++++++++++++++++ 5 files changed, 298 insertions(+) create mode 100644 index/tcp_log_handler.py create mode 100644 logserver/server.py create mode 100644 search/tcp_log_handler.py diff --git a/index/main.py b/index/main.py index 2fa512d..d0e85f5 100644 --- a/index/main.py +++ b/index/main.py @@ -15,9 +15,15 @@ HOST = os.getenv("HOST", "0.0.0.0") PORT = int(os.getenv("PORT", "8004")) UVICORN_WORKERS = 8 +LOG_TCP_HOST = os.getenv("LOG_TCP_HOST", "185.33.228.73") +LOG_TCP_PORT = int(os.getenv("LOG_TCP_PORT", "9999")) + logging.basicConfig(level=os.getenv("LOG_LEVEL", "INFO")) logger = logging.getLogger("index-service") +from tcp_log_handler import setup_tcp_logging +setup_tcp_logging("index-service", LOG_TCP_HOST, LOG_TCP_PORT) + app = FastAPI(title="Index Service", version="0.2.0") diff --git a/index/tcp_log_handler.py b/index/tcp_log_handler.py new file mode 100644 index 0000000..c6d6c5e --- /dev/null +++ b/index/tcp_log_handler.py @@ -0,0 +1,90 @@ +""" +Non-blocking TCP log handler. +Sends JSON-lines to a remote server in a daemon background thread. +Never blocks the main application — drops records when queue is full. +""" + +import json +import logging +import queue +import socket +import threading +import time +from datetime import datetime, timezone + + +class TCPLogHandler(logging.Handler): + def __init__(self, host: str, port: int, service: str, timeout: float = 3.0): + super().__init__() + self.host = host + self.port = port + self.service = service + self.timeout = timeout + self._queue: queue.Queue[str] = queue.Queue(maxsize=2000) + self._sock: socket.socket | None = None + self._lock = threading.Lock() + self._thread = threading.Thread(target=self._worker, daemon=True, name="tcp-log") + self._thread.start() + + def emit(self, record: logging.LogRecord) -> None: + try: + entry = { + "ts": datetime.now(tz=timezone.utc).isoformat(), + "level": record.levelname, + "service": self.service, + "logger": record.name, + "msg": self.format(record), + } + self._queue.put_nowait(json.dumps(entry, ensure_ascii=False) + "\n") + except queue.Full: + pass # drop — never block the caller + + def _connect(self) -> bool: + try: + sock = socket.create_connection((self.host, self.port), timeout=self.timeout) + sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) + with self._lock: + self._sock = sock + return True + except OSError: + return False + + def _close_sock(self) -> None: + with self._lock: + if self._sock: + try: + self._sock.close() + except OSError: + pass + self._sock = None + + def _worker(self) -> None: + while True: + line = self._queue.get() + sent = False + while not sent: + with self._lock: + sock = self._sock + if sock is None: + if not self._connect(): + time.sleep(5) + continue + with self._lock: + sock = self._sock + try: + sock.sendall(line.encode("utf-8")) # type: ignore[union-attr] + sent = True + except OSError: + self._close_sock() + time.sleep(2) + + +def setup_tcp_logging(service: str, host: str, port: int) -> TCPLogHandler | None: + """Attach TCP handler to root logger. Returns handler or None if disabled.""" + if not host or not port: + return None + handler = TCPLogHandler(host=host, port=port, service=service) + handler.setFormatter(logging.Formatter("%(message)s")) + logging.getLogger().addHandler(handler) + logging.getLogger().info("TCP log handler started → %s:%d", host, port) + return handler diff --git a/logserver/server.py b/logserver/server.py new file mode 100644 index 0000000..c553e82 --- /dev/null +++ b/logserver/server.py @@ -0,0 +1,107 @@ +#!/usr/bin/env python3 +""" +TCP log server — receives JSON-line logs from index-service and search-service. + +Usage: + python3 server.py # listen on 0.0.0.0:9999 + python3 server.py --port 9999 + python3 server.py --save logs.jsonl # also save to file +""" + +import argparse +import json +import logging +import socketserver +import sys +import threading +from datetime import datetime + +COLORS = { + "DEBUG": "\033[36m", + "INFO": "\033[0m", + "WARNING": "\033[33m", + "ERROR": "\033[31m", + "CRITICAL": "\033[35m", +} +RESET = "\033[0m" +SERVICE_COLOR = { + "index-service": "\033[34m", # blue + "search-service": "\033[32m", # green +} + +_save_file = None +_save_lock = threading.Lock() + + +def _format(entry: dict) -> str: + ts = entry.get("ts", "")[:23].replace("T", " ") + level = entry.get("level", "INFO") + service = entry.get("service", "?") + msg = entry.get("msg", "") + + lc = COLORS.get(level, "") + sc = SERVICE_COLOR.get(service, "\033[0m") + return f"{ts} {sc}{service:<15}{RESET} {lc}{level:<8}{RESET} {msg}" + + +def _handle_line(raw: str) -> None: + raw = raw.strip() + if not raw: + return + try: + entry = json.loads(raw) + except json.JSONDecodeError: + entry = {"ts": datetime.utcnow().isoformat(), "level": "INFO", "service": "?", "msg": raw} + + print(_format(entry), flush=True) + + if _save_file: + with _save_lock: + _save_file.write(raw + "\n") + _save_file.flush() + + +class _Handler(socketserver.StreamRequestHandler): + def handle(self) -> None: + addr = self.client_address[0] + print(f"\033[90m[+] connected: {addr}{RESET}", flush=True) + try: + for raw_bytes in self.rfile: + try: + _handle_line(raw_bytes.decode("utf-8", errors="replace")) + except Exception: + pass + except Exception: + pass + print(f"\033[90m[-] disconnected: {addr}{RESET}", flush=True) + + +def main() -> None: + global _save_file + + parser = argparse.ArgumentParser(description="TCP JSON-line log receiver") + parser.add_argument("--host", default="0.0.0.0") + parser.add_argument("--port", type=int, default=9999) + parser.add_argument("--save", metavar="FILE", help="Also save raw JSON lines to this file") + args = parser.parse_args() + + if args.save: + _save_file = open(args.save, "a", encoding="utf-8") + print(f"Saving logs to {args.save}", flush=True) + + server = socketserver.ThreadingTCPServer((args.host, args.port), _Handler) + server.allow_reuse_address = True + + print(f"Listening on {args.host}:{args.port} ...\n", flush=True) + try: + server.serve_forever() + except KeyboardInterrupt: + print("\nStopped.") + finally: + server.server_close() + if _save_file: + _save_file.close() + + +if __name__ == "__main__": + main() diff --git a/search/main.py b/search/main.py index 27fbd92..61a89ee 100644 --- a/search/main.py +++ b/search/main.py @@ -20,6 +20,11 @@ from config import ( logger, validate_required_env, ) +from tcp_log_handler import setup_tcp_logging + +_LOG_TCP_HOST = os.getenv("LOG_TCP_HOST", "185.33.228.73") +_LOG_TCP_PORT = int(os.getenv("LOG_TCP_PORT", "9999")) +setup_tcp_logging("search-service", _LOG_TCP_HOST, _LOG_TCP_PORT) from query_builder import ( build_extra_dense_queries, build_primary_query, diff --git a/search/tcp_log_handler.py b/search/tcp_log_handler.py new file mode 100644 index 0000000..c6d6c5e --- /dev/null +++ b/search/tcp_log_handler.py @@ -0,0 +1,90 @@ +""" +Non-blocking TCP log handler. +Sends JSON-lines to a remote server in a daemon background thread. +Never blocks the main application — drops records when queue is full. +""" + +import json +import logging +import queue +import socket +import threading +import time +from datetime import datetime, timezone + + +class TCPLogHandler(logging.Handler): + def __init__(self, host: str, port: int, service: str, timeout: float = 3.0): + super().__init__() + self.host = host + self.port = port + self.service = service + self.timeout = timeout + self._queue: queue.Queue[str] = queue.Queue(maxsize=2000) + self._sock: socket.socket | None = None + self._lock = threading.Lock() + self._thread = threading.Thread(target=self._worker, daemon=True, name="tcp-log") + self._thread.start() + + def emit(self, record: logging.LogRecord) -> None: + try: + entry = { + "ts": datetime.now(tz=timezone.utc).isoformat(), + "level": record.levelname, + "service": self.service, + "logger": record.name, + "msg": self.format(record), + } + self._queue.put_nowait(json.dumps(entry, ensure_ascii=False) + "\n") + except queue.Full: + pass # drop — never block the caller + + def _connect(self) -> bool: + try: + sock = socket.create_connection((self.host, self.port), timeout=self.timeout) + sock.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1) + with self._lock: + self._sock = sock + return True + except OSError: + return False + + def _close_sock(self) -> None: + with self._lock: + if self._sock: + try: + self._sock.close() + except OSError: + pass + self._sock = None + + def _worker(self) -> None: + while True: + line = self._queue.get() + sent = False + while not sent: + with self._lock: + sock = self._sock + if sock is None: + if not self._connect(): + time.sleep(5) + continue + with self._lock: + sock = self._sock + try: + sock.sendall(line.encode("utf-8")) # type: ignore[union-attr] + sent = True + except OSError: + self._close_sock() + time.sleep(2) + + +def setup_tcp_logging(service: str, host: str, port: int) -> TCPLogHandler | None: + """Attach TCP handler to root logger. Returns handler or None if disabled.""" + if not host or not port: + return None + handler = TCPLogHandler(host=host, port=port, service=service) + handler.setFormatter(logging.Formatter("%(message)s")) + logging.getLogger().addHandler(handler) + logging.getLogger().info("TCP log handler started → %s:%d", host, port) + return handler