net-kingdom/identity-provisioner/provisioner.py
tegwick c8e07615c3
All checks were successful
CI Smoke / host-smoke (push) Successful in 0s
CI Smoke / container-smoke (push) Successful in 3s
Identity provider journey acceptance / provider (push) Successful in 6s
Build and Publish identity-provisioner / build-and-push (push) Successful in 10s
Surface redacted directory bind failures before native onboarding
Map uncaught HTTPError from LLDAP login to a structured
dependency_unavailable response, add /readyz as the provisioner-to-directory
preflight, keep /healthz as process liveness, and run the contract in CI.
Auth rejection is not retried during cooldown.

NK-WP-0036-T05 remains in progress until the immutable image is published,
pinned with /readyz, and one native login/create/password-setup journey is
verified.

Assistant: grok
Assistant-Session: 01a09dc6-3f0e-78f1-a884-c8c703c24ddf
2026-09-14 04:46:29 +02:00

450 lines
18 KiB
Python

"""NetKingdom's idempotent LLDAP lifecycle adapter."""
from __future__ import annotations
from dataclasses import dataclass
import json
import re
import time
from typing import Any, Callable
from urllib.error import HTTPError, URLError
from urllib.request import Request, urlopen
@dataclass(frozen=True)
class Result:
provider: str
external_subject: str
status: str
resumed: bool
@dataclass(frozen=True)
class DriftResult:
provider: str
external_subject: str
status: str
drift: tuple[str, ...]
changed: tuple[str, ...] = ()
class DependencyFailure(RuntimeError):
"""Directory dependency failed. Messages never carry credentials or upstream bodies."""
def __init__(self, code: str, *, dependency: str = "directory") -> None:
self.code = code
self.dependency = dependency
super().__init__("dependency_unavailable")
def payload(self) -> dict[str, str]:
return {
"error": "dependency_unavailable",
"dependency": self.dependency,
"reason": self.code,
}
class LLDAPProvisioner:
def __init__(
self,
*,
base_url: str,
admin_password: str,
opener: Callable = urlopen,
auth_rejected_cooldown: float = 30.0,
clock: Callable[[], float] = time.monotonic,
) -> None:
self.base_url = base_url.rstrip("/")
self.admin_password = admin_password
self.opener = opener
self.auth_rejected_cooldown = auth_rejected_cooldown
self.clock = clock
self._auth_rejected_until = 0.0
def provision(self, payload: dict[str, Any]) -> Result:
_required(payload, "user_id", "tenant", "primary_email", "idempotency_key", "correlation_id")
email = str(payload["primary_email"]).strip().lower()
username = _username(email, payload.get("preferred_username"))
token = self._login()
users, groups = self._directory(token)
existing = next((user for user in users if user.get("id") == username), None)
resumed = existing is not None
created = existing is None
if created:
self._gql(token, """
mutation CreateUser($id: String!, $email: String!, $display: String!) {
createUser(user: {id: $id, email: $email, displayName: $display}) { id }
}""", {
"id": username,
"email": email,
"display": str(payload.get("display_name") or username),
})
elif str(existing.get("email", "")).lower() != email:
raise ValueError("directory username collision")
roles = {str(role) for role in payload.get("roles", ())}
group_names = [f"{payload['tenant']}:users"]
if "tenant-admin" in roles:
group_names.append(f"{payload['tenant']}:admins")
try:
for name in group_names:
group_id = self._ensure_group(token, groups, name)
self._add_group(token, username, group_id)
except Exception:
if created:
try:
self._delete(token, username)
except Exception:
pass
raise
return Result("netkingdom-lldap", _oidc_subject(username), "password_setup_required", resumed)
def tenant_access(self, payload: dict[str, Any]) -> Result:
"""Change only the two role groups owned by the specified tenant."""
_required(payload, "external_subject", "tenant", "idempotency_key", "correlation_id")
tenant = str(payload["tenant"])
if not re.fullmatch(r"tenant:[a-z0-9][a-z0-9._-]*(?::[a-z0-9][a-z0-9._-]*)*", tenant):
raise ValueError("invalid tenant identifier")
if tenant == "tenant:platform:root":
raise ValueError("platform root is not a tenant access target")
if not isinstance(payload.get("enabled"), bool):
raise ValueError("enabled must be boolean")
roles = payload.get("roles", [])
if not isinstance(roles, (list, tuple)) or any(r not in {"user", "tenant-admin"} for r in roles):
raise ValueError("unsupported tenant role")
subject = _directory_username(str(payload["external_subject"]))
if not re.fullmatch(r"[A-Za-z0-9._-]+", subject):
raise ValueError("invalid directory subject")
token = self._login()
user = self._user(token, subject)
if user is None:
raise ValueError("login identity not found; create it before changing access")
managed = {f"{tenant}:users", f"{tenant}:admins"}
desired = {f"{tenant}:users"} if payload["enabled"] else set()
if payload["enabled"] and "tenant-admin" in roles:
desired.add(f"{tenant}:admins")
current = {str(g["displayName"]): int(g["id"]) for g in user.get("groups", ())}
groups = list(self._directory(token)[1])
for name in sorted(desired - current.keys()):
self._add_group(token, subject, self._ensure_group(token, groups, name))
for name in sorted((managed & current.keys()) - desired):
self._remove_group(token, subject, current[name])
checked = self._user(token, subject)
if checked is None or ({g["displayName"] for g in checked.get("groups", ())} & managed) != desired:
raise RuntimeError("tenant access readback did not confirm the requested state")
return Result("netkingdom-lldap", _oidc_subject(subject),
"tenant_active" if payload["enabled"] else "tenant_disabled", False)
def suspend(self, subject: str) -> Result:
subject = _directory_username(subject)
token = self._login()
_, groups = self._directory(token)
group_id = self._ensure_group(token, groups, "netkingdom-suspended")
self._add_group(token, subject, group_id)
return Result("netkingdom-lldap", subject, "suspended", False)
def reactivate(self, subject: str) -> Result:
subject = _directory_username(subject)
token = self._login()
_, groups = self._directory(token)
group = next((item for item in groups if item.get("displayName") == "netkingdom-suspended"), None)
if group:
self._gql(token, """
mutation Remove($userId: String!, $groupId: Int!) {
removeUserFromGroup(userId: $userId, groupId: $groupId) { ok }
}""", {"userId": subject, "groupId": int(group["id"])})
return Result("netkingdom-lldap", subject, "active", False)
def deprovision(self, subject: str) -> Result:
subject = _directory_username(subject)
token = self._login()
if self._user(token, subject) is None:
return Result("netkingdom-lldap", subject, "deprovisioned", True)
self._delete(token, subject)
return Result("netkingdom-lldap", subject, "deprovisioned", False)
def drift(self, payload: dict[str, Any]) -> DriftResult:
subject, email, desired_groups, desired_status = _desired(payload)
token = self._login()
user = self._user(token, subject)
drift = self._drift(user, email, desired_groups, desired_status, str(payload["tenant"]))
status = "in_sync" if not drift else "drifted"
return DriftResult("netkingdom-lldap", subject, status, tuple(drift))
def reconcile(self, payload: dict[str, Any]) -> DriftResult:
subject, email, desired_groups, desired_status = _desired(payload)
token = self._login()
user = self._user(token, subject)
changed: list[str] = []
if user is None:
result = self.provision(payload)
changed.append("user:created")
subject = result.external_subject
token = self._login()
user = self._user(token, subject)
if user is None:
raise RuntimeError("directory reconciliation did not create the identity")
if str(user.get("email", "")).lower() != email:
raise ValueError("directory email drift requires explicit identity repair")
groups = list(self._directory(token)[1])
current = {str(item["displayName"]): int(item["id"]) for item in user.get("groups", ())}
tenant = str(payload["tenant"])
managed = {
name for name in current
if name in {f"{tenant}:users", f"{tenant}:admins", "netkingdom-suspended"}
}
for name in sorted(desired_groups - managed):
self._add_group(token, subject, self._ensure_group(token, groups, name))
changed.append(f"group:added:{name}")
for name in sorted(managed - desired_groups):
self._remove_group(token, subject, current[name])
changed.append(f"group:removed:{name}")
user = self._user(token, subject)
remaining = self._drift(user, email, desired_groups, desired_status, tenant)
status = "reconciled" if not remaining else "drifted"
return DriftResult("netkingdom-lldap", subject, status, tuple(remaining), tuple(changed))
def preflight(self) -> dict[str, str]:
"""One login plus one directory read. Auth rejection is not retried during cooldown."""
now = self.clock()
if now < self._auth_rejected_until:
raise DependencyFailure("auth_rejected")
try:
token = self._login(timeout=3)
self._gql(token, "query { groups { id } }", {}, timeout=3)
except DependencyFailure as exc:
if exc.code == "auth_rejected":
self._auth_rejected_until = now + self.auth_rejected_cooldown
raise
return {"status": "ready", "dependency": "directory"}
def _login(self, timeout: float = 10) -> str:
return directory_login(
base_url=self.base_url,
admin_password=self.admin_password,
opener=self.opener,
timeout=timeout,
)
def _directory(self, token: str) -> tuple[list[dict], list[dict]]:
value = self._gql(token, "query { users { id email displayName } groups { id displayName } }", {})
return list(value["users"]), list(value["groups"])
def _user(self, token: str, subject: str) -> dict[str, Any] | None:
try:
value = self._gql(token, """
query User($id: String!) {
user(userId: $id) { id email displayName groups { id displayName } }
}""", {"id": subject})
except ValueError as exc:
if "not found" in str(exc).lower():
return None
raise
user = value.get("user")
return dict(user) if user else None
def _ensure_group(self, token: str, groups: list[dict], name: str) -> int:
existing = next((group for group in groups if group.get("displayName") == name), None)
if existing:
return int(existing["id"])
created = self._gql(
token,
"mutation CreateGroup($name: String!) { createGroup(name: $name) { id displayName } }",
{"name": name},
)["createGroup"]
groups.append(created)
return int(created["id"])
def _add_group(self, token: str, username: str, group_id: int) -> None:
try:
self._gql(token, """
mutation Add($userId: String!, $groupId: Int!) {
addUserToGroup(userId: $userId, groupId: $groupId) { ok }
}""", {"userId": username, "groupId": group_id})
except ValueError as exc:
if "already" not in str(exc).lower() and "unique" not in str(exc).lower():
raise
def _remove_group(self, token: str, username: str, group_id: int) -> None:
self._gql(token, """
mutation Remove($userId: String!, $groupId: Int!) {
removeUserFromGroup(userId: $userId, groupId: $groupId) { ok }
}""", {"userId": username, "groupId": group_id})
def _delete(self, token: str, subject: str) -> None:
self._gql(
token,
"mutation Delete($id: String!) { deleteUser(userId: $id) { ok } }",
{"id": subject},
)
@staticmethod
def _drift(
user: dict[str, Any] | None,
email: str,
desired_groups: set[str],
desired_status: str,
tenant: str,
) -> list[str]:
if user is None:
return ["user:missing"]
drift: list[str] = []
if str(user.get("email", "")).lower() != email:
drift.append("email:mismatch")
current = {str(item["displayName"]) for item in user.get("groups", ())}
managed = {
name for name in current
if name in {f"{tenant}:users", f"{tenant}:admins", "netkingdom-suspended"}
}
drift.extend(f"group:missing:{name}" for name in sorted(desired_groups - managed))
drift.extend(f"group:unexpected:{name}" for name in sorted(managed - desired_groups))
suspended = "netkingdom-suspended" in current
if suspended != (desired_status == "suspended"):
drift.append("status:mismatch")
return drift
def _gql(self, token: str, query: str, variables: dict[str, Any], timeout: float = 15) -> dict:
request = Request(
self.base_url + "/api/graphql",
data=json.dumps({"query": query, "variables": variables}).encode(),
headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"},
method="POST",
)
try:
with self.opener(request, timeout=timeout) as response:
payload = json.loads(response.read())
except HTTPError as exc:
status = int(getattr(exc, "code", 0) or 0)
_discard(exc)
if status in {401, 403}:
raise DependencyFailure("auth_rejected") from None
raise DependencyFailure("protocol_error") from None
except (TimeoutError, URLError, OSError):
raise DependencyFailure("unreachable") from None
except (json.JSONDecodeError, TypeError, ValueError, UnicodeDecodeError):
raise DependencyFailure("protocol_error") from None
if payload.get("errors"):
raise ValueError(str(payload["errors"][0].get("message", "LLDAP GraphQL error")))
return dict(payload.get("data") or {})
def directory_login(
*,
base_url: str,
admin_password: str,
opener: Callable = urlopen,
timeout: float = 10,
) -> str:
request = Request(
base_url.rstrip("/") + "/auth/simple/login",
data=json.dumps({"username": "admin", "password": admin_password}).encode(),
headers={"Content-Type": "application/json"},
method="POST",
)
try:
with opener(request, timeout=timeout) as response:
payload = json.loads(response.read())
except HTTPError as exc:
status = int(getattr(exc, "code", 0) or 0)
_discard(exc)
if status in {401, 403}:
raise DependencyFailure("auth_rejected") from None
raise DependencyFailure("protocol_error") from None
except (TimeoutError, URLError, OSError):
raise DependencyFailure("unreachable") from None
except (json.JSONDecodeError, TypeError, ValueError, UnicodeDecodeError):
raise DependencyFailure("protocol_error") from None
if not isinstance(payload, dict):
raise DependencyFailure("protocol_error")
token = str(payload.get("token") or "")
if not token:
raise DependencyFailure("protocol_error")
return token
def _discard(exc: BaseException) -> None:
read = getattr(exc, "read", None)
if callable(read):
try:
read(65536)
except Exception:
pass
def dispatch(
provisioner: LLDAPProvisioner, path: str, payload: dict[str, Any]
) -> Result | DriftResult:
if path == "/v1/identities/tenant-access":
return provisioner.tenant_access(payload)
if path == "/v1/identities/provision":
return provisioner.provision(payload)
if path == "/v1/identities/drift":
return provisioner.drift(payload)
if path == "/v1/identities/reconcile":
return provisioner.reconcile(payload)
_required(payload, "external_subject", "idempotency_key", "correlation_id")
subject = _directory_username(str(payload["external_subject"]))
if path == "/v1/identities/suspend":
return provisioner.suspend(subject)
if path == "/v1/identities/reactivate":
return provisioner.reactivate(subject)
if path == "/v1/identities/deprovision":
return provisioner.deprovision(subject)
raise KeyError(path)
def _username(email: str, preferred: object = None) -> str:
if preferred is not None:
value = str(preferred).strip().lower()
if not re.fullmatch(r"[a-z][a-z0-9._-]{2,31}", value):
raise ValueError("preferred_username is invalid")
if value in {"admin", "administrator", "platform-root", "root", "system"}:
raise ValueError("preferred_username is reserved")
return value
local = email.partition("@")[0].lower()
value = re.sub(r"[^a-z0-9._-]+", "-", local).strip("-")
if not value or "@" not in email:
raise ValueError("valid primary_email is required")
return value[:64]
def _required(payload: dict[str, Any], *fields: str) -> None:
missing = [field for field in fields if not payload.get(field)]
if missing:
raise ValueError("missing required fields: " + ", ".join(missing))
if len(str(payload.get("idempotency_key", ""))) < 16:
raise ValueError("idempotency_key must contain at least 16 characters")
def _desired(payload: dict[str, Any]) -> tuple[str, str, set[str], str]:
_required(
payload,
"external_subject",
"tenant",
"primary_email",
"idempotency_key",
"correlation_id",
)
subject = str(payload["external_subject"])
email = str(payload["primary_email"]).strip().lower()
if _username(email) != subject:
raise ValueError("external_subject does not match canonical email username")
desired_status = str(payload.get("desired_status", "active"))
if desired_status not in {"active", "suspended"}:
raise ValueError("desired_status must be active or suspended")
tenant = str(payload["tenant"])
groups = {f"{tenant}:users"}
if "tenant-admin" in {str(role) for role in payload.get("roles", ())}:
groups.add(f"{tenant}:admins")
if desired_status == "suspended":
groups.add("netkingdom-suspended")
return subject, email, groups, desired_status
def _directory_username(subject: str) -> str:
return subject[4:].split(",", 1)[0] if subject.startswith("uid=") else subject
def _oidc_subject(username: str) -> str:
return f"uid={username},ou=people,dc=netkingdom,dc=local"