Aggiunto un log migliore
This commit is contained in:
@@ -28,6 +28,7 @@ STARTUP_STABILIZATION_SECONDS = float(os.getenv("STARTUP_STABILIZATION_SECONDS",
|
||||
TCP_CONNECT_TIMEOUT_SECONDS = float(os.getenv("TCP_CONNECT_TIMEOUT_SECONDS", "1.5"))
|
||||
INITIAL_REQUEST_TIMEOUT_SECONDS = float(os.getenv("INITIAL_REQUEST_TIMEOUT_SECONDS", "3"))
|
||||
MAX_REQUEST_HEADER_BYTES = int(os.getenv("MAX_REQUEST_HEADER_BYTES", "65536"))
|
||||
GATEWAY_ROUTE_TTL_SECONDS = int(os.getenv("GATEWAY_ROUTE_TTL_SECONDS", "3600"))
|
||||
|
||||
DEFAULTS: Dict[str, Any] = {
|
||||
"container_prefix": "fis-",
|
||||
@@ -44,6 +45,7 @@ DEFAULTS: Dict[str, Any] = {
|
||||
_docker_client = None
|
||||
user_activity: Dict[str, float] = {}
|
||||
users_by_id: Dict[str, Dict[str, Any]] = {}
|
||||
gateway_client_routes: Dict[str, Dict[str, Any]] = {}
|
||||
user_container_locks: Dict[str, threading.Lock] = {}
|
||||
user_container_locks_guard = threading.Lock()
|
||||
|
||||
@@ -246,6 +248,28 @@ def rewrite_http_request_target(data: bytes, target: str) -> bytes:
|
||||
return rewritten_line + data[header_end:]
|
||||
|
||||
|
||||
def remember_gateway_route(client_ip: str, user_id: str) -> None:
|
||||
if not client_ip:
|
||||
return
|
||||
gateway_client_routes[client_ip] = {"user_id": user_id, "updated_at": time.time()}
|
||||
|
||||
|
||||
def get_remembered_gateway_user(client_ip: str) -> Optional[Dict[str, Any]]:
|
||||
if not client_ip:
|
||||
return None
|
||||
|
||||
route = gateway_client_routes.get(client_ip)
|
||||
if not route:
|
||||
return None
|
||||
|
||||
if time.time() - route["updated_at"] > GATEWAY_ROUTE_TTL_SECONDS:
|
||||
gateway_client_routes.pop(client_ip, None)
|
||||
return None
|
||||
|
||||
user_id = route["user_id"]
|
||||
return users_by_id.get(user_id)
|
||||
|
||||
|
||||
def extract_user_from_http_request(data: bytes, known_users: Dict[str, Dict[str, Any]]):
|
||||
if not data:
|
||||
return None, data, None
|
||||
@@ -466,6 +490,7 @@ async def handle_connection(reader, writer, user: Dict[str, Any], initial_data:
|
||||
|
||||
async def handle_gateway_connection(reader, writer):
|
||||
peer = writer.get_extra_info("peername")
|
||||
client_ip = peer[0] if isinstance(peer, tuple) and peer else None
|
||||
try:
|
||||
initial_data = await read_initial_http_request(reader)
|
||||
except asyncio.TimeoutError:
|
||||
@@ -476,6 +501,9 @@ async def handle_gateway_connection(reader, writer):
|
||||
return
|
||||
|
||||
user, rewritten_data, redirect_target = extract_user_from_http_request(initial_data, users_by_id)
|
||||
if user:
|
||||
remember_gateway_route(client_ip, user["id"])
|
||||
|
||||
if redirect_target:
|
||||
writer.write(build_http_redirect(redirect_target))
|
||||
await writer.drain()
|
||||
@@ -483,6 +511,13 @@ async def handle_gateway_connection(reader, writer):
|
||||
await writer.wait_closed()
|
||||
return
|
||||
|
||||
if not user:
|
||||
remembered_user = get_remembered_gateway_user(client_ip)
|
||||
if remembered_user:
|
||||
user = remembered_user
|
||||
rewritten_data = initial_data
|
||||
log(f"[gateway] Reused remembered route for {client_ip} -> user {user['id']}")
|
||||
|
||||
if not user:
|
||||
writer.write(
|
||||
build_http_error(
|
||||
@@ -506,6 +541,10 @@ async def idle_checker():
|
||||
await asyncio.sleep(IDLE_CHECK_INTERVAL)
|
||||
now = time.time()
|
||||
|
||||
for client_ip, route in list(gateway_client_routes.items()):
|
||||
if now - route["updated_at"] > GATEWAY_ROUTE_TTL_SECONDS:
|
||||
gateway_client_routes.pop(client_ip, None)
|
||||
|
||||
for user_id, last_active in list(user_activity.items()):
|
||||
if now - last_active <= INACTIVITY_SECONDS:
|
||||
continue
|
||||
|
||||
Reference in New Issue
Block a user