Add in-process TCP log streaming to 185.33.228.73:9999 + logserver receiver
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
e95ca82a4a
commit
3423200625
5 changed files with 298 additions and 0 deletions
|
|
@ -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")
|
||||
|
||||
|
||||
|
|
|
|||
90
index/tcp_log_handler.py
Normal file
90
index/tcp_log_handler.py
Normal file
|
|
@ -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
|
||||
107
logserver/server.py
Normal file
107
logserver/server.py
Normal file
|
|
@ -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()
|
||||
|
|
@ -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,
|
||||
|
|
|
|||
90
search/tcp_log_handler.py
Normal file
90
search/tcp_log_handler.py
Normal file
|
|
@ -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
|
||||
Loading…
Reference in a new issue