Feat: Add community plane

- Implemented relevant ping ping routes, regions, cloudflare, auth as a proper abstract to Core for a community cloud
This commit is contained in:
2026-09-15 05:46:18 +05:45
parent b25c6cd24e
commit 812a1318e9
39 changed files with 3657 additions and 0 deletions
+55
View File
@@ -0,0 +1,55 @@
COMMUNITY_DATABASE_URL=mysql+pymysql://community_user:community_password@community-db:3306/community
PUBLIC_BASE_URL=http://localhost:5002
SIGNALING_PUBLIC_URL=ws://localhost:5003
SIGNALING_BROWSER_URL=
SIGNALING_INTERNAL_URL=http://community-signaling:5003
SIGNALING_SHARED_SECRET=change-me
STUN_URLS=stun:stun.l.google.com:19302
WORKER_BINARY_URL=
WORKER_BINARY_SHA256=
NSCONTROLLER_IMAGE_URL=
NSCONTROLLER_IMAGE_SHA256=
COMMUNITY_SECRET_KEY=
CORE_API_BASE_URL=http://api-server:5000
CORE_API_KEY=change-me-generate-a-real-key
CORE_PUBLIC_BASE_URL=
CORE_PUBLIC_WS_URL=
CORE_REGION_NAME_PREFIX=Community
HEARTBEAT_TIMEOUT_SECONDS=120
PING_TOKEN_TTL_SECONDS=120
LOCAL_USER_EMAIL=local@xcloudify.dev
BOOTSTRAP_ADMIN_EMAILS=
# production (default when unset) | development.
APP_ENV=production
FLASK_ENV=development
# Development only: no sign-in, identity from the X-Dev-User header (default
# LOCAL_USER_EMAIL). On by default when APP_ENV=development; set false there
# to test real sign-in. Ignored (and logged) in any other APP_ENV.
AUTH_DEV_BYPASS=
# Sign-in outside development. Same shape as cloud/.env.example -- community
# can use the same IdP with its own client id.
OIDC_ISSUER=https://secuird.tech/
OIDC_JWKS_URL=
OIDC_AUDIENCE=change-me-oidc-client-id
OIDC_ADDITIONAL_AUDIENCES=
OIDC_SUBJECT_CLAIM=sub
OIDC_EMAIL_CLAIM=email
OAUTH2_PROXY_CLIENT_ID=change-me-oidc-client-id
OAUTH2_PROXY_CLIENT_SECRET=change-me
OAUTH2_PROXY_COOKIE_SECRET=
OAUTH2_PROXY_REDIRECT_URL=http://localhost:8091/oauth2/callback
OAUTH2_PROXY_COOKIE_SECURE=false
OAUTH2_PROXY_UPSTREAMS=http://community-api:5002/api/,http://host.docker.internal:8081/
+20
View File
@@ -0,0 +1,20 @@
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1
WORKDIR /app
RUN apt-get update && apt-get install -y --no-install-recommends \
build-essential libmariadb-dev curl \
&& rm -rf /var/lib/apt/lists/*
COPY packages/pyshared /packages/pyshared
RUN pip install --upgrade pip && pip install -e /packages/pyshared
COPY community/requirements.txt .
RUN pip install -r requirements.txt
COPY community/ .
CMD ["bash"]
+4
View File
@@ -0,0 +1,4 @@
from app import app # noqa: F401
if __name__ == "__main__":
app.run(debug=True, port=5002, host="0.0.0.0")
+139
View File
@@ -0,0 +1,139 @@
import re
import uuid
from flask import Flask, g, request
from flask_sqlalchemy import SQLAlchemy
from flask_migrate import Migrate
from flask_cors import CORS
from xcloudify_shared import api_response, logger
from xcloudify_shared.env import dev_bypass_enabled
from runtime_urls import (
COMMUNITY_DATABASE_URL, PUBLIC_BASE_URL, WORKER_BINARY_URL, WORKER_BINARY_SHA256,
NSCONTROLLER_IMAGE_URL, NSCONTROLLER_IMAGE_SHA256,
SIGNALING_PUBLIC_URL, SIGNALING_BROWSER_URL, SIGNALING_INTERNAL_URL,
SIGNALING_SHARED_SECRET, STUN_URLS,
CORE_API_BASE_URL, CORE_API_KEY, CORE_PUBLIC_BASE_URL, CORE_PUBLIC_WS_URL,
CORE_REGION_NAME_PREFIX,
HEARTBEAT_TIMEOUT_SECONDS, PING_TOKEN_TTL_SECONDS,
LOCAL_USER_EMAIL, BOOTSTRAP_ADMIN_EMAILS,
APP_ENV, AUTH_DEV_BYPASS, OIDC_ISSUER, OIDC_JWKS_URL, OIDC_AUDIENCE,
OIDC_ADDITIONAL_AUDIENCES, OIDC_SUBJECT_CLAIM, OIDC_EMAIL_CLAIM,
)
app = Flask(__name__)
CORS(app, supports_credentials=True)
app.config["SQLALCHEMY_DATABASE_URI"] = COMMUNITY_DATABASE_URL
app.config["SQLALCHEMY_TRACK_MODIFICATIONS"] = False
app.config["PUBLIC_BASE_URL"] = PUBLIC_BASE_URL
app.config["SIGNALING_PUBLIC_URL"] = SIGNALING_PUBLIC_URL
app.config["SIGNALING_BROWSER_URL"] = SIGNALING_BROWSER_URL
app.config["SIGNALING_INTERNAL_URL"] = SIGNALING_INTERNAL_URL
app.config["SIGNALING_SHARED_SECRET"] = SIGNALING_SHARED_SECRET
app.config["STUN_URLS"] = STUN_URLS
app.config["WORKER_BINARY_URL"] = WORKER_BINARY_URL
app.config["WORKER_BINARY_SHA256"] = WORKER_BINARY_SHA256
app.config["NSCONTROLLER_IMAGE_URL"] = NSCONTROLLER_IMAGE_URL
app.config["NSCONTROLLER_IMAGE_SHA256"] = NSCONTROLLER_IMAGE_SHA256
app.config["CORE_API_BASE_URL"] = CORE_API_BASE_URL
app.config["CORE_API_KEY"] = CORE_API_KEY
app.config["CORE_PUBLIC_BASE_URL"] = CORE_PUBLIC_BASE_URL
app.config["CORE_PUBLIC_WS_URL"] = CORE_PUBLIC_WS_URL
app.config["CORE_REGION_NAME_PREFIX"] = CORE_REGION_NAME_PREFIX
app.config["HEARTBEAT_TIMEOUT_SECONDS"] = HEARTBEAT_TIMEOUT_SECONDS
app.config["PING_TOKEN_TTL_SECONDS"] = PING_TOKEN_TTL_SECONDS
app.config["LOCAL_USER_EMAIL"] = LOCAL_USER_EMAIL
app.config["BOOTSTRAP_ADMIN_EMAILS"] = BOOTSTRAP_ADMIN_EMAILS
app.config["OIDC_ISSUER"] = OIDC_ISSUER
app.config["OIDC_JWKS_URL"] = OIDC_JWKS_URL
app.config["OIDC_AUDIENCE"] = OIDC_AUDIENCE
app.config["OIDC_ADDITIONAL_AUDIENCES"] = OIDC_ADDITIONAL_AUDIENCES
app.config["OIDC_SUBJECT_CLAIM"] = OIDC_SUBJECT_CLAIM
app.config["OIDC_EMAIL_CLAIM"] = OIDC_EMAIL_CLAIM
app.config["APP_ENV"] = APP_ENV
app.config["AUTH_DEV_BYPASS"] = dev_bypass_enabled(APP_ENV, AUTH_DEV_BYPASS)
if app.config["AUTH_DEV_BYPASS"]:
logger.warning("AUTH DEV BYPASS ON (APP_ENV=%s): requests act as X-Dev-User or %s", APP_ENV, LOCAL_USER_EMAIL)
elif not OIDC_ISSUER:
logger.error("OIDC_ISSUER is empty and the dev bypass is off -- every signed-in route will return 401")
if not CORE_API_KEY:
logger.error("CORE_API_KEY is empty -- region creation and host lookup against core will fail")
db = SQLAlchemy(app)
migrate = Migrate(app, db)
from app import models # noqa: E402 needed for db.metadata / migrations
_PUBLIC_PATHS = {
"/api/healthz",
"/api/nodes/enroll",
"/api/nodes/heartbeat",
"/api/nodes/ping-auth",
"/api/nodes/worker-binary",
"/api/nodes/worker-binary/sha256",
"/api/nodes/nscontroller-image",
"/api/nodes/nscontroller-image/sha256",
}
_PUBLIC_PREFIXES = ("/api/nodes/install/", "/api/nodes/worker-runtime/")
_SOFT_PUBLIC_PATHS = {"/api/nodes/available"}
_PING_SESSION_RE = re.compile(r"^/api/nodes/[^/]+/ping-session$")
def _is_soft_public(path: str) -> bool:
return path in _SOFT_PUBLIC_PATHS or bool(_PING_SESSION_RE.match(path))
@app.before_request
def assign_request_id():
g.request_id = str(uuid.uuid4())
@app.before_request
def resolve_identity():
"""Seed g.current_user and refuse /api calls that need one and have none.
Node-credential routes authenticate the node themselves. The marketplace
and latency-demo routes are open to anonymous visitors, so identity is
only resolved there when the request carries some.
"""
if request.method == "OPTIONS":
return
g.current_user = None
path = request.path
if not path.startswith("/api"):
return
if path in _PUBLIC_PATHS or path.startswith(_PUBLIC_PREFIXES):
return
from app.auth_utils import authenticate_request, has_credentials
if _is_soft_public(path):
if has_credentials():
authenticate_request()
return
authenticate_request()
if g.current_user is None:
return api_response(success=False, status=401, message="Authentication required", error_type="UNAUTHORIZED")
@app.after_request
def log_request_id(response):
response.headers["X-Request-ID"] = g.request_id
return response
from app.routes import api_bp # noqa: E402
@api_bp.route("/healthz", methods=["GET"])
def healthz():
from flask import jsonify
return jsonify("ok")
app.register_blueprint(api_bp)
logger.debug("community app init complete")
+174
View File
@@ -0,0 +1,174 @@
"""Who a community request is.
Production: an OIDC access token (a bearer, or oauth2-proxy's forwarded one),
verified against the issuer and mapped to a User by its subject.
Development (APP_ENV=development, bypass not turned off): no token needed.
The X-Dev-User header names the account, defaulting to LOCAL_USER_EMAIL, so
one browser can act as both a provider and a customer -- you cannot lease
your own node. A token, when sent, still wins, so real sign-in can be tested
in dev too.
"""
from datetime import datetime, timedelta
from functools import wraps
from typing import Optional
from flask import current_app, g, has_request_context, request as flask_request
from xcloudify_shared import api_response, logger
from xcloudify_shared.oidc import OIDCVerifier, bearer_token
LOCAL_USER_EMAIL = "local@xcloudify.dev"
DEV_USER_HEADER = "X-Dev-User"
def dev_mode() -> bool:
return bool(current_app.config.get("AUTH_DEV_BYPASS"))
def _verifier() -> OIDCVerifier:
verifier = current_app.extensions.get("oidc_verifier")
if verifier is None:
verifier = OIDCVerifier.from_config(current_app.config, user_agent="xcloudify-community")
current_app.extensions["oidc_verifier"] = verifier
return verifier
def _is_bootstrap_admin(email: str) -> bool:
"""Whether this address is listed in BOOTSTRAP_ADMIN_EMAILS. Grant-only:
never revokes, so a deploy shipping a different value cannot silently
demote an operator."""
raw = current_app.config.get("BOOTSTRAP_ADMIN_EMAILS") or ""
listed = {e.strip().lower() for e in raw.split(",") if e.strip()}
return bool(email) and email.lower() in listed
def upsert_user_from_claims(claims: dict):
from app import db
from app.models import User
subject_claim = current_app.config.get("OIDC_SUBJECT_CLAIM", "sub")
email_claim = current_app.config.get("OIDC_EMAIL_CLAIM", "email")
external_id = f"oidc:{claims.get(subject_claim) or claims['sub']}"
email = (claims.get(email_claim) or f"{external_id[5:]}@users.noreply").strip().lower()
full_name = (claims.get("name") or "").strip()
given = claims.get("given_name") or (full_name.split(" ")[0] if full_name else None)
family = claims.get("family_name") or (" ".join(full_name.split(" ")[1:]) if full_name else None)
user = User.query.filter_by(external_id=external_id).first()
if user is None:
# An existing row with this email (created in dev, or before the IdP
# was wired) is claimed only on a verified address -- otherwise anyone
# who could set that email at the IdP would inherit the account.
by_email = User.query.filter_by(email=email).first()
if by_email is not None:
if claims.get("email_verified") is not True:
raise PermissionError(f"email {email} already has an account and is not verified by the IdP")
user = by_email
if user is None:
user = User(external_id=external_id, email=email)
db.session.add(user)
user.external_id = external_id
user.email = email
if given:
user.first_name = given
if family:
user.last_name = family
if claims.get("picture"):
user.avatar_url = claims["picture"]
if _is_bootstrap_admin(email):
user.is_platform_admin = True
now = datetime.utcnow()
if user.last_login_at is None or (now - user.last_login_at) > timedelta(minutes=5):
user.last_login_at = now
db.session.commit()
return user
def _dev_identity() -> str:
"""Anything with an `@` is accepted verbatim -- in dev the header *is* the identity."""
requested = (flask_request.headers.get(DEV_USER_HEADER) or "").strip().lower()
if requested and "@" in requested and len(requested) <= 255:
return requested
return current_app.config.get("LOCAL_USER_EMAIL", LOCAL_USER_EMAIL).strip().lower()
def _dev_user():
from app import db
from app.models import User
email = _dev_identity()
default_email = current_app.config.get("LOCAL_USER_EMAIL", LOCAL_USER_EMAIL).strip().lower()
user = User.query.filter_by(email=email).first()
if user is None:
user = User(external_id=f"dev:{email}", email=email,
first_name="Dev", last_name=email.split("@")[0],
is_platform_admin=(email == default_email) or _is_bootstrap_admin(email))
db.session.add(user)
user.last_login_at = datetime.utcnow()
db.session.commit()
return user
def has_credentials() -> bool:
"""Whether the request carries anything identity could be resolved from --
lets public routes stay anonymous for visitors who sent nothing."""
return bool(bearer_token(flask_request)) or (dev_mode() and bool(flask_request.headers.get(DEV_USER_HEADER)))
def authenticate_request():
"""Resolve g.current_user, or leave it None."""
g.current_user = None
token = bearer_token(flask_request)
if token:
verifier = _verifier()
try:
g.current_user = upsert_user_from_claims(verifier.verify(token))
return
except Exception as exc:
logger.warning(
"OIDC token verification failed on %s: %s (expected aud=%s iss=%s)",
flask_request.path, exc, verifier.audiences, verifier.issuer,
)
if dev_mode():
g.current_user = _dev_user()
def require_auth(fn):
"""401 unless the request resolved to a user. app/__init__.py already
refuses unauthenticated /api calls; this states it on the route too."""
@wraps(fn)
def wrapper(*args, **kwargs):
if getattr(g, "current_user", None) is None:
return api_response(success=False, status=401, message="Authentication required", error_type="UNAUTHORIZED")
return fn(*args, **kwargs)
return wrapper
def anonymous_marketing_user():
"""Shared placeholder attributed to ping sessions from the public site.
The latency demo is open to unauthenticated visitors, but
PingToken.user_id is not nullable, so every anonymous ping needs some
user row to point at. One shared row rather than one per visitor: this
is not an account and must never accumulate into a user list.
"""
from app import db
from app.models import User
email = "anonymous@marketing.internal"
user = User.query.filter_by(email=email).first()
if user is None:
user = User(external_id="anonymous:marketing", email=email,
first_name="Anonymous", last_name="Visitor")
db.session.add(user)
db.session.commit()
return user
def get_request_user_id() -> Optional[str]:
if has_request_context() and getattr(g, "current_user", None):
return g.current_user.id
return None
+210
View File
@@ -0,0 +1,210 @@
from typing import Dict, List
from app import db
from app.core_client import core_request
from app.models import CloudflareCredential, CloudflareDNSRecord, CloudflareTunnel, Lease
from xcloudify_shared.cloudflare import CloudflareTunnelManager
from xcloudify_shared import logger
def _cf_manager_for_lease(lease: Lease) -> CloudflareTunnelManager | None:
cred = CloudflareCredential.query.filter_by(user_id=lease.customer_id).first()
if not cred:
logger.warning(
"Lease %s requested public exposure but customer %s has no Cloudflare "
"credentials on file (PUT /me/cloudflare) -- skipping",
lease.id, lease.customer_id,
)
return None
api_token, account_id, zone_id = cred.credentials()
return CloudflareTunnelManager(api_token, account_id, zone_id, logger), cred
def provision_exposure(lease_id: str, pod_id: str, exposures: List[Dict]) -> None:
"""exposures: [{"container_workload_id": str, "container_name": str, "internal_port": int}, ...]
for containers in this pod whose port had use_dns=true."""
if not exposures:
return
lease = Lease.query.get(lease_id)
if not lease:
logger.error("provision_exposure: lease %s not found", lease_id)
return
resolved = _cf_manager_for_lease(lease)
if not resolved:
return
cf_mgr, cred = resolved
tunnel = CloudflareTunnel.query.filter_by(pod_id=pod_id, deleted=False).first()
ingress_mappings = [
{
"dns_hostname": f"{e['container_name']}-{e['internal_port']}.{cred.domain}",
"local_ip": "127.0.0.1",
"local_port": str(e["internal_port"]),
}
for e in exposures
]
try:
cf_rsp = cf_mgr.setup_tunnel(
tunnel_name=tunnel.name if tunnel else f"lease-tun-{pod_id[:8]}",
ingress_mappings=ingress_mappings,
)
except Exception as exc:
logger.error("provision_exposure: Cloudflare setup_tunnel failed for pod %s: %s", pod_id, exc)
return
is_new_tunnel = tunnel is None
if is_new_tunnel:
tunnel = CloudflareTunnel(
lease_id=lease_id,
account_id=cf_mgr.account_id,
tunnel_id=cf_rsp["tunnel_id"],
name=cf_rsp["tunnel_name"],
tunnel_secret=cf_rsp.get("tunnel_secret") or "",
token=cf_rsp["token"],
pod_id=pod_id,
)
db.session.add(tunnel)
db.session.flush()
for e, mapping in zip(exposures, ingress_mappings):
hostname = mapping["dns_hostname"]
existing = CloudflareDNSRecord.query.filter_by(
tunnel_id=tunnel.id, container_workload_id=e["container_workload_id"],
internal_port=e["internal_port"], deleted=False,
).first()
if existing:
continue
dns_id = next(
(d["response"]["id"] for d in cf_rsp.get("dns_records_created", []) if d["hostname"] == hostname),
None,
)
if not dns_id:
found = cf_mgr.get_dns_record(hostname)
dns_id = found["id"] if found else None
if not dns_id:
logger.error("provision_exposure: no Cloudflare DNS record id for %s", hostname)
continue
db.session.add(CloudflareDNSRecord(
tunnel_id=tunnel.id,
zone_id=cf_mgr.zone_id,
dns_record_id=dns_id,
hostname=hostname,
content=f"{tunnel.tunnel_id}.cfargotunnel.com",
container_workload_id=e["container_workload_id"],
internal_port=e["internal_port"],
))
db.session.commit()
if is_new_tunnel:
_create_cloudflared_sidecar(pod_id=pod_id, lease_id=lease_id, tunnel_token=tunnel.token)
def _create_cloudflared_sidecar(*, pod_id: str, lease_id: str, tunnel_token: str) -> None:
"""Add the cloudflared sidecar to the pod via core's normal container API,
with the platform's own service key -- not the customer's lease gateway,
since this container is not something the gateway's allowlist or per-lease
network logic needs to apply to."""
resp = core_request("POST", "workloads/containers", json={
"pod": pod_id,
"tenant_id": lease_id,
"containers": [{
"docker_image": "cloudflare/cloudflared:latest",
"container_name": f"cloudflared-sidecar-{pod_id[:8]}",
"command": f"tunnel --no-autoupdate run --token {tunnel_token}",
"restart_policy": "always",
"cpu": 1,
"mem_limit": 64,
}],
})
if resp.status_code >= 300:
logger.error("Failed to create cloudflared sidecar for pod %s: %s %s", pod_id, resp.status_code, resp.text)
else:
logger.info("Created cloudflared sidecar for pod %s", pod_id)
def cleanup_exposure(container_workload_id: str) -> None:
"""Called after a container is deleted through the gateway. Removes any
DNS record for that container; if its tunnel has no records left,
deletes the tunnel and the cloudflared sidecar too."""
records = CloudflareDNSRecord.query.filter_by(
container_workload_id=container_workload_id, deleted=False
).all()
if not records:
return
tunnel_ids = {r.tunnel_id for r in records if r.tunnel_id}
for record in records:
tunnel = CloudflareTunnel.query.get(record.tunnel_id) if record.tunnel_id else None
if tunnel:
lease = Lease.query.get(tunnel.lease_id)
resolved = _cf_manager_for_lease(lease) if lease else None
if resolved:
cf_mgr, _ = resolved
try:
cf_mgr.delete_dns_record(record.dns_record_id)
except Exception as exc:
logger.error("cleanup_exposure: failed to delete DNS record %s: %s", record.hostname, exc)
record.deleted = True
db.session.add(record)
db.session.commit()
for tunnel_id in tunnel_ids:
tunnel = CloudflareTunnel.query.get(tunnel_id)
if not tunnel:
continue
remaining = CloudflareDNSRecord.query.filter_by(tunnel_id=tunnel_id, deleted=False).count()
if remaining:
continue
_teardown_tunnel(tunnel, delete_sidecar=True)
db.session.commit()
def cleanup_exposure_for_pod(pod_id: str) -> None:
"""Called after a whole pod is deleted through the gateway. core has
already deleted the pod's containers, including any cloudflared sidecar
-- this only tears down the Cloudflare-side tunnel/DNS and local rows."""
tunnels = CloudflareTunnel.query.filter_by(pod_id=pod_id, deleted=False).all()
for tunnel in tunnels:
records = CloudflareDNSRecord.query.filter_by(tunnel_id=tunnel.id, deleted=False).all()
lease = Lease.query.get(tunnel.lease_id)
resolved = _cf_manager_for_lease(lease) if lease else None
cf_mgr = resolved[0] if resolved else None
for record in records:
if cf_mgr:
try:
cf_mgr.delete_dns_record(record.dns_record_id)
except Exception as exc:
logger.error("cleanup_exposure_for_pod: failed to delete DNS record %s: %s", record.hostname, exc)
record.deleted = True
db.session.add(record)
_teardown_tunnel(tunnel, delete_sidecar=False)
db.session.commit()
def _teardown_tunnel(tunnel: CloudflareTunnel, *, delete_sidecar: bool) -> None:
lease = Lease.query.get(tunnel.lease_id)
resolved = _cf_manager_for_lease(lease) if lease else None
if resolved:
cf_mgr, _ = resolved
try:
cf_mgr.cleanup_tunnel(tunnel.name, delete_tunnel=True)
except Exception as exc:
logger.error("Failed to delete Cloudflare tunnel %s: %s", tunnel.name, exc)
if delete_sidecar and lease:
resp = core_request("GET", "workloads/containers", params={"tenant_id": lease.id})
if resp.ok:
for c in resp.json().get("data", []):
if c.get("name") == f"cloudflared-sidecar-{tunnel.pod_id[:8]}":
del_resp = core_request("DELETE", f"workloads/containers/{c['id']}")
if del_resp.status_code >= 300:
logger.error("Failed to delete cloudflared sidecar %s: %s", c["id"], del_resp.text)
break
tunnel.deleted = True
db.session.add(tunnel)
+78
View File
@@ -0,0 +1,78 @@
from xcloudify_shared import logger
from xcloudify_shared.core_client import app_core_request as core_request
def create_region(name: str, country: str, abbreviation: str) -> dict:
"""Create a private, failover-isolated region in core and return it.
`private_region` keeps it out of the public cloud's region picker, and
failover relocation is off because the hosts in a community region are
NAT'd islands under one person's desk -- core relocating a workload
between them is not the recovery it is for a rack of our own machines.
"""
resp = core_request("POST", "regions", json={
"name": name,
"country": country,
"abbreviation": abbreviation,
"private_region": True,
"relocate_on_host_failure": False,
"description": "Community-contributed hardware. Operator-owned and untrusted.",
})
resp.raise_for_status()
return resp.json()["data"]
def get_workload_host_status(core_workload_host_id: str) -> dict | None:
"""Live CPU/RAM utilization and the workloads actually running, for a
node's own dashboard and the operator's "everything" view.
Two core calls, not one: /utilization gives pooled totals (cheap,
aggregate), /active_workloads gives the list a count and a "what is this
node doing right now" table are both built from. Returns None on any
failure -- core being briefly unreachable should degrade a node card to
"usage unknown", not break the page it's on.
"""
try:
util_resp = core_request("GET", f"workload_hosts/{core_workload_host_id}/utilization")
util_resp.raise_for_status()
workloads_resp = core_request("GET", f"workload_hosts/{core_workload_host_id}/active_workloads")
workloads_resp.raise_for_status()
workloads = workloads_resp.json().get("data") or []
return {
"resources": util_resp.json().get("data", {}).get("resources", {}),
"workload_count": len(workloads),
}
except Exception as exc:
logger.warning("Could not fetch utilization for workload_host=%s: %s", core_workload_host_id, exc)
return None
def find_workload_host_by_physical_identifier(region_id: str, physical_identifier: str):
"""Look up the WorkloadHost a node's worker registered as, once it has
self-enrolled into core. Returns the host dict, or None if the worker
hasn't enrolled yet (or core is unreachable).
More than one can share this physical_identifier: core mints a brand new
WorkloadHost id on every enrollment, and this node's own worker discards
and re-enrolls whenever core stops recognising its stored credentials
(core/worker/main.py's credential self-heal). Each such recovery leaves
the previous row behind rather than replacing it, so "first match" would
keep returning the oldest, dead one forever -- which is exactly what
silently broke the reconciliation this function feeds the first time a
node recovered this way. Preferring an online row, and the most recently
created one among ties, is what makes this call actually track the
worker's current identity instead of its first.
"""
try:
resp = core_request("GET", "workload_hosts", params={"region_id": region_id})
resp.raise_for_status()
matches = [h for h in resp.json().get("data", []) if h.get("physical_identifier") == physical_identifier]
if not matches:
return None
online = [h for h in matches if h.get("status") == "online"]
pool = online or matches
return max(pool, key=lambda h: h.get("created_at") or "")
except Exception as exc:
logger.warning("Lookup of workload_host by physical_identifier=%s failed: %s",
physical_identifier, exc)
return None
+7
View File
@@ -0,0 +1,7 @@
from runtime_urls import COMMUNITY_SECRET_KEY
from xcloudify_shared.crypto import SecretBox
_box = SecretBox(COMMUNITY_SECRET_KEY, "COMMUNITY_SECRET_KEY")
encrypt_secret = _box.encrypt
decrypt_secret = _box.decrypt
+365
View File
@@ -0,0 +1,365 @@
import secrets
import uuid
from datetime import datetime
from sqlalchemy import (
Column, String, DateTime, ForeignKey, Boolean, Integer, UniqueConstraint,
)
from sqlalchemy.orm import relationship
from app import db
def _uuid() -> str:
return str(uuid.uuid4())
def _token(nbytes: int = 32) -> str:
return secrets.token_urlsafe(nbytes)
class User(db.Model):
"""Both providers (who lend hardware out) and customers (who rent it).
The same person can be both; the distinction is per-node, not per-user."""
__tablename__ = "users"
id = Column(String(36), primary_key=True, default=_uuid)
external_id = Column(String(255), unique=True, nullable=False)
email = Column(String(255), unique=True, nullable=False)
first_name = Column(String(255), nullable=True)
last_name = Column(String(255), nullable=True)
avatar_url = Column(String(512), nullable=True)
is_platform_admin = Column(Boolean, nullable=False, default=False)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
last_login_at = Column(DateTime, nullable=True)
nodes = relationship("Node", back_populates="owner", cascade="all, delete-orphan")
@property
def name(self) -> str:
full = f"{self.first_name or ''} {self.last_name or ''}".strip()
return full or self.email
def to_json(self) -> dict:
return {
"id": self.id,
"email": self.email,
"name": self.name,
"avatar_url": self.avatar_url,
"is_platform_admin": self.is_platform_admin,
}
class ProviderRegion(db.Model):
"""One core Region per (provider, country). See app/regions.py for why
this is not a single shared "Community" region.
Holds core's region enrollment key, which is a *shared* secret for
everyone enrolling into that region. Scoping a region to one provider is
what keeps a key read off one member's disk from being able to enroll
forged hosts into everybody else's capacity.
"""
__tablename__ = "provider_regions"
__table_args__ = (UniqueConstraint("owner_id", "country", name="uq_provider_region_owner_country"),)
id = Column(String(36), primary_key=True, default=_uuid)
owner_id = Column(String(36), ForeignKey("users.id"), nullable=False, index=True)
country = Column(String(2), nullable=False)
core_region_id = Column(String(36), nullable=False, index=True)
core_region_enrollment_key = Column(String(64), nullable=False)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
owner = relationship("User")
STATUS_PENDING = "pending"
STATUS_ONLINE = "online"
STATUS_OFFLINE = "offline"
STATUS_RETIRED = "retired"
TRUST_COMMUNITY = "community"
class Node(db.Model):
__tablename__ = "nodes"
id = Column(String(36), primary_key=True, default=_uuid)
owner_id = Column(String(36), ForeignKey("users.id"), nullable=False, index=True)
name = Column(String(255), nullable=False)
country = Column(String(2), nullable=False, index=True)
city = Column(String(255), nullable=True)
enrollment_key = Column(String(64), unique=True, nullable=True, index=True)
enrolled_at = Column(DateTime, nullable=True)
node_token = Column(String(64), unique=True, nullable=True, index=True)
status = Column(String(20), nullable=False, default=STATUS_PENDING, index=True)
trust_tier = Column(String(20), nullable=False, default=TRUST_COMMUNITY)
last_heartbeat_at = Column(DateTime, nullable=True)
cpu_cores = Column(Integer, nullable=True)
memory_mb = Column(Integer, nullable=True)
disk_gb = Column(Integer, nullable=True)
gpu_model = Column(String(255), nullable=True)
gpu_count = Column(Integer, nullable=True, default=0)
kernel = Column(String(255), nullable=True)
arch = Column(String(50), nullable=True)
provider_region_id = Column(String(36), ForeignKey("provider_regions.id"), nullable=True, index=True)
core_workload_host_id = Column(String(36), nullable=True, index=True)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
owner = relationship("User", back_populates="nodes")
provider_region = relationship("ProviderRegion")
leases = relationship("Lease", back_populates="node", cascade="all, delete-orphan")
@property
def is_leasable(self) -> bool:
if self.status != STATUS_ONLINE:
return False
return not any(l.is_active for l in self.leases)
def to_json(self, include_secrets: bool = False) -> dict:
data = {
"id": self.id,
"name": self.name,
"owner_id": self.owner_id,
"status": self.status,
"trust_tier": self.trust_tier,
"is_leasable": self.is_leasable,
"country": self.country,
"city": self.city,
"last_heartbeat_at": self.last_heartbeat_at.isoformat() if self.last_heartbeat_at else None,
"enrolled_at": self.enrolled_at.isoformat() if self.enrolled_at else None,
"specs": {
"cpu_cores": self.cpu_cores,
"memory_mb": self.memory_mb,
"disk_gb": self.disk_gb,
"gpu_model": self.gpu_model,
"gpu_count": self.gpu_count,
"kernel": self.kernel,
"arch": self.arch,
"self_reported": True,
},
"created_at": self.created_at.isoformat() if self.created_at else None,
}
if include_secrets:
data["enrollment_key"] = self.enrollment_key
data["core_workload_host_id"] = self.core_workload_host_id
return data
class Lease(db.Model):
"""A customer renting a node for a period."""
__tablename__ = "leases"
id = Column(String(36), primary_key=True, default=_uuid)
node_id = Column(String(36), ForeignKey("nodes.id"), nullable=False, index=True)
customer_id = Column(String(36), ForeignKey("users.id"), nullable=False, index=True)
started_at = Column(DateTime, default=datetime.utcnow, nullable=False)
ended_at = Column(DateTime, nullable=True)
node = relationship("Node", back_populates="leases")
customer = relationship("User")
@property
def is_active(self) -> bool:
return self.ended_at is None
def to_json(self) -> dict:
node = self.node
return {
"id": self.id,
"node_id": self.node_id,
"customer_id": self.customer_id,
"started_at": self.started_at.isoformat() if self.started_at else None,
"ended_at": self.ended_at.isoformat() if self.ended_at else None,
"active": self.is_active,
"node_name": node.name if node else None,
"node_status": node.status if node else None,
"node_country": node.country if node else None,
"node_city": node.city if node else None,
"region_id": (
node.provider_region.core_region_id
if node and node.provider_region else None
),
"node_enrolled": bool(node and node.core_workload_host_id),
}
class PingToken(db.Model):
"""Short-lived credential for one browser->worker ping session.
This service authorises the session and hands the browser a signaling
address; once the WebRTC data channel is up, neither this service nor the
signaling server is in the data path, so the measured RTT is
browser<->worker and not a round trip through our infrastructure.
"""
__tablename__ = "ping_tokens"
id = Column(String(36), primary_key=True, default=_uuid)
node_id = Column(String(36), ForeignKey("nodes.id"), nullable=False, index=True)
user_id = Column(String(36), ForeignKey("users.id"), nullable=False)
token = Column(String(64), unique=True, nullable=False, index=True, default=lambda: _token(24))
expires_at = Column(DateTime, nullable=False)
used_at = Column(DateTime, nullable=True)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
node = relationship("Node")
@property
def is_valid(self) -> bool:
return self.used_at is None and datetime.utcnow() < self.expires_at
def to_json(self) -> dict:
return {
"token": self.token,
"node_id": self.node_id,
"expires_at": self.expires_at.isoformat(),
}
class CloudflareCredential(db.Model):
"""One customer's own Cloudflare account, brought so their leased
workloads can be reached publicly without their node ever needing one.
Per-user rather than per-lease: community has no Project-like grouping
to hang this off (see cloud/app/models.py's Project, which is the model
this mirrors), and a customer's Cloudflare account is naturally a fact
about *them*, reused across however many machines they lease. The token
is verified against Cloudflare's API before being saved (see
routes/cloudflare_routes.py) and stored encrypted (app/crypto_utils.py).
"""
__tablename__ = "cloudflare_credentials"
id = Column(String(36), primary_key=True, default=_uuid)
user_id = Column(String(36), ForeignKey("users.id"), nullable=False, unique=True, index=True)
api_token_encrypted = Column(String(1024), nullable=False)
account_id = Column(String(255), nullable=False)
zone_id = Column(String(255), nullable=False)
domain = Column(String(255), nullable=False)
verified_at = Column(DateTime, nullable=False, default=datetime.utcnow)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
def credentials(self):
"""Decrypt and return (api_token, account_id, zone_id)."""
from app.crypto_utils import decrypt_secret
return decrypt_secret(self.api_token_encrypted), self.account_id, self.zone_id
def to_json(self) -> dict:
return {
"account_id": self.account_id,
"zone_id": self.zone_id,
"domain": self.domain,
"verified_at": self.verified_at.isoformat() if self.verified_at else None,
}
class CloudflareTunnel(db.Model):
"""A Cloudflare tunnel exposing one leased pod's ports publicly.
Mirrors cloud/app/models.py's CloudflareTunnel, keyed by lease_id instead
of vdc_id -- a lease is community's equivalent tenant boundary (the
gateway stamps tenant_id=lease.id on everything it creates), and the
customer's own credentials (above) authorise the Cloudflare side.
"""
__tablename__ = "cloudflare_tunnels"
id = Column(String(36), primary_key=True, default=_uuid)
lease_id = Column(String(36), ForeignKey("leases.id"), nullable=False, index=True)
name = Column(String(255), nullable=False)
account_id = Column(String(255), nullable=False)
tunnel_id = Column(String(255), nullable=False, unique=True)
tunnel_secret = Column(String(255), nullable=False)
token = Column(String(1024), nullable=False)
pod_id = Column(String(36), nullable=True)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
deleted = Column(Boolean, default=False, nullable=False)
deleted_at = Column(DateTime, nullable=True)
lease = relationship("Lease")
def to_json(self) -> dict:
return {
"id": self.id,
"lease_id": self.lease_id,
"name": self.name,
"tunnel_id": self.tunnel_id,
"pod_id": self.pod_id,
}
class CloudflareDNSRecord(db.Model):
"""A DNS record pointing at a CloudflareTunnel, for one exposed
container port. Mirrors cloud/app/models.py's CloudflareDNSRecord."""
__tablename__ = "cloudflare_dns_records"
id = Column(String(36), primary_key=True, default=_uuid)
tunnel_id = Column(String(36), ForeignKey("cloudflare_tunnels.id"), nullable=False, index=True)
zone_id = Column(String(255), nullable=False)
dns_record_id = Column(String(255), nullable=False, unique=True)
hostname = Column(String(255), nullable=False)
content = Column(String(255), nullable=False)
container_workload_id = Column(String(36), nullable=True, index=True)
internal_port = Column(Integer, nullable=True)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
deleted = Column(Boolean, default=False, nullable=False)
deleted_at = Column(DateTime, nullable=True)
tunnel = relationship("CloudflareTunnel", backref="dns_records")
def to_json(self) -> dict:
return {
"id": self.id,
"hostname": self.hostname,
"internal_port": self.internal_port,
"container_workload_id": self.container_workload_id,
}
class SSHKey(db.Model):
"""An account-level SSH public key, injected into VMs at launch.
Mirrors cloud/app/models.py's SSHKey exactly -- core has no key vault of
its own (core migration 0003_drop_ssh_keys), so every plane above it that
wants a saved-key picker keeps one itself. A VM create sends `public_keys`
as a list of ids; this gateway resolves them to raw OpenSSH text before
forwarding to core, the same way cloud's does (see gateway_routes.py's
`_resolve_ssh_keys`).
"""
__tablename__ = "ssh_keys"
id = Column(String(36), primary_key=True, default=_uuid)
user_id = Column(String(36), ForeignKey("users.id"), nullable=False, index=True)
key_name = Column(String(255), nullable=False)
public_key_data = Column(String(4096), nullable=False)
key_fingerprint = Column(String(255), nullable=True)
is_default = Column(Boolean, default=False, nullable=False)
created_at = Column(DateTime, default=datetime.utcnow, nullable=False)
updated_at = Column(DateTime, default=datetime.utcnow, onupdate=datetime.utcnow, nullable=False)
deleted = Column(Boolean, default=False, nullable=False)
deleted_at = Column(DateTime, nullable=True)
user = relationship("User")
def to_json(self) -> dict:
return {
"id": self.id,
"user_id": self.user_id,
"key_name": self.key_name,
"key_fingerprint": self.key_fingerprint,
"is_default": self.is_default,
"created_at": self.created_at.isoformat() if self.created_at else None,
"updated_at": self.updated_at.isoformat() if self.updated_at else None,
}
def to_json_with_key(self) -> dict:
result = self.to_json()
result["public_key_data"] = self.public_key_data
return result
+55
View File
@@ -0,0 +1,55 @@
import re
from flask import current_app
from app import db
from app.core_client import create_region
from app.models import ProviderRegion
from xcloudify_shared import logger
_COUNTRY_RE = re.compile(r"^[A-Za-z]{2}$")
def normalize_country(raw: str) -> str | None:
"""ISO-3166 alpha-2, upper-cased. Shape only -- we do not carry a country
table to validate membership, and a wrong-but-well-formed code is the
provider misdescribing their own hardware, which is already the trust
assumption everywhere else here."""
value = (raw or "").strip()
return value.upper() if _COUNTRY_RE.match(value) else None
def _abbreviation(owner_id: str, country: str) -> str:
"""core's Region.abbreviation is String(16); this fits in 12 and stays
unique per (owner, country) without a lookup."""
return f"C-{country}-{owner_id.replace('-', '')[:6]}"
def ensure_provider_region(user, country: str) -> ProviderRegion:
"""Return this provider's core region for `country`, creating it if needed.
Raises on failure rather than degrading: a node with no region cannot
enroll into core at all, and registering it anyway would leave the owner
with a node that silently never becomes schedulable.
"""
existing = ProviderRegion.query.filter_by(owner_id=user.id, country=country).first()
if existing:
return existing
prefix = current_app.config["CORE_REGION_NAME_PREFIX"]
region = create_region(
name=f"{prefix}/{user.email}/{country}",
country=country,
abbreviation=_abbreviation(user.id, country),
)
provider_region = ProviderRegion(
owner_id=user.id,
country=country,
core_region_id=region["id"],
core_region_enrollment_key=region["enrollment_key"],
)
db.session.add(provider_region)
db.session.commit()
logger.info("Created core region %s for provider %s in %s", region["id"], user.email, country)
return provider_region
+9
View File
@@ -0,0 +1,9 @@
from flask import Blueprint
api_bp = Blueprint("api", __name__, url_prefix="/api")
from app.routes import ( # noqa: E402,F401
auth_routes, node_routes, ping_routes, lease_routes, binary_routes,
gateway_routes, admin_routes, cloudflare_routes, catalog_routes,
provider_routes, ssh_key_routes,
)
+152
View File
@@ -0,0 +1,152 @@
from flask import g
from app.core_client import core_request
from app.models import Lease, Node, User
from app.routes import api_bp
from xcloudify_shared import api_response
from xcloudify_shared import logger
def _require_operator():
if not g.current_user.is_platform_admin:
return api_response(success=False, status=403,
message="Operator access required", error_type="FORBIDDEN")
return None
def _core_workloads(kind: str):
"""VMs or containers as core sees them. Returns [] when core is
unreachable rather than failing the whole page -- the node inventory below
is still worth showing."""
try:
resp = core_request("GET", f"workloads/{kind}")
resp.raise_for_status()
return resp.json().get("data") or []
except Exception as exc:
logger.warning("Could not list %s from core: %s", kind, exc)
return []
@api_bp.route("/admin/workloads", methods=["GET"])
def admin_workloads():
"""Every workload, attributed to the person running it and the machine
it landed on."""
denied = _require_operator()
if denied:
return denied
nodes_by_host = {n.core_workload_host_id: n for n in Node.query.all() if n.core_workload_host_id}
leases = {l.id: l for l in Lease.query.all()}
users = {u.id: u for u in User.query.all()}
rows = []
for kind in ("virtual_machines", "containers"):
for w in _core_workloads(kind):
node = nodes_by_host.get(w.get("workload_host_id"))
lease = leases.get(w.get("tenant_id"))
customer = users.get(lease.customer_id) if lease else None
owner = users.get(node.owner_id) if node else None
rows.append({
"workload_id": w.get("id"),
"name": w.get("name"),
"type": w.get("workload_type") or kind,
"status": w.get("status"),
"created_at": w.get("created_at"),
"host_id": w.get("workload_host_id"),
"node": {
"id": node.id,
"name": node.name,
"country": node.country,
"city": node.city,
"status": node.status,
"owner_email": owner.email if owner else None,
} if node else None,
"is_community_hardware": node is not None,
"lease_id": lease.id if lease else None,
"customer_email": customer.email if customer else None,
})
rows.sort(key=lambda r: (r.get("created_at") or ""), reverse=True)
return api_response(data=rows, meta={
"total": len(rows),
"on_community_hardware": sum(1 for r in rows if r["is_community_hardware"]),
})
@api_bp.route("/admin/users", methods=["GET"])
def admin_users():
"""Every registered person, provider or customer or both, with what they
actually hold -- nodes lent, leases taken. There is no separate signup:
an account exists the moment X-Dev-User names one (app/auth_utils.py), so
this is the only place a full roster is visible at all."""
denied = _require_operator()
if denied:
return denied
nodes_by_owner: dict[str, int] = {}
for n in Node.query.all():
nodes_by_owner[n.owner_id] = nodes_by_owner.get(n.owner_id, 0) + 1
active_leases_by_customer: dict[str, int] = {}
for l in Lease.query.all():
if l.is_active:
active_leases_by_customer[l.customer_id] = active_leases_by_customer.get(l.customer_id, 0) + 1
rows = [
{
"id": u.id,
"email": u.email,
"is_platform_admin": u.is_platform_admin,
"created_at": u.created_at.isoformat() if u.created_at else None,
"nodes_owned": nodes_by_owner.get(u.id, 0),
"active_leases": active_leases_by_customer.get(u.id, 0),
}
for u in User.query.order_by(User.created_at).all()
]
return api_response(data=rows, meta={"total": len(rows)})
@api_bp.route("/admin/leases", methods=["GET"])
def admin_leases():
"""Every lease ever taken, active or ended, attributed to customer and node."""
denied = _require_operator()
if denied:
return denied
users = {u.id: u for u in User.query.all()}
nodes = {n.id: n for n in Node.query.all()}
rows = []
for l in Lease.query.order_by(Lease.started_at.desc()).all():
node = nodes.get(l.node_id)
customer = users.get(l.customer_id)
rows.append({
"id": l.id,
"node_id": l.node_id,
"node_name": node.name if node else None,
"owner_email": users.get(node.owner_id).email if node and node.owner_id in users else None,
"customer_email": customer.email if customer else None,
"started_at": l.started_at.isoformat() if l.started_at else None,
"ended_at": l.ended_at.isoformat() if l.ended_at else None,
"active": l.is_active,
})
return api_response(data=rows, meta={"total": len(rows), "active": sum(1 for r in rows if r["active"])})
@api_bp.route("/admin/overview", methods=["GET"])
def admin_overview():
"""Counts an operator wants at a glance, without loading every table."""
denied = _require_operator()
if denied:
return denied
nodes = Node.query.all()
active_leases = [l for l in Lease.query.all() if l.is_active]
return api_response(data={
"providers": len({n.owner_id for n in nodes}),
"nodes_total": len(nodes),
"nodes_online": sum(1 for n in nodes if n.status == "online"),
"nodes_enrolled_into_core": sum(1 for n in nodes if n.core_workload_host_id),
"leases_active": len(active_leases),
"customers": len({l.customer_id for l in active_leases}),
})
+15
View File
@@ -0,0 +1,15 @@
from flask import current_app, g
from app.auth_utils import require_auth
from app.routes import api_bp
from xcloudify_shared import api_response
@api_bp.route("/auth/me", methods=["GET"])
@require_auth
def auth_me():
"""Who this request is, and whether identity is the dev header or a real
sign-in -- the portals only offer "switch identity" in dev mode."""
data = g.current_user.to_json()
data["auth_mode"] = "dev" if current_app.config.get("AUTH_DEV_BYPASS") else "oidc"
return api_response(data=data)
+76
View File
@@ -0,0 +1,76 @@
import hashlib
from pathlib import Path
from flask import Response, current_app, send_file
from app.routes import api_bp
from xcloudify_shared import api_response
DIST_DIR = Path("/opt/worker-dist")
BINARY_PATH = DIST_DIR / "xcloudify-worker"
RUNTIME_FILES = {
"Dockerfile": DIST_DIR / "Dockerfile",
"docker-compose.yml": DIST_DIR / "docker-compose.yml",
}
NSCONTROLLER_IMAGE_PATH = DIST_DIR / "nscontroller-image.tar.gz"
@api_bp.route("/nodes/worker-binary", methods=["GET"])
def worker_binary():
if not BINARY_PATH.is_file():
return api_response(
success=False, status=404,
message="No worker binary has been built yet. Run core/worker/build_binary.sh.",
error_type="NOT_FOUND",
)
return send_file(BINARY_PATH, as_attachment=True, download_name="xcloudify-worker")
@api_bp.route("/nodes/worker-runtime/<path:filename>", methods=["GET"])
def worker_runtime_file(filename):
"""Dockerfile / docker-compose.yml for the node's worker container.
Public for the same reason as the binary: the installer fetches these
before it holds any credential, and neither is a secret.
"""
path = RUNTIME_FILES.get(filename)
if path is None:
return api_response(success=False, status=404, message=f"Unknown runtime file: {filename}",
error_type="NOT_FOUND")
if not path.is_file():
return api_response(
success=False, status=404,
message=f"{filename} has not been published yet. Run core/worker/build_binary.sh.",
error_type="NOT_FOUND",
)
return Response(path.read_text(), content_type="text/plain")
@api_bp.route("/nodes/worker-binary/sha256", methods=["GET"])
def worker_binary_sha256():
if not BINARY_PATH.is_file():
return api_response(success=False, status=404, message="No worker binary built yet",
error_type="NOT_FOUND")
digest = hashlib.sha256(BINARY_PATH.read_bytes()).hexdigest()
return Response(digest + "\n", content_type="text/plain")
@api_bp.route("/nodes/nscontroller-image", methods=["GET"])
def nscontroller_image():
if not NSCONTROLLER_IMAGE_PATH.is_file():
return api_response(
success=False, status=404,
message="No NSController image has been built yet. Run core/worker/build_nscontroller_image.sh.",
error_type="NOT_FOUND",
)
return send_file(NSCONTROLLER_IMAGE_PATH, as_attachment=True, download_name="nscontroller-image.tar.gz")
@api_bp.route("/nodes/nscontroller-image/sha256", methods=["GET"])
def nscontroller_image_sha256():
if not NSCONTROLLER_IMAGE_PATH.is_file():
return api_response(success=False, status=404, message="No NSController image built yet",
error_type="NOT_FOUND")
digest = hashlib.sha256(NSCONTROLLER_IMAGE_PATH.read_bytes()).hexdigest()
return Response(digest + "\n", content_type="text/plain")
+59
View File
@@ -0,0 +1,59 @@
from flask import g
from app.core_client import core_request
from app.models import Lease
from app.routes import api_bp
from xcloudify_shared import api_response, logger
def _core_list(path: str, **params):
try:
resp = core_request("GET", path, params=params or None)
resp.raise_for_status()
return resp.json().get("data") or [], None
except Exception as exc:
logger.warning("Could not read %s from core: %s", path, exc)
return None, api_response(
success=False, status=502,
message="The compute plane is unreachable right now. Try again shortly.",
error_type="UPSTREAM_ERROR", error_details={"detail": str(exc)},
)
@api_bp.route("/images", methods=["GET"])
def list_images():
"""Boot images a leaseholder can build a VM from.
Unfiltered by lease on purpose: core's image catalog is global (there is
no per-region image table), and an image is a public artifact -- a URL and
a checksum -- not somebody's data.
"""
images, error = _core_list("images")
return error or api_response(data=images, meta={"total": len(images)})
@api_bp.route("/regions", methods=["GET"])
def list_regions():
"""Only the regions this user's own leases actually run in.
Core would happily list every region on the instance, including every
other provider's private one (`app/regions.py` gives each provider their
own per-country region). None of that is a leaseholder's business, and
none of it is usable by them either -- their workloads are pinned to the
one host they lease. So this is a filtered read, not a passthrough:
it answers "what region am I in", which is the only question the portal
asks regions for.
"""
mine = {
lease.node.provider_region.core_region_id
for lease in Lease.query.filter_by(customer_id=g.current_user.id).all()
if lease.is_active and lease.node and lease.node.provider_region
}
if not mine:
return api_response(data=[], meta={"total": 0})
regions, error = _core_list("regions")
if error:
return error
visible = [r for r in regions if r.get("id") in mine]
return api_response(data=visible, meta={"total": len(visible)})
+140
View File
@@ -0,0 +1,140 @@
from datetime import datetime
from flask import g, request
from app import db
from app.models import CloudflareCredential, CloudflareDNSRecord, CloudflareTunnel, Lease
from app.routes import api_bp
from xcloudify_shared import api_response, logger
@api_bp.route("/me/cloudflare", methods=["GET"])
def get_my_cloudflare():
cred = CloudflareCredential.query.filter_by(user_id=g.current_user.id).first()
return api_response(data=cred.to_json() if cred else None)
@api_bp.route("/me/cloudflare", methods=["PUT"])
def set_my_cloudflare():
"""Register or update this customer's Cloudflare account. The token is
verified against Cloudflare's API before being saved, and encrypted at
rest -- see app/crypto_utils.py."""
data = request.get_json(force=True) or {}
api_token = (data.get("api_token") or "").strip()
account_id = (data.get("account_id") or "").strip()
zone_id = (data.get("zone_id") or "").strip()
missing = [f for f, v in (("api_token", api_token), ("account_id", account_id), ("zone_id", zone_id)) if not v]
if missing:
return api_response(
success=False, status=400,
message=f"Missing required field(s): {', '.join(missing)}",
error_type="VALIDATION_ERROR",
)
from xcloudify_shared.cloudflare import CloudflareTunnelManager
cf_mgr = CloudflareTunnelManager(api_token, account_id, zone_id, logger)
try:
domain = cf_mgr.verify_credentials()
except Exception as exc:
return api_response(
success=False, status=400,
message=f"Could not verify Cloudflare credentials: {exc}",
error_type="CLOUDFLARE_VERIFICATION_FAILED",
)
try:
from app.crypto_utils import encrypt_secret
encrypted_token = encrypt_secret(api_token)
except Exception as exc:
logger.exception("Could not encrypt Cloudflare token for user %s", g.current_user.id)
return api_response(
success=False, status=500,
message=f"Server is not configured to store secrets: {exc}",
error_type="SECRET_STORAGE_UNAVAILABLE",
)
cred = CloudflareCredential.query.filter_by(user_id=g.current_user.id).first()
if cred is None:
cred = CloudflareCredential(user_id=g.current_user.id)
cred.api_token_encrypted = encrypted_token
cred.account_id = account_id
cred.zone_id = zone_id
cred.domain = domain
cred.verified_at = datetime.utcnow()
db.session.add(cred)
db.session.commit()
return api_response(data=cred.to_json(), message="Cloudflare credentials verified and saved")
@api_bp.route("/me/cloudflare", methods=["DELETE"])
def delete_my_cloudflare():
cred = CloudflareCredential.query.filter_by(user_id=g.current_user.id).first()
if cred is not None:
db.session.delete(cred)
db.session.commit()
return api_response(message="Cloudflare credentials removed")
@api_bp.route("/leases/<lease_id>/exposures", methods=["GET"])
def list_lease_exposures(lease_id):
"""The public hostnames for this lease's `use_dns` container ports.
The community counterpart of cloud's `/vdcs/<id>/exposures`, and the same
two halves for the same reason:
`cloudflare_domain` is the customer's own zone, and the hostname a port
will get is fully derived from it -- `{container_name}-{internal_port}.
{domain}` (app/cloudflare_exposure.py builds exactly that string). So a
caller holding the domain can name every URL without asking.
`exposures` is what was actually provisioned. Provisioning is best-effort
and deliberately non-fatal (gateway_routes._maybe_expose_containers
swallows Cloudflare failures rather than turning a working container into
an error), so a derived hostname with no row here is a port that asked to
be public and is not -- which the portal should say, instead of linking to
a URL nobody serves.
Built by hand, not from to_json(): CloudflareTunnel holds `token` and
`tunnel_secret`, which are the sidecar's credential for the tunnel and
have no business in a browser.
"""
lease = Lease.query.get(lease_id)
if lease is None:
return api_response(success=False, status=404, message="Lease not found", error_type="NOT_FOUND")
if lease.customer_id != g.current_user.id:
return api_response(success=False, status=403, message="Not your lease", error_type="FORBIDDEN")
cred = CloudflareCredential.query.filter_by(user_id=g.current_user.id).first()
records = (
CloudflareDNSRecord.query
.join(CloudflareTunnel, CloudflareDNSRecord.tunnel_id == CloudflareTunnel.id)
.filter(
CloudflareTunnel.lease_id == lease_id,
CloudflareDNSRecord.deleted == False, # noqa: E712
CloudflareTunnel.deleted == False, # noqa: E712
)
.all()
)
return api_response(data={
"cloudflare_domain": cred.domain if cred else None,
"exposures": [
{
"id": r.id,
"hostname": r.hostname,
"url": f"https://{r.hostname}",
"container_workload_id": r.container_workload_id,
"internal_port": r.internal_port,
"proxied": True,
"created_at": r.created_at.isoformat() if r.created_at else None,
"tunnel": {
"id": r.tunnel.id,
"name": r.tunnel.name,
"tunnel_id": r.tunnel.tunnel_id,
"pod_id": r.tunnel.pod_id,
} if r.tunnel else None,
}
for r in records
],
})
+220
View File
@@ -0,0 +1,220 @@
from flask import g, request
from app.core_client import core_request
from app.models import Lease, SSHKey, STATUS_ONLINE
from app.routes import api_bp
from xcloudify_shared import api_response
from xcloudify_shared import logger
from xcloudify_shared.gateway import WORKLOAD_ROUTES, deleted_workload, pending_dns_exposures
from xcloudify_shared.ssh_keys import resolve_ssh_keys
# A customer picks a boot image from core's catalog through their lease.
ROUTES = WORKLOAD_ROUTES.extend(exact={("GET", "images"): False})
def _active_lease(lease_id: str):
"""The caller's own active lease, or an error response."""
lease = Lease.query.get(lease_id)
if lease is None:
return None, api_response(success=False, status=404, message="Lease not found", error_type="NOT_FOUND")
if lease.customer_id != g.current_user.id:
return None, api_response(success=False, status=403, message="Not your lease", error_type="FORBIDDEN")
if not lease.is_active:
return None, api_response(success=False, status=409, message="Lease has ended", error_type="CONFLICT")
return lease, None
def _placement_for(lease):
"""Where this lease's workloads must run, or an error response.
Both values have to exist before anything is forwarded: without a core
region there is nothing to schedule into, and without a workload host id
the request would schedule anywhere in that region -- on hardware the
customer is not paying for.
"""
node = lease.node
if node.status != STATUS_ONLINE:
return None, api_response(success=False, status=409,
message="The leased node is not online",
error_type="CONFLICT")
if node.provider_region is None:
return None, api_response(success=False, status=409,
message="The leased node has no core region yet",
error_type="CONFLICT")
if not node.core_workload_host_id:
return None, api_response(
success=False, status=409,
message="The leased node has not finished enrolling into the compute plane yet. "
"It usually takes a minute after the installer runs.",
error_type="CONFLICT")
return {
"region_id": node.provider_region.core_region_id,
"worker_id": node.core_workload_host_id,
}, None
def _lease_subnet(lease_id: str) -> str:
"""A private /24 for this lease's network, stable across calls.
Derived from the lease id rather than fixed, so two leases on the same
node (sequential, one ended) do not fight over the same range if core
ever keeps both networks around. Each network lives on its own OVS
bridge/namespace, so cross-lease collision is not a correctness problem
either way -- this just avoids it being confusing to look at.
"""
import hashlib
h = hashlib.sha256(lease_id.encode()).digest()
return f"10.{h[0]}.{h[1]}.0/24"
def _ensure_lease_network(lease, region_id: str):
"""The network a lease's containers run on, creating it on first use.
A container workload cannot be created at all without a `networks` entry
(core: create_workload_container.py), but nothing about registering a
node or leasing it produces one -- a customer would otherwise have to
call POST networks themselves before their first container, a step nobody
asks for and the VM path has no equivalent of. So this makes "create a
container" self-contained the way "create a VM" already is: one lazily-
created network per lease, reused after.
"""
resp = core_request("GET", "networks", params={"tenant_id": lease.id})
if resp.ok:
existing = resp.json().get("data") or []
if existing:
return existing[0]["id"], None
cidr = _lease_subnet(lease.id)
base = cidr.rsplit(".", 1)[0]
create_resp = core_request("POST", "networks", json={
"name": f"lease-{lease.id[:8]}",
"region_id": region_id,
"tenant_id": lease.id,
"ipv4_cidr": cidr,
"ipv4_gateway": f"{base}.1",
"dhcp_range_start": f"{base}.10",
"dhcp_range_end": f"{base}.250",
})
if not create_resp.ok:
logger.error("Could not create a network for lease %s: %s", lease.id, create_resp.text)
return None, api_response(
success=False, status=502,
message="Could not prepare a network for this container. Try again shortly.",
error_type="UPSTREAM_ERROR",
)
return create_resp.json()["data"]["id"], None
def _maybe_expose_containers(lease, subpath: str, method: str, json_body, resp) -> None:
"""Hand any use_dns port off to Cloudflare tunnel/DNS provisioning --
synchronously, since community has no celery (see cloudflare_exposure.py)."""
pending = pending_dns_exposures(subpath, method, json_body, resp)
if pending is None:
return
pod_id, exposures = pending
from app.cloudflare_exposure import provision_exposure
provision_exposure(lease.id, pod_id, exposures)
def _maybe_cleanup_exposure(subpath: str, method: str) -> None:
"""Tear down any Cloudflare exposure that pointed at a deleted container/pod."""
deleted = deleted_workload(subpath, method)
if deleted is None:
return
kind, obj_id = deleted
from app.cloudflare_exposure import cleanup_exposure, cleanup_exposure_for_pod
(cleanup_exposure if kind == "container" else cleanup_exposure_for_pod)(obj_id)
def _owns_workload(lease, workload: dict) -> bool:
"""Whether a workload core returned belongs to this lease.
Core scopes nothing by lease -- it has never heard of one -- so a
single-object read or delete is checked here against the tenant id this
gateway stamps on everything it creates.
"""
return workload.get("tenant_id") == lease.id
@api_bp.route("/leases/<lease_id>/core/<path:subpath>", methods=["GET", "POST", "PUT", "DELETE"])
def lease_core_gateway(lease_id, subpath):
lease, error = _active_lease(lease_id)
if error:
return error
if ROUTES.needs_write(request.method, subpath) is None:
return api_response(success=False, status=403,
message=f"{request.method} {subpath} is not available through this gateway",
error_type="FORBIDDEN")
placement, error = _placement_for(lease)
if error:
return error
params = dict(request.args)
body = None
if request.method in ("POST", "PUT", "PATCH"):
body = request.get_json(silent=True) or {}
body["region_id"] = placement["region_id"]
body["tenant_id"] = lease.id
body["placement"] = {"worker_id": placement["worker_id"]}
if request.method == "POST" and subpath == "workloads/containers":
network_id, error = _ensure_lease_network(lease, placement["region_id"])
if error:
return error
for container in body.get("containers", []):
if not container.get("networks"):
container["networks"] = [network_id]
if request.method == "POST" and subpath == "workloads/virtual_machines":
network_id, error = _ensure_lease_network(lease, placement["region_id"])
if error:
return error
for vm in body.get("virtual-machines", []):
if not vm.get("networks"):
vm["networks"] = [{"id": network_id}]
resolve_ssh_keys(body, SSHKey, g.current_user.id)
else:
params["tenant_id"] = lease.id
resp = core_request(request.method, subpath, params=params, json=body)
if resp.ok and request.method == "GET" and subpath == "workload_hosts":
try:
payload = resp.json()
except ValueError:
payload = None
if isinstance(payload, dict) and isinstance(payload.get("data"), list):
leased_host_id = lease.node.core_workload_host_id
mine = [h for h in payload["data"] if h.get("id") == leased_host_id]
payload["data"] = mine
if isinstance(payload.get("meta"), dict):
payload["meta"]["total"] = len(mine)
return payload, resp.status_code
if (
resp.ok and request.method == "GET"
and subpath.startswith("workloads/") and subpath.count("/") == 2
):
try:
data = resp.json().get("data")
except ValueError:
data = None
if isinstance(data, dict) and data and not _owns_workload(lease, data):
logger.warning("Lease %s tried to reach workload %s it does not own", lease.id, subpath)
return api_response(success=False, status=404, message="Not found", error_type="NOT_FOUND")
if resp.ok:
try:
_maybe_expose_containers(lease, subpath, request.method, body, resp)
_maybe_cleanup_exposure(subpath, request.method)
except Exception:
logger.exception("Cloudflare exposure hook failed for %s %s (lease %s)", request.method, subpath, lease.id)
try:
return resp.json(), resp.status_code
except ValueError:
return api_response(success=False, status=502,
message="Compute plane returned an unreadable response",
error_type="UPSTREAM_ERROR")
+60
View File
@@ -0,0 +1,60 @@
from datetime import datetime
from flask import g, request
from app import db
from app.models import Lease, Node, STATUS_ONLINE
from app.routes import api_bp
from xcloudify_shared import api_response
@api_bp.route("/nodes/<node_id>/lease", methods=["POST"])
def lease_node(node_id):
node = Node.query.get(node_id)
if node is None:
return api_response(success=False, status=404, message="Node not found", error_type="NOT_FOUND")
if node.owner_id == g.current_user.id:
return api_response(success=False, status=400, message="You cannot lease your own node", error_type="VALIDATION_ERROR")
if node.status != STATUS_ONLINE:
return api_response(success=False, status=409, message="Node is not online", error_type="CONFLICT")
if any(l.is_active for l in node.leases):
return api_response(success=False, status=409, message="Node is already leased", error_type="CONFLICT")
lease = Lease(node_id=node.id, customer_id=g.current_user.id)
db.session.add(lease)
db.session.commit()
return api_response(data=lease.to_json(), status=201, message="Node leased")
@api_bp.route("/leases", methods=["GET"])
def list_leases():
"""Leases this user is party to, as customer or as node owner."""
as_customer = Lease.query.filter_by(customer_id=g.current_user.id).all()
as_owner = (
Lease.query.join(Node, Lease.node_id == Node.id)
.filter(Node.owner_id == g.current_user.id)
.all()
)
seen, combined = set(), []
for lease in as_customer + as_owner:
if lease.id not in seen:
seen.add(lease.id)
combined.append(lease)
return api_response(data=[l.to_json() for l in combined])
@api_bp.route("/leases/<lease_id>", methods=["DELETE"])
def end_lease(lease_id):
lease = Lease.query.get(lease_id)
if lease is None:
return api_response(success=False, status=404, message="Lease not found", error_type="NOT_FOUND")
if not lease.is_active:
return api_response(success=False, status=409, message="Lease already ended", error_type="CONFLICT")
if lease.customer_id != g.current_user.id and lease.node.owner_id != g.current_user.id:
return api_response(success=False, status=403, message="Not your lease", error_type="FORBIDDEN")
lease.ended_at = datetime.utcnow()
db.session.commit()
return api_response(data=lease.to_json(), message="Lease ended")
+299
View File
@@ -0,0 +1,299 @@
import secrets
from datetime import datetime
from pathlib import Path
from flask import Response, current_app, g, request
from app import db
from app.core_client import find_workload_host_by_physical_identifier, get_workload_host_status
from app.models import (
Node, STATUS_PENDING, STATUS_ONLINE, STATUS_OFFLINE, STATUS_RETIRED,
)
from app.regions import ensure_provider_region, normalize_country
from app.routes import api_bp
from app.signaling_client import connected_node_ids
from xcloudify_shared import api_response
from xcloudify_shared import logger
_TEMPLATE = Path(__file__).resolve().parents[2] / "templates" / "install.sh"
def _refresh_status(node: Node) -> None:
"""Mark a node offline once its heartbeat goes stale.
Derived on read rather than by a background sweeper: there is no other
reason to run a scheduler in this service yet, and a node's status only
matters at the moment someone looks at it or tries to lease it.
"""
if node.status not in (STATUS_ONLINE, STATUS_OFFLINE):
return
timeout = current_app.config["HEARTBEAT_TIMEOUT_SECONDS"]
if node.last_heartbeat_at is None:
node.status = STATUS_OFFLINE
return
stale = (datetime.utcnow() - node.last_heartbeat_at).total_seconds() > timeout
node.status = STATUS_OFFLINE if stale else STATUS_ONLINE
def _install_command(node: Node) -> str:
base = current_app.config["PUBLIC_BASE_URL"].rstrip("/")
return f"curl -fsSL {base}/api/nodes/install/{node.enrollment_key} | sudo bash"
def _serialize(node: Node, *, owned: bool, pingable: set[str] | None = None, with_status: bool = False) -> dict:
data = node.to_json(include_secrets=owned)
if pingable is not None:
data["pingable"] = node.id in pingable
if owned and node.status == STATUS_PENDING and node.enrollment_key:
data["install_command"] = _install_command(node)
if with_status and node.core_workload_host_id:
data["status_detail"] = get_workload_host_status(node.core_workload_host_id)
return data
@api_bp.route("/nodes", methods=["POST"])
def register_node():
"""Step 1 of enrollment: reserve the node and mint its single-use key.
Nothing is installed yet -- this only produces the curl|bash command the
owner runs on the machine itself. The owner's core region for this
country is created here too, so the install script can carry a region to
enroll into rather than the node discovering one later.
"""
data = request.get_json(force=True) or {}
name = (data.get("name") or "").strip()
if not name:
return api_response(success=False, status=400, message="'name' is required", error_type="VALIDATION_ERROR")
country = normalize_country(data.get("country"))
if not country:
return api_response(success=False, status=400,
message="'country' is required, as an ISO-3166 alpha-2 code (e.g. NP)",
error_type="VALIDATION_ERROR")
try:
provider_region = ensure_provider_region(g.current_user, country)
except Exception as exc:
logger.exception("Could not provision a core region for %s in %s", g.current_user.email, country)
return api_response(success=False, status=502,
message="Could not reach core to prepare a region for this node. Try again shortly.",
error_type="UPSTREAM_ERROR", error_details={"detail": str(exc)})
node = Node(
owner_id=g.current_user.id,
name=name,
country=country,
city=(data.get("city") or "").strip() or None,
enrollment_key=secrets.token_urlsafe(32),
status=STATUS_PENDING,
provider_region_id=provider_region.id,
)
db.session.add(node)
db.session.commit()
return api_response(data=_serialize(node, owned=True), status=201,
message="Node registered. Run the install command on the machine.")
@api_bp.route("/nodes", methods=["GET"])
def list_my_nodes():
nodes = Node.query.filter_by(owner_id=g.current_user.id).order_by(Node.created_at.desc()).all()
for n in nodes:
_refresh_status(n)
db.session.commit()
pingable = connected_node_ids()
return api_response(data=[_serialize(n, owned=True, pingable=pingable, with_status=True) for n in nodes])
@api_bp.route("/admin/nodes", methods=["GET"])
def list_all_nodes():
"""Operator view: every node and who owns it, not just the caller's own."""
if not g.current_user.is_platform_admin:
return api_response(success=False, status=403, message="Operator access required", error_type="FORBIDDEN")
nodes = Node.query.order_by(Node.created_at.desc()).all()
for n in nodes:
_refresh_status(n)
db.session.commit()
pingable = connected_node_ids()
return api_response(data=[
{**_serialize(n, owned=True, pingable=pingable, with_status=True), "owner_email": n.owner.email}
for n in nodes
])
@api_bp.route("/nodes/available", methods=["GET"])
def list_available_nodes():
"""The public listing: online, unleased nodes. Open to anonymous visitors,
which is what lets the marketing site's latency demo be real rather than a
mock -- so it must never include an owner's credentials."""
nodes = Node.query.filter(Node.status.in_((STATUS_ONLINE, STATUS_OFFLINE))).all()
for n in nodes:
_refresh_status(n)
db.session.commit()
pingable = connected_node_ids()
out = [
_serialize(n, owned=False, pingable=pingable)
for n in nodes
if n.is_leasable
]
country = normalize_country(request.args.get("country"))
if country:
out = [n for n in out if n["country"] == country]
return api_response(data=out)
@api_bp.route("/nodes/<node_id>", methods=["GET"])
def get_node(node_id):
node = Node.query.get(node_id)
if node is None:
return api_response(success=False, status=404, message="Node not found", error_type="NOT_FOUND")
_refresh_status(node)
db.session.commit()
caller = getattr(g, "current_user", None)
owned = caller is not None and node.owner_id == caller.id
return api_response(data=_serialize(node, owned=owned, pingable=connected_node_ids()))
@api_bp.route("/nodes/<node_id>", methods=["DELETE"])
def retire_node(node_id):
node = Node.query.get(node_id)
if node is None:
return api_response(success=False, status=404, message="Node not found", error_type="NOT_FOUND")
if node.owner_id != g.current_user.id:
return api_response(success=False, status=403, message="Not your node", error_type="FORBIDDEN")
if any(l.is_active for l in node.leases):
return api_response(success=False, status=409, message="Node has an active lease; end it first", error_type="CONFLICT")
node.status = STATUS_RETIRED
node.node_token = None
node.enrollment_key = None
db.session.commit()
return api_response(message="Node retired")
@api_bp.route("/nodes/install/<enrollment_key>", methods=["GET"])
def install_script(enrollment_key):
"""Served unauthenticated: the machine running it has only the key.
The key alone is the credential, which is why it is single-use and why
the script's job is to trade it for a per-node token immediately.
"""
node = Node.query.filter_by(enrollment_key=enrollment_key).first()
if node is None:
return Response("#!/bin/sh\necho 'Invalid or already-used enrollment key' >&2\nexit 1\n",
status=404, content_type="text/x-shellscript")
region = node.provider_region
script = _TEMPLATE.read_text()
for placeholder, value in {
"{{ base_url }}": current_app.config["PUBLIC_BASE_URL"].rstrip("/"),
"{{ signaling_url }}": current_app.config["SIGNALING_PUBLIC_URL"].rstrip("/"),
"{{ stun_urls }}": current_app.config["STUN_URLS"],
"{{ enrollment_key }}": enrollment_key,
"{{ worker_binary_url }}": current_app.config.get("WORKER_BINARY_URL") or "",
"{{ worker_binary_sha256 }}": current_app.config.get("WORKER_BINARY_SHA256") or "",
"{{ nscontroller_image_url }}": current_app.config.get("NSCONTROLLER_IMAGE_URL") or "",
"{{ nscontroller_image_sha256 }}": current_app.config.get("NSCONTROLLER_IMAGE_SHA256") or "",
"{{ core_public_url }}": current_app.config.get("CORE_PUBLIC_BASE_URL") or "",
"{{ core_ws_url }}": current_app.config.get("CORE_PUBLIC_WS_URL") or "",
"{{ core_region_id }}": (region.core_region_id if region else ""),
"{{ core_region_enrollment_key }}": (region.core_region_enrollment_key if region else ""),
}.items():
script = script.replace(placeholder, value)
return Response(script, content_type="text/x-shellscript")
@api_bp.route("/nodes/enroll", methods=["POST"])
def enroll_node():
"""Step 2: the machine trades its single-use key for a per-node token."""
data = request.get_json(force=True) or {}
key = (data.get("enrollment_key") or "").strip()
if not key:
return api_response(success=False, status=400, message="'enrollment_key' is required", error_type="VALIDATION_ERROR")
node = Node.query.filter_by(enrollment_key=key).first()
if node is None:
return api_response(success=False, status=404, message="Invalid or already-used enrollment key", error_type="NOT_FOUND")
def _int(field):
try:
return int(data.get(field) or 0)
except (TypeError, ValueError):
return 0
node.cpu_cores = _int("cpu_cores")
node.memory_mb = _int("memory_mb")
node.disk_gb = _int("disk_gb")
node.gpu_count = _int("gpu_count")
node.gpu_model = (data.get("gpu_model") or "").strip() or None
node.kernel = (data.get("kernel") or "").strip() or None
node.arch = (data.get("arch") or "").strip() or None
node.node_token = secrets.token_urlsafe(32)
node.enrollment_key = None
node.enrolled_at = datetime.utcnow()
node.last_heartbeat_at = datetime.utcnow()
node.status = STATUS_ONLINE
db.session.commit()
logger.info("Node %s enrolled (%s cores, %s MB)", node.id, node.cpu_cores, node.memory_mb)
return api_response(
data={"node_id": node.id, "node_token": node.node_token},
status=201,
message="Enrolled",
)
def _node_from_token() -> Node | None:
token = request.headers.get("X-Node-Token") or (request.get_json(silent=True) or {}).get("node_token")
if not token:
return None
return Node.query.filter_by(node_token=token).first()
def _reconcile_core_workload_host(node: Node) -> None:
"""Learn the WorkloadHost id core assigned this node's worker.
The worker self-enrolls into core directly (it already does this for
core-owned hosts, see core/worker/main.py::enroll_worker) using this
node's own id as PHYSICAL_IDENTIFIER. So rather than this service calling
core's enroll endpoint on the node's behalf, it looks for the host core
already created and remembers its id -- lazily, on heartbeat, since that
is already a periodic touchpoint and self-heals if core was unreachable
on an earlier attempt.
Re-checked even once already set, not just the first time: core mints a
brand new WorkloadHost id on every enrollment, so a worker whose stored
WORKER_ID/SECRET went stale and re-enrolled (main.py's own credential
self-heal) comes back under a *different* id core has never told this
node about. Left as "resolve once and stop," that id would sit stale
forever, silently breaking the lease gateway for a node that is actually
healthy. Cheap to redo -- one lookup by this node's own physical
identifier -- so there is no reason to gate it on being unset.
"""
if node.provider_region is None:
return
host = find_workload_host_by_physical_identifier(node.provider_region.core_region_id, node.id)
if host and host["id"] != node.core_workload_host_id:
logger.info(
"Node %s re-matched to core workload_host %s (was %s)",
node.id, host["id"], node.core_workload_host_id,
)
node.core_workload_host_id = host["id"]
@api_bp.route("/nodes/heartbeat", methods=["POST"])
def heartbeat():
"""Authenticated by the node's own token, not a user session."""
node = _node_from_token()
if node is None:
return api_response(success=False, status=401, message="Invalid node token", error_type="UNAUTHORIZED")
if node.status == STATUS_RETIRED:
return api_response(success=False, status=410, message="Node retired", error_type="GONE")
node.last_heartbeat_at = datetime.utcnow()
node.status = STATUS_ONLINE
_reconcile_core_workload_host(node)
db.session.commit()
return api_response(data={"node_id": node.id}, message="ok")
+81
View File
@@ -0,0 +1,81 @@
from datetime import datetime, timedelta
from flask import current_app, g, request
from app import db
from app.models import Node, PingToken, STATUS_ONLINE
from app.routes import api_bp
from app.signaling_client import connected_node_ids
from xcloudify_shared import api_response
@api_bp.route("/nodes/<node_id>/ping-session", methods=["POST"])
def create_ping_session(node_id):
"""Issue a short-lived token for one browser->worker ping session."""
node = Node.query.get(node_id)
if node is None:
return api_response(success=False, status=404, message="Node not found", error_type="NOT_FOUND")
if node.status != STATUS_ONLINE:
return api_response(success=False, status=409, message="Node is not online", error_type="CONFLICT")
if node.id not in connected_node_ids():
return api_response(
success=False, status=409,
message="Node's worker is not connected to the signaling server",
error_type="CONFLICT",
)
if getattr(g, "current_user", None) is not None:
user_id = g.current_user.id
else:
from app.auth_utils import anonymous_marketing_user
user_id = anonymous_marketing_user().id
ttl = current_app.config["PING_TOKEN_TTL_SECONDS"]
token = PingToken(
node_id=node.id,
user_id=user_id,
expires_at=datetime.utcnow() + timedelta(seconds=ttl),
)
db.session.add(token)
db.session.commit()
data = token.to_json()
data["signaling_url"] = current_app.config["SIGNALING_BROWSER_URL"].rstrip("/")
data["worker_id"] = node.id
data["ice_servers"] = [
{"urls": url.strip()}
for url in (current_app.config["STUN_URLS"] or "").split(",") if url.strip()
]
data["path"] = "webrtc"
return api_response(data=data, status=201, message="Ping session authorized")
@api_bp.route("/nodes/ping-auth", methods=["POST"])
def ping_auth():
"""Called by the worker to validate a token a browser just presented.
Not authenticated by user session on purpose -- the caller is the worker,
and the token itself is the credential being checked. The worker calls
this before answering the offer, so an unauthorised visitor never gets a
peer connection at all.
"""
data = request.get_json(force=True) or {}
token_value = (data.get("token") or "").strip()
node_token = (data.get("node_token") or request.headers.get("X-Node-Token") or "").strip()
if not token_value or not node_token:
return api_response(success=False, status=400, message="'token' and node token are required", error_type="VALIDATION_ERROR")
node = Node.query.filter_by(node_token=node_token).first()
if node is None:
return api_response(success=False, status=401, message="Invalid node token", error_type="UNAUTHORIZED")
ping_token = PingToken.query.filter_by(token=token_value).first()
if ping_token is None or not ping_token.is_valid:
return api_response(data={"allowed": False}, message="Token invalid or expired")
if ping_token.node_id != node.id:
return api_response(data={"allowed": False}, message="Token is not for this node")
ping_token.used_at = datetime.utcnow()
db.session.commit()
return api_response(data={"allowed": True, "user_id": ping_token.user_id}, message="ok")
+151
View File
@@ -0,0 +1,151 @@
from datetime import datetime
from flask import g
from app.core_client import core_request, get_workload_host_status
from app.models import Lease, Node, User, STATUS_ONLINE, STATUS_PENDING
from app.routes import api_bp
from xcloudify_shared import api_response, logger
def _hours(lease: Lease) -> float:
"""How long this lease has run, in hours, still counting if it is live."""
end = lease.ended_at or datetime.utcnow()
return max(0.0, (end - lease.started_at).total_seconds() / 3600.0)
def _lease_row(lease: Lease, customers: dict) -> dict:
customer = customers.get(lease.customer_id)
return {
"id": lease.id,
"node_id": lease.node_id,
"customer_email": customer.email if customer else None,
"started_at": lease.started_at.isoformat() if lease.started_at else None,
"ended_at": lease.ended_at.isoformat() if lease.ended_at else None,
"active": lease.is_active,
"hours": round(_hours(lease), 2),
}
@api_bp.route("/provider/summary", methods=["GET"])
def provider_summary():
"""The provider dashboard in one call.
Deliberately one round trip rather than letting the page assemble it from
/nodes plus a lease call per node: a provider with a rack of machines
would otherwise open N+1 connections to draw a summary card.
"""
nodes = Node.query.filter_by(owner_id=g.current_user.id).all()
node_ids = {n.id for n in nodes}
leases = (
Lease.query.filter(Lease.node_id.in_(node_ids)).order_by(Lease.started_at.desc()).all()
if node_ids else []
)
customers = {u.id: u for u in User.query.filter(
User.id.in_({l.customer_id for l in leases})
).all()} if leases else {}
active = [l for l in leases if l.is_active]
return api_response(data={
"nodes_total": len(nodes),
"nodes_online": sum(1 for n in nodes if n.status == STATUS_ONLINE),
"nodes_pending": sum(1 for n in nodes if n.status == STATUS_PENDING),
"nodes_leased": len({l.node_id for l in active}),
"leases_active": len(active),
"leases_total": len(leases),
"hours_served_total": round(sum(_hours(l) for l in leases), 2),
"hours_served_active": round(sum(_hours(l) for l in active), 2),
"customers": len({l.customer_id for l in leases}),
"recent_leases": [_lease_row(l, customers) for l in leases[:10]],
})
@api_bp.route("/nodes/<node_id>/leases", methods=["GET"])
def node_leases(node_id):
"""Every lease ever taken on one of my machines.
Owner-only. A customer holding the current lease has no business seeing
who rented the box before them, and a passer-by none at all.
"""
node = Node.query.get(node_id)
if node is None:
return api_response(success=False, status=404, message="Node not found", error_type="NOT_FOUND")
if node.owner_id != g.current_user.id:
return api_response(success=False, status=403, message="Not your node", error_type="FORBIDDEN")
leases = sorted(node.leases, key=lambda l: l.started_at or datetime.min, reverse=True)
customers = {u.id: u for u in User.query.filter(
User.id.in_({l.customer_id for l in leases})
).all()} if leases else {}
return api_response(
data=[_lease_row(l, customers) for l in leases],
meta={
"total": len(leases),
"active": sum(1 for l in leases if l.is_active),
"hours_served_total": round(sum(_hours(l) for l in leases), 2),
},
)
@api_bp.route("/nodes/<node_id>/workloads", methods=["GET"])
def node_workloads(node_id):
"""What is running on my machine right now.
Deliberately redacted. A provider is entitled to know their hardware is
running two containers and a VM, how big those are, and which customer
they belong to -- that is what running a machine responsibly requires.
They are *not* entitled to `launch_params`: the docker image, the
environment, the cloud-init metadata, the SSH keys. That asymmetry is the
whole trust story in community/app/models.py's docstring -- the host
operator is untrusted, and owning the disk a workload sits on is not the
same as being handed its contents by the control plane.
"""
node = Node.query.get(node_id)
if node is None:
return api_response(success=False, status=404, message="Node not found", error_type="NOT_FOUND")
if node.owner_id != g.current_user.id:
return api_response(success=False, status=403, message="Not your node", error_type="FORBIDDEN")
if not node.core_workload_host_id:
return api_response(data=[], meta={"total": 0, "enrolled": False})
try:
resp = core_request("GET", f"workload_hosts/{node.core_workload_host_id}/active_workloads")
resp.raise_for_status()
workloads = resp.json().get("data") or []
except Exception as exc:
logger.warning("Could not list workloads on node %s: %s", node.id, exc)
return api_response(
success=False, status=502,
message="The compute plane is unreachable right now. Try again shortly.",
error_type="UPSTREAM_ERROR", error_details={"detail": str(exc)},
)
leases = {l.id: l for l in Lease.query.filter_by(node_id=node.id).all()}
customers = {u.id: u for u in User.query.filter(
User.id.in_({l.customer_id for l in leases.values()})
).all()} if leases else {}
rows = []
for w in workloads:
lease = leases.get(w.get("tenant_id"))
customer = customers.get(lease.customer_id) if lease else None
rows.append({
"id": w.get("id"),
"name": w.get("name"),
"type": w.get("workload_type"),
"status": w.get("status"),
"created_at": w.get("created_at"),
"is_sidecar": bool(w.get("is_sidecar")),
"resource_usage": w.get("resource_usage"),
"lease_id": lease.id if lease else None,
"customer_email": customer.email if customer else None,
})
status = get_workload_host_status(node.core_workload_host_id)
return api_response(data=rows, meta={
"total": len(rows),
"enrolled": True,
"resources": (status or {}).get("resources", {}),
})
+15
View File
@@ -0,0 +1,15 @@
from flask import g
from app import db
from app.auth_utils import require_auth
from app.models import SSHKey
from app.routes import api_bp
from xcloudify_shared.ssh_keys import register_ssh_key_routes
register_ssh_key_routes(
api_bp,
db=db,
SSHKey=SSHKey,
current_user_id=lambda: g.current_user.id,
decorators=[require_auth],
)
+26
View File
@@ -0,0 +1,26 @@
import requests
from flask import current_app
from xcloudify_shared import logger
_TIMEOUT = 3
def connected_node_ids() -> set[str]:
"""Node ids with a live worker signaling session. Empty set if the
signaling server is unreachable -- callers degrade to "not pingable",
which is the truthful answer when we cannot set up a session anyway."""
base = (current_app.config.get("SIGNALING_INTERNAL_URL") or "").rstrip("/")
if not base:
return set()
try:
resp = requests.get(
f"{base}/internal/workers",
headers={"X-Signaling-Secret": current_app.config.get("SIGNALING_SHARED_SECRET") or ""},
timeout=_TIMEOUT,
)
resp.raise_for_status()
return set(resp.json().get("workers", []))
except Exception as exc:
logger.warning("Could not reach the signaling server at %s: %s", base, exc)
return set()
+115
View File
@@ -0,0 +1,115 @@
name: core
include:
- path: ../core/docker-compose.yml
services:
community-db:
image: mariadb:10.11
container_name: xcloudify-community-db
environment:
MYSQL_ROOT_PASSWORD: password
MYSQL_DATABASE: community
MYSQL_USER: community_user
MYSQL_PASSWORD: community_password
ports:
- "3308:3306"
volumes:
- community_db_data:/var/lib/mysql
healthcheck:
test: ["CMD", "healthcheck.sh", "--connect", "--innodb_initialized"]
interval: 10s
timeout: 5s
retries: 5
community-db-migrate:
build:
context: ..
dockerfile: community/Dockerfile
image: xcloudify-community
env_file: .env
volumes:
- .:/app
- ../packages/pyshared:/packages/pyshared
working_dir: /app
command: python manage.py migrations:apply
depends_on:
community-db:
condition: service_healthy
community-signaling:
build:
context: ..
dockerfile: community/signaling/Dockerfile
image: xcloudify-community-signaling
container_name: xcloudify-community-signaling
env_file: .env
ports:
- "5003:5003"
healthcheck:
test: ["CMD", "python", "-c", "import urllib.request;urllib.request.urlopen('http://localhost:5003/healthz')"]
interval: 10s
timeout: 5s
retries: 5
community-api:
build:
context: ..
dockerfile: community/Dockerfile
image: xcloudify-community
container_name: xcloudify-community-api
env_file: .env
environment:
FLASK_APP: api_server.py
volumes:
- .:/app
- ../packages/pyshared:/packages/pyshared
- ../core/worker/dist:/opt/worker-dist:ro
working_dir: /app
command: >
sh -c "flask run --host=0.0.0.0 --port=5002 --debug"
ports:
- "5002:5002"
depends_on:
community-db:
condition: service_healthy
community-db-migrate:
condition: service_completed_successfully
api-server:
condition: service_healthy
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:5002/api/healthz"]
interval: 10s
timeout: 5s
retries: 5
# Real sign-in outside development: `docker compose --profile auth up`.
# Fronts community-api's /api and the customer portal on one origin, so
# the portal is built with VITE_COMMUNITY_API_BASE_URL=/api and the session
# cookie rides along. Uses cloud's oauth2-proxy.cfg unchanged.
community-oauth2-proxy:
image: quay.io/oauth2-proxy/oauth2-proxy:v7.6.0
container_name: xcloudify-community-oauth2-proxy
profiles: ["auth"]
command: ["--config=/oauth2-proxy.cfg", "--http-address=0.0.0.0:8091"]
volumes:
- ../cloud/oauth2-proxy.cfg:/oauth2-proxy.cfg:ro
environment:
OAUTH2_PROXY_OIDC_ISSUER_URL: ${OIDC_ISSUER:-}
OAUTH2_PROXY_CLIENT_ID: ${OAUTH2_PROXY_CLIENT_ID:-}
OAUTH2_PROXY_CLIENT_SECRET: ${OAUTH2_PROXY_CLIENT_SECRET:-}
OAUTH2_PROXY_COOKIE_SECRET: ${OAUTH2_PROXY_COOKIE_SECRET:-}
OAUTH2_PROXY_REDIRECT_URL: ${OAUTH2_PROXY_REDIRECT_URL:-http://localhost:8091/oauth2/callback}
OAUTH2_PROXY_COOKIE_SECURE: ${OAUTH2_PROXY_COOKIE_SECURE:-false}
OAUTH2_PROXY_UPSTREAMS: ${OAUTH2_PROXY_UPSTREAMS:-http://community-api:5002/api/,http://host.docker.internal:8081/}
extra_hosts:
- "host.docker.internal:host-gateway"
ports:
- "8091:8091"
depends_on:
community-api:
condition: service_healthy
volumes:
community_db_data:
+220
View File
@@ -0,0 +1,220 @@
#!/usr/bin/env python3
import logging
import re
import click
from flask.cli import with_appcontext
from flask_migrate import (
init as alembic_init,
migrate as alembic_migrate,
upgrade as alembic_upgrade,
downgrade as alembic_downgrade,
current as alembic_current,
history as alembic_history,
stamp as alembic_stamp,
)
from app import app, db # noqa: F401
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
@click.group()
def command_line_interface():
pass
@command_line_interface.command("migrations:init")
@with_appcontext
def init_migrations():
from pathlib import Path
if Path("migrations").exists():
logger.info("migrations/ already exists")
else:
alembic_init()
alembic_stamp(revision="head")
@command_line_interface.command("migrations:generate")
@click.option("--message", "-m", default="schema update")
@with_appcontext
def generate_migration(message):
alembic_migrate(message=message)
@command_line_interface.command("migrations:apply")
@with_appcontext
def apply_migrations():
alembic_upgrade()
@command_line_interface.command("migrations:current")
@with_appcontext
def current_migration():
alembic_current(verbose=True)
@command_line_interface.command("migrations:history")
@with_appcontext
def migration_history():
alembic_history(verbose=True)
@command_line_interface.command("migrations:downgrade")
@click.option("--revision", "-r", required=True)
@with_appcontext
def downgrade_migration(revision):
alembic_downgrade(revision=revision)
def _set_admin(email: str, value: bool):
from app.models import User
email = email.strip().lower()
user = User.query.filter_by(email=email).first()
if user is None:
raise click.ClickException(
f"No user with email {email}. They must log in once before they can be "
f"granted operator access -- accounts are created from the IdP identity, "
f"not by this command."
)
user.is_platform_admin = value
db.session.commit()
logger.info("%s is_platform_admin=%s", email, value)
@command_line_interface.command("admin:grant")
@click.argument("email")
@with_appcontext
def grant_admin(email):
"""Give an existing user platform-operator access."""
_set_admin(email, True)
@command_line_interface.command("admin:revoke")
@click.argument("email")
@with_appcontext
def revoke_admin(email):
"""Remove platform-operator access."""
_set_admin(email, False)
@command_line_interface.command("admin:list")
@with_appcontext
def list_admins():
"""Show who currently holds operator access."""
from app.models import User
admins = User.query.filter_by(is_platform_admin=True).all()
if not admins:
click.echo("No platform operators. Grant one with: manage.py admin:grant <email>")
return
for u in admins:
click.echo(f"{u.email}\t{u.name}")
DEFAULT_IMAGES = [
{
"name": "Ubuntu 24.04 LTS (Noble)",
"location": "https://cloud-images.ubuntu.com/noble/current/noble-server-cloudimg-amd64.img",
"location_type": "http",
"checksum": "d0fe84bb5f80853425fa6be28e2c106f30104c3cfe8611933f2e65c9b63f0e30",
"os_family": "ubuntu",
"os_version": "24.04",
"format": "qcow2",
"size": 3.5,
},
{
"name": "Ubuntu 22.04 LTS (Jammy)",
"location": "https://cloud-images.ubuntu.com/jammy/current/jammy-server-cloudimg-amd64.img",
"location_type": "http",
"checksum": "",
"os_family": "ubuntu",
"os_version": "22.04",
"format": "qcow2",
"size": 3.5,
},
]
_CHECKSUM_SOURCES = {
"https://cloud-images.ubuntu.com/noble/current/noble-server-cloudimg-amd64.img":
("https://cloud-images.ubuntu.com/noble/current/SHA256SUMS", "noble-server-cloudimg-amd64.img"),
"https://cloud-images.ubuntu.com/jammy/current/jammy-server-cloudimg-amd64.img":
("https://cloud-images.ubuntu.com/jammy/current/SHA256SUMS", "jammy-server-cloudimg-amd64.img"),
}
_SHA256_RE = re.compile(r"[0-9a-f]{64}", re.IGNORECASE)
def _resolve_checksum(location: str) -> str | None:
"""Read the publisher's own checksum for an image, or None."""
import requests
source = _CHECKSUM_SOURCES.get(location)
if not source:
return None
url, filename = source
try:
resp = requests.get(url, timeout=30)
resp.raise_for_status()
except Exception as exc:
logger.warning("Could not fetch %s: %s", url, exc)
return None
for line in resp.text.splitlines():
parts = line.split()
if len(parts) == 2 and parts[1].lstrip("*") == filename:
digest = parts[0]
if not _SHA256_RE.fullmatch(digest):
logger.error("%s lists a non-sha256 digest for %s; the worker cannot verify it",
url, filename)
return None
return digest
logger.warning("%s not listed in %s", filename, url)
return None
@command_line_interface.command("images:seed")
@click.option("--force", is_flag=True, help="Re-create images that already exist by name.")
@with_appcontext
def seed_images(force):
"""Put a usable boot-image catalog in core, if it has none.
Idempotent by image name: run it on every deploy. Community writes to
core with the platform service key it already holds for region creation,
so this needs no extra credential.
"""
from app.core_client import core_request
resp = core_request("GET", "images")
if not resp.ok:
raise click.ClickException(f"Could not read core's image catalog: HTTP {resp.status_code}")
existing = {i.get("name"): i for i in (resp.json().get("data") or [])}
created, skipped, failed = 0, 0, 0
for spec in DEFAULT_IMAGES:
if spec["name"] in existing and not force:
skipped += 1
continue
payload = dict(spec)
checksum = _resolve_checksum(payload["location"]) or payload.get("checksum")
if not checksum:
logger.error("No checksum available for %s -- skipping", payload["name"])
failed += 1
continue
payload["checksum"] = checksum
create = core_request("POST", "images", json=payload)
if create.ok:
created += 1
click.echo(f" + {payload['name']}")
else:
failed += 1
logger.error("Core rejected %s: HTTP %s %s", payload["name"], create.status_code, create.text[:200])
click.echo(f"images:seed — {created} created, {skipped} already present, {failed} failed")
if failed and not created:
raise click.ClickException("No images could be seeded; VM creation will not work.")
if __name__ == "__main__":
command_line_interface()
+53
View File
@@ -0,0 +1,53 @@
# A generic, single database configuration.
[alembic]
# Path to migration scripts (relative to this file's directory)
script_location = .
# template used to generate migration files
# file_template = %%(rev)s_%%(slug)s
# set to 'true' to run the environment during
# the 'revision' command, regardless of autogenerate
# revision_environment = false
# Logging configuration
[loggers]
keys = root,sqlalchemy,alembic,flask_migrate
[handlers]
keys = console
[formatters]
keys = generic
[logger_root]
level = WARN
handlers = console
qualname =
[logger_sqlalchemy]
level = WARN
handlers =
qualname = sqlalchemy.engine
[logger_alembic]
level = INFO
handlers =
qualname = alembic
[logger_flask_migrate]
level = INFO
handlers =
qualname = flask_migrate
[handler_console]
class = StreamHandler
args = (sys.stderr,)
level = NOTSET
formatter = generic
[formatter_generic]
format = %(levelname)-5.5s [%(name)s] %(message)s
datefmt = %H:%M:%S
+113
View File
@@ -0,0 +1,113 @@
import logging
from logging.config import fileConfig
from flask import current_app
from alembic import context
# this is the Alembic Config object, which provides
# access to the values within the .ini file in use.
config = context.config
# Interpret the config file for Python logging.
# This line sets up loggers basically.
fileConfig(config.config_file_name)
logger = logging.getLogger('alembic.env')
def get_engine():
try:
# this works with Flask-SQLAlchemy<3 and Alchemical
return current_app.extensions['migrate'].db.get_engine()
except (TypeError, AttributeError):
# this works with Flask-SQLAlchemy>=3
return current_app.extensions['migrate'].db.engine
def get_engine_url():
try:
return get_engine().url.render_as_string(hide_password=False).replace(
'%', '%%')
except AttributeError:
return str(get_engine().url).replace('%', '%%')
# add your model's MetaData object here
# for 'autogenerate' support
# from myapp import mymodel
# target_metadata = mymodel.Base.metadata
config.set_main_option('sqlalchemy.url', get_engine_url())
target_db = current_app.extensions['migrate'].db
# other values from the config, defined by the needs of env.py,
# can be acquired:
# my_important_option = config.get_main_option("my_important_option")
# ... etc.
def get_metadata():
if hasattr(target_db, 'metadatas'):
return target_db.metadatas[None]
return target_db.metadata
def run_migrations_offline():
"""Run migrations in 'offline' mode.
This configures the context with just a URL
and not an Engine, though an Engine is acceptable
here as well. By skipping the Engine creation
we don't even need a DBAPI to be available.
Calls to context.execute() here emit the given string to the
script output.
"""
url = config.get_main_option("sqlalchemy.url")
context.configure(
url=url, target_metadata=get_metadata(), literal_binds=True
)
with context.begin_transaction():
context.run_migrations()
def run_migrations_online():
"""Run migrations in 'online' mode.
In this scenario we need to create an Engine
and associate a connection with the context.
"""
# this callback is used to prevent an auto-migration from being generated
# when there are no changes to the schema
# reference: http://alembic.zzzcomputing.com/en/latest/cookbook.html
def process_revision_directives(context, revision, directives):
if getattr(config.cmd_opts, 'autogenerate', False):
script = directives[0]
if script.upgrade_ops.is_empty():
directives[:] = []
logger.info('No changes in schema detected.')
conf_args = current_app.extensions['migrate'].configure_args
if conf_args.get("process_revision_directives") is None:
conf_args["process_revision_directives"] = process_revision_directives
connectable = get_engine()
with connectable.connect() as connection:
context.configure(
connection=connection,
target_metadata=get_metadata(),
**conf_args
)
with context.begin_transaction():
context.run_migrations()
if context.is_offline_mode():
run_migrations_offline()
else:
run_migrations_online()
+24
View File
@@ -0,0 +1,24 @@
"""${message}
Revision ID: ${up_revision}
Revises: ${down_revision | comma,n}
Create Date: ${create_date}
"""
from alembic import op
import sqlalchemy as sa
${imports if imports else ""}
# revision identifiers, used by Alembic.
revision = ${repr(up_revision)}
down_revision = ${repr(down_revision)}
branch_labels = ${repr(branch_labels)}
depends_on = ${repr(depends_on)}
def upgrade():
${upgrades if upgrades else "pass"}
def downgrade():
${downgrades if downgrades else "pass"}
@@ -0,0 +1,87 @@
from alembic import op
import sqlalchemy as sa
revision = '0001_initial'
down_revision = None
branch_labels = None
depends_on = None
def upgrade():
op.create_table(
'users',
sa.Column('id', sa.String(length=36), primary_key=True),
sa.Column('oidc_id', sa.String(length=255), nullable=False, unique=True),
sa.Column('email', sa.String(length=255), nullable=False, unique=True),
sa.Column('first_name', sa.String(length=255), nullable=True),
sa.Column('last_name', sa.String(length=255), nullable=True),
sa.Column('avatar_url', sa.String(length=512), nullable=True),
sa.Column('is_platform_admin', sa.Boolean(), nullable=False, server_default=sa.false()),
sa.Column('created_at', sa.DateTime(), nullable=False, server_default=sa.func.now()),
sa.Column('last_login_at', sa.DateTime(), nullable=True),
)
op.create_table(
'provider_regions',
sa.Column('id', sa.String(length=36), primary_key=True),
sa.Column('owner_id', sa.String(length=36), sa.ForeignKey('users.id'), nullable=False, index=True),
sa.Column('country', sa.String(length=2), nullable=False),
sa.Column('core_region_id', sa.String(length=36), nullable=False, index=True),
sa.Column('core_region_enrollment_key', sa.String(length=64), nullable=False),
sa.Column('created_at', sa.DateTime(), nullable=False, server_default=sa.func.now()),
sa.UniqueConstraint('owner_id', 'country', name='uq_provider_region_owner_country'),
)
op.create_table(
'nodes',
sa.Column('id', sa.String(length=36), primary_key=True),
sa.Column('owner_id', sa.String(length=36), sa.ForeignKey('users.id'), nullable=False, index=True),
sa.Column('name', sa.String(length=255), nullable=False),
sa.Column('country', sa.String(length=2), nullable=False, index=True),
sa.Column('city', sa.String(length=255), nullable=True),
sa.Column('enrollment_key', sa.String(length=64), nullable=True, unique=True, index=True),
sa.Column('enrolled_at', sa.DateTime(), nullable=True),
sa.Column('node_token', sa.String(length=64), nullable=True, unique=True, index=True),
sa.Column('status', sa.String(length=20), nullable=False, server_default='pending', index=True),
sa.Column('trust_tier', sa.String(length=20), nullable=False, server_default='community'),
sa.Column('last_heartbeat_at', sa.DateTime(), nullable=True),
sa.Column('cpu_cores', sa.Integer(), nullable=True),
sa.Column('memory_mb', sa.Integer(), nullable=True),
sa.Column('disk_gb', sa.Integer(), nullable=True),
sa.Column('gpu_model', sa.String(length=255), nullable=True),
sa.Column('gpu_count', sa.Integer(), nullable=True, server_default='0'),
sa.Column('kernel', sa.String(length=255), nullable=True),
sa.Column('arch', sa.String(length=50), nullable=True),
sa.Column('hourly_price_cents', sa.Integer(), nullable=False, server_default='0'),
sa.Column('provider_region_id', sa.String(length=36), sa.ForeignKey('provider_regions.id'), nullable=True, index=True),
sa.Column('core_workload_host_id', sa.String(length=36), nullable=True, index=True),
sa.Column('created_at', sa.DateTime(), nullable=False, server_default=sa.func.now()),
)
op.create_table(
'leases',
sa.Column('id', sa.String(length=36), primary_key=True),
sa.Column('node_id', sa.String(length=36), sa.ForeignKey('nodes.id'), nullable=False, index=True),
sa.Column('customer_id', sa.String(length=36), sa.ForeignKey('users.id'), nullable=False, index=True),
sa.Column('started_at', sa.DateTime(), nullable=False, server_default=sa.func.now()),
sa.Column('ended_at', sa.DateTime(), nullable=True),
)
op.create_table(
'ping_tokens',
sa.Column('id', sa.String(length=36), primary_key=True),
sa.Column('node_id', sa.String(length=36), sa.ForeignKey('nodes.id'), nullable=False, index=True),
sa.Column('user_id', sa.String(length=36), sa.ForeignKey('users.id'), nullable=False),
sa.Column('token', sa.String(length=64), nullable=False, unique=True, index=True),
sa.Column('expires_at', sa.DateTime(), nullable=False),
sa.Column('used_at', sa.DateTime(), nullable=True),
sa.Column('created_at', sa.DateTime(), nullable=False, server_default=sa.func.now()),
)
def downgrade():
op.drop_table('ping_tokens')
op.drop_table('leases')
op.drop_table('nodes')
op.drop_table('provider_regions')
op.drop_table('users')
@@ -0,0 +1,39 @@
from alembic import op
import sqlalchemy as sa
revision = '0002_drop_oidc_naming'
down_revision = '0001_initial'
branch_labels = None
depends_on = None
def _columns():
bind = op.get_bind()
return {c["name"] for c in sa.inspect(bind).get_columns("users")}
def upgrade():
cols = _columns()
if "external_id" in cols:
return
if "oidc_id" not in cols:
raise RuntimeError("users has neither oidc_id nor external_id")
op.alter_column(
'users', 'oidc_id',
new_column_name='external_id',
existing_type=sa.String(length=255),
existing_nullable=False,
)
def downgrade():
cols = _columns()
if "oidc_id" in cols:
return
op.alter_column(
'users', 'external_id',
new_column_name='oidc_id',
existing_type=sa.String(length=255),
existing_nullable=False,
)
@@ -0,0 +1,25 @@
from alembic import op
import sqlalchemy as sa
revision = '0003_drop_hourly_price'
down_revision = '0002_drop_oidc_naming'
branch_labels = None
depends_on = None
def _columns():
bind = op.get_bind()
return {c["name"] for c in sa.inspect(bind).get_columns("nodes")}
def upgrade():
if "hourly_price_cents" in _columns():
op.drop_column('nodes', 'hourly_price_cents')
def downgrade():
if "hourly_price_cents" not in _columns():
op.add_column(
'nodes',
sa.Column('hourly_price_cents', sa.Integer(), nullable=False, server_default='0'),
)
@@ -0,0 +1,71 @@
from alembic import op
import sqlalchemy as sa
revision = "0004_cloudflare"
down_revision = "0003_drop_hourly_price"
branch_labels = None
depends_on = None
def _tables():
bind = op.get_bind()
return set(sa.inspect(bind).get_table_names())
def upgrade():
tables = _tables()
if "cloudflare_credentials" not in tables:
op.create_table(
"cloudflare_credentials",
sa.Column("id", sa.String(36), primary_key=True),
sa.Column("user_id", sa.String(36), sa.ForeignKey("users.id"), nullable=False, unique=True),
sa.Column("api_token_encrypted", sa.String(1024), nullable=False),
sa.Column("account_id", sa.String(255), nullable=False),
sa.Column("zone_id", sa.String(255), nullable=False),
sa.Column("domain", sa.String(255), nullable=False),
sa.Column("verified_at", sa.DateTime(), nullable=False),
sa.Column("created_at", sa.DateTime(), nullable=False),
)
if "cloudflare_tunnels" not in tables:
op.create_table(
"cloudflare_tunnels",
sa.Column("id", sa.String(36), primary_key=True),
sa.Column("lease_id", sa.String(36), sa.ForeignKey("leases.id"), nullable=False, index=True),
sa.Column("name", sa.String(255), nullable=False),
sa.Column("account_id", sa.String(255), nullable=False),
sa.Column("tunnel_id", sa.String(255), nullable=False, unique=True),
sa.Column("tunnel_secret", sa.String(255), nullable=False),
sa.Column("token", sa.String(1024), nullable=False),
sa.Column("pod_id", sa.String(36), nullable=True),
sa.Column("created_at", sa.DateTime(), nullable=False),
sa.Column("deleted", sa.Boolean(), nullable=False, server_default=sa.false()),
sa.Column("deleted_at", sa.DateTime(), nullable=True),
)
if "cloudflare_dns_records" not in tables:
op.create_table(
"cloudflare_dns_records",
sa.Column("id", sa.String(36), primary_key=True),
sa.Column("tunnel_id", sa.String(36), sa.ForeignKey("cloudflare_tunnels.id"), nullable=False, index=True),
sa.Column("zone_id", sa.String(255), nullable=False),
sa.Column("dns_record_id", sa.String(255), nullable=False, unique=True),
sa.Column("hostname", sa.String(255), nullable=False),
sa.Column("content", sa.String(255), nullable=False),
sa.Column("container_workload_id", sa.String(36), nullable=True, index=True),
sa.Column("internal_port", sa.Integer(), nullable=True),
sa.Column("created_at", sa.DateTime(), nullable=False),
sa.Column("deleted", sa.Boolean(), nullable=False, server_default=sa.false()),
sa.Column("deleted_at", sa.DateTime(), nullable=True),
)
def downgrade():
tables = _tables()
if "cloudflare_dns_records" in tables:
op.drop_table("cloudflare_dns_records")
if "cloudflare_tunnels" in tables:
op.drop_table("cloudflare_tunnels")
if "cloudflare_credentials" in tables:
op.drop_table("cloudflare_credentials")
@@ -0,0 +1,34 @@
from alembic import op
import sqlalchemy as sa
revision = "0005_ssh_keys"
down_revision = "0004_cloudflare"
branch_labels = None
depends_on = None
def _tables():
bind = op.get_bind()
return set(sa.inspect(bind).get_table_names())
def upgrade():
if "ssh_keys" not in _tables():
op.create_table(
"ssh_keys",
sa.Column("id", sa.String(36), primary_key=True),
sa.Column("user_id", sa.String(36), sa.ForeignKey("users.id"), nullable=False, index=True),
sa.Column("key_name", sa.String(255), nullable=False),
sa.Column("public_key_data", sa.String(4096), nullable=False),
sa.Column("key_fingerprint", sa.String(255), nullable=True),
sa.Column("is_default", sa.Boolean(), nullable=False, server_default=sa.false()),
sa.Column("created_at", sa.DateTime(), nullable=False),
sa.Column("updated_at", sa.DateTime(), nullable=False),
sa.Column("deleted", sa.Boolean(), nullable=False, server_default=sa.false()),
sa.Column("deleted_at", sa.DateTime(), nullable=True),
)
def downgrade():
if "ssh_keys" in _tables():
op.drop_table("ssh_keys")
+9
View File
@@ -0,0 +1,9 @@
flask
flask_sqlalchemy
flask_migrate
flask_cors
pymysql
requests
colorlog
click
gunicorn
+65
View File
@@ -0,0 +1,65 @@
import os
def _env(key: str, default: str = "") -> str:
return os.environ.get(key, default)
def _env_int(key: str, default: int) -> int:
try:
return int(os.environ.get(key, default))
except (TypeError, ValueError):
return default
COMMUNITY_DATABASE_URL = _env(
"COMMUNITY_DATABASE_URL",
"mysql+pymysql://community_user:community_password@127.0.0.1:3308/community",
)
PUBLIC_BASE_URL = _env("PUBLIC_BASE_URL", "http://localhost:5002")
SIGNALING_PUBLIC_URL = _env("SIGNALING_PUBLIC_URL", "ws://localhost:5003")
SIGNALING_BROWSER_URL = _env("SIGNALING_BROWSER_URL", "") or SIGNALING_PUBLIC_URL
SIGNALING_INTERNAL_URL = _env("SIGNALING_INTERNAL_URL", "http://community-signaling:5003")
SIGNALING_SHARED_SECRET = _env("SIGNALING_SHARED_SECRET", "")
STUN_URLS = _env("STUN_URLS", "stun:stun.l.google.com:19302")
WORKER_BINARY_URL = _env("WORKER_BINARY_URL", "")
WORKER_BINARY_SHA256 = _env("WORKER_BINARY_SHA256", "")
NSCONTROLLER_IMAGE_URL = _env("NSCONTROLLER_IMAGE_URL", "")
NSCONTROLLER_IMAGE_SHA256 = _env("NSCONTROLLER_IMAGE_SHA256", "")
CORE_API_BASE_URL = _env("CORE_API_BASE_URL", "http://api-server:5000")
CORE_API_KEY = _env("CORE_API_KEY", "")
CORE_PUBLIC_BASE_URL = _env("CORE_PUBLIC_BASE_URL", "")
CORE_PUBLIC_WS_URL = _env("CORE_PUBLIC_WS_URL", "")
CORE_REGION_NAME_PREFIX = _env("CORE_REGION_NAME_PREFIX", "Community")
HEARTBEAT_TIMEOUT_SECONDS = _env_int("HEARTBEAT_TIMEOUT_SECONDS", 120)
PING_TOKEN_TTL_SECONDS = _env_int("PING_TOKEN_TTL_SECONDS", 120)
LOCAL_USER_EMAIL = _env("LOCAL_USER_EMAIL", "local@xcloudify.dev")
BOOTSTRAP_ADMIN_EMAILS = _env("BOOTSTRAP_ADMIN_EMAILS", "")
# Unset means production. The dev bypass (X-Dev-User, no sign-in) is only
# possible in a development APP_ENV, where it is on unless AUTH_DEV_BYPASS=false.
APP_ENV = _env("APP_ENV", "production")
AUTH_DEV_BYPASS = (
os.environ["AUTH_DEV_BYPASS"].strip().lower() in ("1", "true", "yes", "on")
if os.environ.get("AUTH_DEV_BYPASS") else None
)
OIDC_ISSUER = _env("OIDC_ISSUER", "")
OIDC_JWKS_URL = _env("OIDC_JWKS_URL", "")
OIDC_AUDIENCE = _env("OIDC_AUDIENCE", "")
OIDC_ADDITIONAL_AUDIENCES = _env("OIDC_ADDITIONAL_AUDIENCES", "")
OIDC_SUBJECT_CLAIM = _env("OIDC_SUBJECT_CLAIM", "sub")
OIDC_EMAIL_CLAIM = _env("OIDC_EMAIL_CLAIM", "email")
FLASK_ENV = _env("FLASK_ENV", "development")
COMMUNITY_SECRET_KEY = _env("COMMUNITY_SECRET_KEY", "")
+13
View File
@@ -0,0 +1,13 @@
FROM python:3.11-slim
ENV PYTHONDONTWRITEBYTECODE=1 \
PYTHONUNBUFFERED=1
WORKDIR /app
COPY community/signaling/requirements.txt .
RUN pip install --upgrade pip && pip install -r requirements.txt
COPY community/signaling/ .
CMD ["uvicorn", "main:app", "--host", "0.0.0.0", "--port", "5003"]
+95
View File
@@ -0,0 +1,95 @@
import json
import logging
import os
from fastapi import FastAPI, Header, HTTPException, WebSocket, WebSocketDisconnect
app = FastAPI()
logger = logging.getLogger("uvicorn.error")
workers: dict[str, WebSocket] = {}
browsers: dict[str, WebSocket] = {}
SHARED_SECRET = os.environ.get("SIGNALING_SHARED_SECRET", "")
MAX_MESSAGE_BYTES = 64 * 1024
@app.get("/healthz")
async def healthz():
return {"status": "ok"}
@app.get("/internal/workers")
async def list_workers(x_signaling_secret: str = Header(default="")):
"""Which nodes currently have a worker session.
The community API calls this to label a node pingable. Guarded by a
shared secret because it enumerates online hardware; not a security
boundary for the signaling itself, which has none to protect.
"""
if SHARED_SECRET and x_signaling_secret != SHARED_SECRET:
raise HTTPException(status_code=401, detail="bad signaling secret")
return {"workers": list(workers.keys())}
@app.websocket("/worker/{node_id}")
async def worker_socket(ws: WebSocket, node_id: str):
await ws.accept()
previous = workers.get(node_id)
if previous is not None:
logger.warning(
"Replacing an existing signaling session for node %s. If this repeats, "
"more than one agent is running for that node.", node_id,
)
try:
await previous.close()
except Exception:
pass
workers[node_id] = ws
try:
while True:
msg = await ws.receive_json()
browser = browsers.get(msg.get("client_id"))
if browser is not None:
await browser.send_json({**msg, "worker_id": node_id})
except WebSocketDisconnect:
pass
finally:
if workers.get(node_id) is ws:
workers.pop(node_id, None)
@app.websocket("/browser/{client_id}")
async def browser_socket(ws: WebSocket, client_id: str):
await ws.accept()
browsers[client_id] = ws
try:
while True:
raw = await ws.receive_text()
if len(raw) > MAX_MESSAGE_BYTES:
await ws.send_json({"type": "error", "message": "message too large"})
continue
try:
msg = json.loads(raw)
except ValueError:
await ws.send_json({"type": "error", "message": "malformed message"})
continue
worker = workers.get(msg.get("worker_id"))
if worker is None:
await ws.send_json({
"type": "error",
"message": "That node's worker is not connected.",
})
continue
await worker.send_json({**msg, "client_id": client_id})
except WebSocketDisconnect:
pass
finally:
if browsers.get(client_id) is ws:
browsers.pop(client_id, None)
+3
View File
@@ -0,0 +1,3 @@
fastapi
uvicorn[standard]
websockets
+281
View File
@@ -0,0 +1,281 @@
#!/usr/bin/env bash
set -euo pipefail
BASE_URL="{{ base_url }}"
SIGNALING_URL="{{ signaling_url }}"
STUN_URLS="{{ stun_urls }}"
ENROLLMENT_KEY="{{ enrollment_key }}"
WORKER_BINARY_URL="{{ worker_binary_url }}"
WORKER_BINARY_SHA256="{{ worker_binary_sha256 }}"
NSCONTROLLER_IMAGE_URL="{{ nscontroller_image_url }}"
NSCONTROLLER_IMAGE_SHA256="{{ nscontroller_image_sha256 }}"
CORE_PUBLIC_URL="{{ core_public_url }}"
CORE_REGION_ID="{{ core_region_id }}"
CORE_REGION_ENROLLMENT_KEY="{{ core_region_enrollment_key }}"
CORE_WS_URL="{{ core_ws_url }}"
INSTALL_DIR="/opt/xcloudify"
BIN_PATH="${INSTALL_DIR}/xcloudify-worker"
CONF_PATH="${INSTALL_DIR}/worker.env"
log() { printf '\033[0;32m==>\033[0m %s\n' "$*"; }
warn() { printf '\033[0;33m!!\033[0m %s\n' "$*" >&2; }
fail() { printf '\033[0;31m!!\033[0m %s\n' "$*" >&2; exit 1; }
SKIP_PROVISION=0
for arg in "$@"; do
case "$arg" in
--no-provision) SKIP_PROVISION=1 ;;
*) printf 'unknown argument: %s\n' "$arg" >&2; exit 1 ;;
esac
done
[ "$(id -u)" -eq 0 ] || fail "Run as root: curl -fsSL ... | sudo bash"
for cmd in curl systemctl; do
command -v "$cmd" >/dev/null 2>&1 || fail "required command not found: $cmd"
done
have() { command -v "$1" >/dev/null 2>&1; }
has_docker() { [ -S /var/run/docker.sock ] || have dockerd; }
has_libvirt() { [ -S /run/libvirt/libvirt-sock ] || have libvirtd; }
has_ovs() { [ -S /run/openvswitch/db.sock ] || have ovs-vswitchd; }
COMPOSE=""
detect_compose() {
if docker compose version >/dev/null 2>&1; then
COMPOSE="docker compose"
elif have docker-compose; then
COMPOSE="docker-compose"
else
return 1
fi
}
provision_host() {
local missing=()
has_docker || missing+=("docker")
has_libvirt || missing+=("libvirt")
has_ovs || missing+=("openvswitch")
if [ ${#missing[@]} -eq 0 ]; then
log "Host already provides docker, libvirt and openvswitch"
return 0
fi
log "Installing host dependencies: ${missing[*]}"
have apt-get || fail "This installer provisions with apt-get, which is not present here.
Install docker, libvirt (with a qemu emulator) and openvswitch by hand, then
re-run this command with: | sudo bash -s -- --no-provision"
export DEBIAN_FRONTEND=noninteractive
apt-get update -qq || fail "apt-get update failed"
if ! has_libvirt; then
apt-get install -y -qq qemu-system-x86 >/dev/null 2>&1 \
|| apt-get install -y -qq qemu-kvm >/dev/null 2>&1 \
|| warn "No qemu emulator could be installed; VM workloads will not run here."
fi
local pkgs=()
has_docker || pkgs+=(docker.io)
has_libvirt || pkgs+=(libvirt-daemon-system libvirt-clients)
has_ovs || pkgs+=(openvswitch-switch)
if [ ${#pkgs[@]} -gt 0 ]; then
apt-get install -y -qq "${pkgs[@]}" >/dev/null || fail "Failed to install: ${pkgs[*]}"
fi
for svc in docker libvirtd openvswitch-switch; do
systemctl enable --now "$svc" >/dev/null 2>&1 || true
done
if ! detect_compose; then
for pkg in docker-compose-v2 docker-compose-plugin docker-compose; do
apt-get install -y -qq "$pkg" >/dev/null 2>&1 && detect_compose && break
done
fi
local broken=()
has_docker || broken+=("docker")
has_libvirt || broken+=("libvirt")
has_ovs || broken+=("openvswitch")
[ ${#broken[@]} -eq 0 ] || fail "Installed, but the daemon is not up: ${broken[*]}
Check: systemctl status docker libvirtd openvswitch-switch"
log "Host dependencies ready"
}
if [ "$SKIP_PROVISION" -eq 1 ]; then
log "Skipping host provisioning (--no-provision)"
else
provision_host
fi
have docker || fail "docker is required but not present. Install it, or re-run without --no-provision."
detect_compose || fail "docker compose is required but not present.
Install one of: docker-compose-v2, docker-compose-plugin, docker-compose"
[ -e /dev/kvm ] || warn "/dev/kvm is missing -- VM workloads cannot run on this node.
Enable VT-x/AMD-V in the BIOS, or turn on nested virtualization if this is a VM."
log "Collecting system specs"
digits() { local v; v="$(printf '%s' "${1:-}" | tr -dc '0-9')"; printf '%s' "${v:-0}"; }
jsonstr() { printf '%s' "${1:-}" | tr -d '\000-\037"\\' | cut -c1-250; }
CPU_CORES="$(digits "$(nproc 2>/dev/null || true)")"
MEMORY_MB="$(digits "$(awk '/MemTotal/ {printf "%d", $2/1024}' /proc/meminfo 2>/dev/null || true)")"
DISK_GB="$(digits "$(df -BG --output=size / 2>/dev/null | tail -1 || true)")"
KERNEL="$(jsonstr "$(uname -r)")"
ARCH="$(jsonstr "$(uname -m)")"
GPU_MODEL=""
GPU_COUNT=0
if command -v nvidia-smi >/dev/null 2>&1 && nvidia-smi >/dev/null 2>&1; then
GPU_MODEL="$(jsonstr "$(nvidia-smi --query-gpu=name --format=csv,noheader 2>/dev/null | head -1 || true)")"
GPU_COUNT="$(digits "$(nvidia-smi --list-gpus 2>/dev/null | wc -l || true)")"
fi
mkdir -p "$INSTALL_DIR"
if [ -n "$WORKER_BINARY_URL" ]; then
log "Downloading worker binary"
curl -fsSL "$WORKER_BINARY_URL" -o "${BIN_PATH}.tmp" || fail "Failed to download worker binary"
if [ -n "$WORKER_BINARY_SHA256" ]; then
log "Verifying checksum"
echo "${WORKER_BINARY_SHA256} ${BIN_PATH}.tmp" | sha256sum -c - \
|| { rm -f "${BIN_PATH}.tmp"; fail "Checksum mismatch -- refusing to install"; }
else
warn "No WORKER_BINARY_SHA256 configured; binary integrity NOT verified"
fi
mv "${BIN_PATH}.tmp" "$BIN_PATH"
chmod 0755 "$BIN_PATH"
log "Downloading runtime files"
for f in Dockerfile docker-compose.yml; do
curl -fsSL "${BASE_URL}/api/nodes/worker-runtime/${f}" -o "${INSTALL_DIR}/${f}" \
|| fail "Failed to download ${f}"
done
log "Building worker image"
docker build -t xcloudify-worker:local "$INSTALL_DIR" >/dev/null \
|| fail "Failed to build the worker image. Re-run with the output visible:
docker build -t xcloudify-worker:local ${INSTALL_DIR}"
WORKER_INSTALLED=1
else
warn "WORKER_BINARY_URL not configured on the server; skipping worker install."
warn " The node is enrolled and will appear in the portal, but cannot run workloads yet."
WORKER_INSTALLED=0
fi
if [ -n "$NSCONTROLLER_IMAGE_URL" ]; then
log "Downloading NSController image"
curl -fsSL "$NSCONTROLLER_IMAGE_URL" -o "${INSTALL_DIR}/nscontroller-image.tar.gz.tmp" \
|| fail "Failed to download the NSController image"
if [ -n "$NSCONTROLLER_IMAGE_SHA256" ]; then
log "Verifying checksum"
echo "${NSCONTROLLER_IMAGE_SHA256} ${INSTALL_DIR}/nscontroller-image.tar.gz.tmp" | sha256sum -c - \
|| { rm -f "${INSTALL_DIR}/nscontroller-image.tar.gz.tmp"; fail "Checksum mismatch -- refusing to load"; }
else
warn "No NSCONTROLLER_IMAGE_SHA256 configured; image integrity NOT verified"
fi
mv "${INSTALL_DIR}/nscontroller-image.tar.gz.tmp" "${INSTALL_DIR}/nscontroller-image.tar.gz"
log "Loading NSController image into Docker"
gunzip -c "${INSTALL_DIR}/nscontroller-image.tar.gz" | docker load >/dev/null \
|| fail "Failed to load the NSController image"
else
warn "NSCONTROLLER_IMAGE_URL not configured on the server; skipping."
warn " The node is enrolled and can run VM workloads, but container/pod"
warn " workloads will fail until this image is present -- run"
warn " core/worker/build_nscontroller_image.sh and re-run this installer."
fi
log "Enrolling with ${BASE_URL}"
ENROLL_BODY="$(printf \
'{"enrollment_key":"%s","cpu_cores":%s,"memory_mb":%s,"disk_gb":%s,"gpu_model":"%s","gpu_count":%s,"kernel":"%s","arch":"%s"}' \
"$ENROLLMENT_KEY" "$CPU_CORES" "$MEMORY_MB" "$DISK_GB" "$GPU_MODEL" "$GPU_COUNT" "$KERNEL" "$ARCH")"
ENROLL_RAW="$(curl -sS -X POST "${BASE_URL}/api/nodes/enroll" \
-H 'Content-Type: application/json' \
-d "$ENROLL_BODY" \
-w $'\n%{http_code}')" || fail "Could not reach ${BASE_URL} -- is it accessible from this machine?"
ENROLL_STATUS="$(printf '%s' "$ENROLL_RAW" | tail -n1)"
ENROLL_RESPONSE="$(printf '%s' "$ENROLL_RAW" | sed '$d')"
if [ "$ENROLL_STATUS" != "201" ]; then
case "$ENROLL_STATUS" in
404) fail "Enrollment key is invalid or already used. Register the node again for a fresh command." ;;
*) fail "Enrollment failed (HTTP ${ENROLL_STATUS}): ${ENROLL_RESPONSE}" ;;
esac
fi
NODE_TOKEN="$(printf '%s' "$ENROLL_RESPONSE" | grep -o '"node_token"[[:space:]]*:[[:space:]]*"[^"]*"' | sed 's/.*"\([^"]*\)"$/\1/' || true)"
NODE_ID="$(printf '%s' "$ENROLL_RESPONSE" | grep -o '"node_id"[[:space:]]*:[[:space:]]*"[^"]*"' | sed 's/.*"\([^"]*\)"$/\1/' || true)"
[ -n "$NODE_TOKEN" ] || fail "No node_token in enrollment response: ${ENROLL_RESPONSE}"
log "Enrolled as node ${NODE_ID}"
log "Writing configuration"
umask 077
cat > "$CONF_PATH" <<EOF
XCLOUDIFY_NODE_ID=${NODE_ID}
XCLOUDIFY_NODE_TOKEN=${NODE_TOKEN}
XCLOUDIFY_COMMUNITY_URL=${BASE_URL}
XCLOUDIFY_SIGNALING_URL=${SIGNALING_URL}
XCLOUDIFY_STUN_URLS=${STUN_URLS}
EOF
if [ -n "$CORE_PUBLIC_URL" ] && [ -n "$CORE_REGION_ID" ] && [ -n "$CORE_REGION_ENROLLMENT_KEY" ]; then
cat >> "$CONF_PATH" <<EOF
API_BASE_URL=${CORE_PUBLIC_URL}/api
WEBSOCKET_SERVER_URL=${CORE_WS_URL}
WEBSOCKET_SERVER_PATH=/ws
REGION_ID=${CORE_REGION_ID}
REGION_ENROLLMENT_KEY=${CORE_REGION_ENROLLMENT_KEY}
PHYSICAL_IDENTIFIER=${NODE_ID}
EOF
WORKLOADS_ENABLED=1
else
warn "No core endpoint configured on the server (CORE_PUBLIC_BASE_URL)."
warn " This node will enroll, appear in the portal and answer latency tests,"
warn " but cannot run workloads."
WORKLOADS_ENABLED=0
fi
chmod 0600 "$CONF_PATH"
mkdir -p /var/lib/xcloudify/worker /tmp/xcloudify
if [ "$WORKLOADS_ENABLED" -eq 1 ]; then
rm -f "${INSTALL_DIR}/docker-compose.override.yml"
else
cat > "${INSTALL_DIR}/docker-compose.override.yml" <<'EOF'
services:
worker:
command: ["--community-agent"]
EOF
fi
if [ "$WORKER_INSTALLED" -eq 1 ]; then
log "Starting worker"
( cd "$INSTALL_DIR" && $COMPOSE up -d --force-recreate ) >/dev/null \
|| fail "Failed to start the worker container. Logs:
cd ${INSTALL_DIR} && ${COMPOSE} logs"
log "Worker started."
else
log "Worker configured but not started (no image built)."
fi
if [ "$WORKLOADS_ENABLED" -eq 1 ]; then
log "Done. Node ${NODE_ID} is registered and will enroll into core as a schedulable host."
else
log "Done. Node ${NODE_ID} is registered (latency tests only)."
fi
log " logs: cd ${INSTALL_DIR} && ${COMPOSE} logs -f"
log " restart: cd ${INSTALL_DIR} && ${COMPOSE} restart"