Compare commits
3 Commits
1e6eadeed2
...
85bcb304f6
| Author | SHA1 | Date | |
|---|---|---|---|
| 85bcb304f6 | |||
| 53e18d92ac | |||
| 2e0d4c78ca |
+38
-62
@@ -120,49 +120,56 @@ def try_acquire_maintenance_leader() -> bool:
|
|||||||
return True
|
return True
|
||||||
|
|
||||||
|
|
||||||
def run_maintenance_service() -> None:
|
def _init_schema() -> None:
|
||||||
logger.info("maintenance_service_bootstrap_started")
|
|
||||||
with open("/tmp/portal-schema.lock", "w") as lock_file:
|
with open("/tmp/portal-schema.lock", "w") as lock_file:
|
||||||
fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
|
fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
|
||||||
Base.metadata.create_all(bind=engine)
|
Base.metadata.create_all(bind=engine)
|
||||||
ensure_schema_compatibility()
|
ensure_schema_compatibility()
|
||||||
fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
|
fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
|
||||||
|
|
||||||
ensure_icons_dir()
|
ensure_icons_dir()
|
||||||
bootstrap_admin()
|
bootstrap_admin()
|
||||||
|
|
||||||
|
|
||||||
|
def _startup_pools(db) -> None:
|
||||||
|
ensure_universal_pool()
|
||||||
|
ensure_web_pool()
|
||||||
|
for svc in db.scalars(
|
||||||
|
select(Service).where(
|
||||||
|
Service.active == True,
|
||||||
|
Service.type.in_([ServiceType.WEB, ServiceType.RDP]),
|
||||||
|
)
|
||||||
|
).all():
|
||||||
|
if svc.type == ServiceType.WEB and WEB_POOL_SIZE <= 0:
|
||||||
|
ensure_warm_pool(svc)
|
||||||
|
elif svc.type == ServiceType.RDP:
|
||||||
|
slots = db.scalars(select(RdpSlot).where(RdpSlot.service_id == svc.id)).all()
|
||||||
|
for slot in slots:
|
||||||
|
try:
|
||||||
|
cname = _rdp_slot_container_name(svc.slug, slot.id)
|
||||||
|
try:
|
||||||
|
c = docker_client().containers.get(cname)
|
||||||
|
if c.status != "running":
|
||||||
|
c.start()
|
||||||
|
except docker.errors.NotFound:
|
||||||
|
start_rdp_slot_container(slot, svc)
|
||||||
|
slot.container_name = cname
|
||||||
|
except Exception:
|
||||||
|
logger.exception("startup_rdp_slot_start_failed slot_id=%s", slot.id)
|
||||||
|
if slots:
|
||||||
|
db.commit()
|
||||||
|
|
||||||
|
|
||||||
|
def run_maintenance_service() -> None:
|
||||||
|
logger.info("maintenance_service_bootstrap_started")
|
||||||
|
_init_schema()
|
||||||
|
|
||||||
maintenance_lock = open("/tmp/portal-maintenance.lock", "w")
|
maintenance_lock = open("/tmp/portal-maintenance.lock", "w")
|
||||||
fcntl.flock(maintenance_lock.fileno(), fcntl.LOCK_EX)
|
fcntl.flock(maintenance_lock.fileno(), fcntl.LOCK_EX)
|
||||||
logger.info("maintenance_service_leader_acquired")
|
logger.info("maintenance_service_leader_acquired")
|
||||||
|
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
try:
|
try:
|
||||||
ensure_universal_pool()
|
_startup_pools(db)
|
||||||
ensure_web_pool()
|
|
||||||
for svc in db.scalars(
|
|
||||||
select(Service).where(
|
|
||||||
Service.active == True,
|
|
||||||
Service.type.in_([ServiceType.WEB, ServiceType.RDP]),
|
|
||||||
)
|
|
||||||
).all():
|
|
||||||
if svc.type == ServiceType.WEB and WEB_POOL_SIZE <= 0:
|
|
||||||
ensure_warm_pool(svc)
|
|
||||||
elif svc.type == ServiceType.RDP:
|
|
||||||
slots = db.scalars(select(RdpSlot).where(RdpSlot.service_id == svc.id)).all()
|
|
||||||
for slot in slots:
|
|
||||||
try:
|
|
||||||
cname = _rdp_slot_container_name(svc.slug, slot.id)
|
|
||||||
try:
|
|
||||||
c = docker_client().containers.get(cname)
|
|
||||||
if c.status != "running":
|
|
||||||
c.start()
|
|
||||||
except docker.errors.NotFound:
|
|
||||||
start_rdp_slot_container(slot, svc)
|
|
||||||
slot.container_name = cname
|
|
||||||
except Exception:
|
|
||||||
logger.exception("startup_rdp_slot_start_failed slot_id=%s", slot.id)
|
|
||||||
if slots:
|
|
||||||
db.commit()
|
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
db.close()
|
||||||
|
|
||||||
@@ -171,13 +178,7 @@ def run_maintenance_service() -> None:
|
|||||||
|
|
||||||
|
|
||||||
def on_startup() -> None:
|
def on_startup() -> None:
|
||||||
with open("/tmp/portal-schema.lock", "w") as lock_file:
|
_init_schema()
|
||||||
fcntl.flock(lock_file.fileno(), fcntl.LOCK_EX)
|
|
||||||
Base.metadata.create_all(bind=engine)
|
|
||||||
ensure_schema_compatibility()
|
|
||||||
fcntl.flock(lock_file.fileno(), fcntl.LOCK_UN)
|
|
||||||
ensure_icons_dir()
|
|
||||||
bootstrap_admin()
|
|
||||||
if not try_acquire_maintenance_leader():
|
if not try_acquire_maintenance_leader():
|
||||||
logger.info("maintenance_leader_skipped")
|
logger.info("maintenance_leader_skipped")
|
||||||
return
|
return
|
||||||
@@ -185,32 +186,7 @@ def on_startup() -> None:
|
|||||||
if ENABLE_STARTUP_MAINTENANCE:
|
if ENABLE_STARTUP_MAINTENANCE:
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
try:
|
try:
|
||||||
ensure_universal_pool()
|
_startup_pools(db)
|
||||||
ensure_web_pool()
|
|
||||||
for svc in db.scalars(
|
|
||||||
select(Service).where(
|
|
||||||
Service.active == True,
|
|
||||||
Service.type.in_([ServiceType.WEB, ServiceType.RDP]),
|
|
||||||
)
|
|
||||||
).all():
|
|
||||||
if svc.type == ServiceType.WEB and WEB_POOL_SIZE <= 0:
|
|
||||||
ensure_warm_pool(svc)
|
|
||||||
elif svc.type == ServiceType.RDP:
|
|
||||||
slots = db.scalars(select(RdpSlot).where(RdpSlot.service_id == svc.id)).all()
|
|
||||||
for slot in slots:
|
|
||||||
try:
|
|
||||||
cname = _rdp_slot_container_name(svc.slug, slot.id)
|
|
||||||
try:
|
|
||||||
c = docker_client().containers.get(cname)
|
|
||||||
if c.status != "running":
|
|
||||||
c.start()
|
|
||||||
except docker.errors.NotFound:
|
|
||||||
start_rdp_slot_container(slot, svc)
|
|
||||||
slot.container_name = cname
|
|
||||||
except Exception:
|
|
||||||
logger.exception("startup_rdp_slot_start_failed slot_id=%s", slot.id)
|
|
||||||
if slots:
|
|
||||||
db.commit()
|
|
||||||
finally:
|
finally:
|
||||||
db.close()
|
db.close()
|
||||||
|
|
||||||
|
|||||||
+25
-47
@@ -194,16 +194,15 @@ def ensure_web_pool(target_size: Optional[int] = None) -> None:
|
|||||||
break
|
break
|
||||||
|
|
||||||
|
|
||||||
def get_universal_pool_status() -> dict:
|
def _get_pool_status(desired_size: int, name_fn) -> dict:
|
||||||
desired = max(0, UNIVERSAL_POOL_SIZE)
|
desired = max(0, desired_size)
|
||||||
if desired <= 0:
|
if desired <= 0:
|
||||||
return {"desired": 0, "running": 0, "total": 0, "health": "down", "names": []}
|
return {"desired": 0, "running": 0, "total": 0, "health": "down", "names": []}
|
||||||
d = docker_client()
|
d = docker_client()
|
||||||
names = [universal_container_name(i) for i in range(desired)]
|
|
||||||
containers = []
|
containers = []
|
||||||
for name in names:
|
for i in range(desired):
|
||||||
try:
|
try:
|
||||||
containers.append(d.containers.get(name))
|
containers.append(d.containers.get(name_fn(i)))
|
||||||
except Exception:
|
except Exception:
|
||||||
continue
|
continue
|
||||||
running = sum(1 for c in containers if c.status == "running")
|
running = sum(1 for c in containers if c.status == "running")
|
||||||
@@ -217,27 +216,12 @@ def get_universal_pool_status() -> dict:
|
|||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def get_universal_pool_status() -> dict:
|
||||||
|
return _get_pool_status(UNIVERSAL_POOL_SIZE, universal_container_name)
|
||||||
|
|
||||||
|
|
||||||
def get_web_pool_status() -> dict:
|
def get_web_pool_status() -> dict:
|
||||||
desired = max(0, WEB_POOL_SIZE)
|
return _get_pool_status(WEB_POOL_SIZE, web_pool_container_name)
|
||||||
if desired <= 0:
|
|
||||||
return {"desired": 0, "running": 0, "total": 0, "health": "down", "names": []}
|
|
||||||
d = docker_client()
|
|
||||||
names = [web_pool_container_name(i) for i in range(desired)]
|
|
||||||
containers = []
|
|
||||||
for name in names:
|
|
||||||
try:
|
|
||||||
containers.append(d.containers.get(name))
|
|
||||||
except Exception:
|
|
||||||
continue
|
|
||||||
running = sum(1 for c in containers if c.status == "running")
|
|
||||||
health = "ok" if running >= min(desired, 1) else "down"
|
|
||||||
return {
|
|
||||||
"desired": desired,
|
|
||||||
"running": running,
|
|
||||||
"total": len(containers),
|
|
||||||
"names": sorted(c.name for c in containers),
|
|
||||||
"health": health,
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
def acquire_universal_slot(db: Session) -> int:
|
def acquire_universal_slot(db: Session) -> int:
|
||||||
@@ -303,6 +287,20 @@ def sanitize_client_resolution(width: Optional[int], height: Optional[int]) -> t
|
|||||||
return clamped_width, clamped_height
|
return clamped_width, clamped_height
|
||||||
|
|
||||||
|
|
||||||
|
def _dispatch_post(url: str, payload: dict) -> None:
|
||||||
|
last_exc = None
|
||||||
|
for _ in range(max(1, POOL_DISPATCH_RETRIES)):
|
||||||
|
try:
|
||||||
|
resp = requests.post(url, json=payload, timeout=POOL_DISPATCH_REQUEST_TIMEOUT_SECONDS)
|
||||||
|
resp.raise_for_status()
|
||||||
|
return
|
||||||
|
except Exception as exc:
|
||||||
|
last_exc = exc
|
||||||
|
time.sleep(max(0.0, POOL_DISPATCH_SLEEP_SECONDS))
|
||||||
|
if last_exc:
|
||||||
|
raise last_exc
|
||||||
|
|
||||||
|
|
||||||
def dispatch_universal_target(slot: int, service: Service, width: Optional[int] = None, height: Optional[int] = None) -> None:
|
def dispatch_universal_target(slot: int, service: Service, width: Optional[int] = None, height: Optional[int] = None) -> None:
|
||||||
name = universal_container_name(slot)
|
name = universal_container_name(slot)
|
||||||
url = ""
|
url = ""
|
||||||
@@ -328,17 +326,7 @@ def dispatch_universal_target(slot: int, service: Service, width: Optional[int]
|
|||||||
else:
|
else:
|
||||||
raise HTTPException(status_code=400, detail="Universal pool supports WEB/RDP only")
|
raise HTTPException(status_code=400, detail="Universal pool supports WEB/RDP only")
|
||||||
|
|
||||||
last_exc = None
|
_dispatch_post(url, payload)
|
||||||
for _ in range(max(1, POOL_DISPATCH_RETRIES)):
|
|
||||||
try:
|
|
||||||
resp = requests.post(url, json=payload, timeout=POOL_DISPATCH_REQUEST_TIMEOUT_SECONDS)
|
|
||||||
resp.raise_for_status()
|
|
||||||
return
|
|
||||||
except Exception as exc:
|
|
||||||
last_exc = exc
|
|
||||||
time.sleep(max(0.0, POOL_DISPATCH_SLEEP_SECONDS))
|
|
||||||
if last_exc:
|
|
||||||
raise last_exc
|
|
||||||
|
|
||||||
|
|
||||||
def dispatch_web_pool_target(slot: int, service: Service, width: Optional[int] = None, height: Optional[int] = None) -> None:
|
def dispatch_web_pool_target(slot: int, service: Service, width: Optional[int] = None, height: Optional[int] = None) -> None:
|
||||||
@@ -350,17 +338,7 @@ def dispatch_web_pool_target(slot: int, service: Service, width: Optional[int] =
|
|||||||
if width and height:
|
if width and height:
|
||||||
payload["width"] = width
|
payload["width"] = width
|
||||||
payload["height"] = height
|
payload["height"] = height
|
||||||
last_exc = None
|
_dispatch_post(url, payload)
|
||||||
for _ in range(max(1, POOL_DISPATCH_RETRIES)):
|
|
||||||
try:
|
|
||||||
resp = requests.post(url, json=payload, timeout=POOL_DISPATCH_REQUEST_TIMEOUT_SECONDS)
|
|
||||||
resp.raise_for_status()
|
|
||||||
return
|
|
||||||
except Exception as exc:
|
|
||||||
last_exc = exc
|
|
||||||
time.sleep(max(0.0, POOL_DISPATCH_SLEEP_SECONDS))
|
|
||||||
if last_exc:
|
|
||||||
raise last_exc
|
|
||||||
|
|
||||||
|
|
||||||
def create_runtime_container(service: Service, session_id: str):
|
def create_runtime_container(service: Service, session_id: str):
|
||||||
|
|||||||
Reference in New Issue
Block a user