fix the manager

This commit is contained in:
2026-07-06 12:18:41 +02:00
parent 8891e78b1d
commit de468929e1
22 changed files with 1069 additions and 468 deletions

View File

@@ -12,9 +12,10 @@ import os
import socket
import threading
import time
from typing import Any, Dict, Iterable, Optional
from typing import Any, Dict, Optional
import docker
from config_model import load_config
CONFIG_PATH = os.getenv("CONFIG_PATH", "config.json")
IMAGE_NAME = os.getenv("IMAGE_NAME", "fis")
@@ -23,23 +24,12 @@ GATEWAY_PORT = int(os.getenv("GATEWAY_PORT", "0"))
INACTIVITY_SECONDS = int(os.getenv("INACTIVITY_SECONDS", "10"))
IDLE_CHECK_INTERVAL = int(os.getenv("IDLE_CHECK_INTERVAL", "10"))
STARTUP_TIMEOUT_SECONDS = int(os.getenv("STARTUP_TIMEOUT_SECONDS", "60"))
STARTUP_STABILIZATION_SECONDS = float(os.getenv("STARTUP_STABILIZATION_SECONDS", "5"))
TCP_CONNECT_TIMEOUT_SECONDS = float(os.getenv("TCP_CONNECT_TIMEOUT_SECONDS", "1.5"))
HTTP_READY_PATH = os.getenv("HTTP_READY_PATH", "/")
HTTP_READY_TIMEOUT_SECONDS = float(os.getenv("HTTP_READY_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"))
DEFAULTS: Dict[str, Any] = {
"container_prefix": "fis-",
"environment": {},
"volumes": {},
"mem_limit": None,
"memswap_limit": None,
"privileged": False,
"tmpfs": {},
"dns": [],
"extra_hosts": {},
}
_docker_client = None
user_activity: Dict[str, float] = {}
users_by_id: Dict[str, Dict[str, Any]] = {}
@@ -69,104 +59,6 @@ def get_user_container_lock(user_id: str) -> threading.Lock:
return user_container_locks[user_id]
def merge_dict(base: Dict[str, Any], override: Optional[Dict[str, Any]]) -> Dict[str, Any]:
result = dict(base or {})
result.update(override or {})
return result
def resolve_placeholders(value: Any, variables: Dict[str, Any]) -> Any:
if isinstance(value, str):
return value.format(**variables)
if isinstance(value, dict):
return {
resolve_placeholders(key, variables): resolve_placeholders(item, variables)
for key, item in value.items()
}
if isinstance(value, list):
return [resolve_placeholders(item, variables) for item in value]
return value
def merge_defaults(config: Dict[str, Any], user: Dict[str, Any]) -> Dict[str, Any]:
defaults = merge_dict(DEFAULTS, config.get("defaults"))
merged = dict(user)
merged["id"] = str(user["id"])
merged["port"] = int(user["port"])
variables = {"id": merged["id"], "port": merged["port"]}
merged["environment"] = merge_dict(
resolve_placeholders(defaults.get("environment", {}), variables),
resolve_placeholders(user.get("environment"), variables),
)
merged["volumes"] = merge_dict(
resolve_placeholders(defaults.get("volumes", {}), variables),
resolve_placeholders(user.get("volumes"), variables),
)
merged["extra_hosts"] = merge_dict(
resolve_placeholders(defaults.get("extra_hosts", {}), variables),
resolve_placeholders(user.get("extra_hosts"), variables),
)
merged["tmpfs"] = merge_dict(
resolve_placeholders(defaults.get("tmpfs", {}), variables),
resolve_placeholders(user.get("tmpfs"), variables),
)
merged["dns"] = resolve_placeholders(user.get("dns", defaults.get("dns", [])), variables)
merged["mem_limit"] = user.get("mem_limit", defaults.get("mem_limit"))
merged["memswap_limit"] = user.get("memswap_limit", defaults.get("memswap_limit"))
merged["privileged"] = user.get("privileged", defaults.get("privileged", False))
if not merged.get("container_name"):
prefix = defaults.get("container_prefix", "fis-")
merged["container_name"] = f"{prefix}{merged['id']}"
if not merged["volumes"]:
merged["volumes"] = {f"v-conf-{merged['id']}": "/config"}
return merged
def validate_users(users: Iterable[Dict[str, Any]]) -> None:
seen_ids = set()
seen_ports = set()
seen_names = set()
for user in users:
uid = user["id"]
port = user["port"]
name = user["container_name"]
if uid in seen_ids:
raise ValueError(f"Duplicate user id: {uid}")
if port in seen_ports:
raise ValueError(f"Duplicate port: {port}")
if name in seen_names:
raise ValueError(f"Duplicate container_name: {name}")
if not user["volumes"]:
raise ValueError(f"User {uid} must define at least one volume")
if "PASSWORD" not in user["environment"] or not user["environment"]["PASSWORD"]:
raise ValueError(f"User {uid} must define environment.PASSWORD")
seen_ids.add(uid)
seen_ports.add(port)
seen_names.add(name)
def load_config(path: str = CONFIG_PATH) -> Dict[str, Any]:
with open(path, encoding="utf-8") as handle:
raw = json.load(handle)
raw_users = raw.get("users", [])
if not raw_users:
raise ValueError("Configuration must contain at least one user")
users = [merge_defaults(raw, user) for user in raw_users]
validate_users(users)
return {"defaults": merge_dict(DEFAULTS, raw.get("defaults")), "users": users}
def get_container_ip(container_name: str) -> Optional[str]:
"""Get IP address of a running container."""
try:
@@ -193,6 +85,30 @@ def is_port_open(host: str, port: int, timeout_seconds: float = TCP_CONNECT_TIME
return False
def is_http_ready(host: str, port: int, path: str = HTTP_READY_PATH) -> bool:
try:
with socket.create_connection((host, port), timeout=HTTP_READY_TIMEOUT_SECONDS) as sock:
sock.settimeout(HTTP_READY_TIMEOUT_SECONDS)
request = (
f"GET {path or '/'} HTTP/1.1\r\n"
f"Host: {host}:{port}\r\n"
"Connection: close\r\n"
"\r\n"
)
sock.sendall(request.encode("ascii"))
response = sock.recv(4096)
except OSError:
return False
status_line = response.split(b"\r\n", 1)[0]
parts = status_line.split()
if len(parts) < 2 or not parts[1].isdigit():
return False
status_code = int(parts[1])
return status_code in {200, 204, 301, 302, 303, 307, 308, 401, 403}
async def read_initial_http_request(reader) -> bytes:
buffer = bytearray()
while len(buffer) < MAX_REQUEST_HEADER_BYTES:
@@ -300,20 +216,16 @@ def extract_user_from_http_request(data: bytes, known_users: Dict[str, Dict[str,
def wait_until_ready(container_name: str) -> bool:
stabilized = False
waiting_for_http = False
deadline = time.time() + STARTUP_TIMEOUT_SECONDS
while time.time() < deadline:
ip = get_container_ip(container_name)
if ip and is_port_open(ip, APP_PORT):
if STARTUP_STABILIZATION_SECONDS > 0 and not stabilized:
log(
f"[{container_name}] Port {APP_PORT} is open, waiting "
f"{STARTUP_STABILIZATION_SECONDS}s for backend stabilization"
)
time.sleep(STARTUP_STABILIZATION_SECONDS)
stabilized = True
continue
return True
if is_http_ready(ip, APP_PORT):
return True
if not waiting_for_http:
log(f"[{container_name}] Port {APP_PORT} is open, waiting for HTTP readiness")
waiting_for_http = True
time.sleep(0.5)
return False
@@ -364,7 +276,7 @@ def start_user_container(user: Dict[str, Any]):
extra_hosts=user.get("extra_hosts") or None,
tty=True,
stdin_open=True,
restart_policy={"Name": "no"},
restart_policy={"Name": "always"},
)
except docker.errors.APIError as exc:
if "already in use" not in str(exc).lower():