api+ui: kind-scoped save/delete endpoints and links
API: /config/{services,jobs,tools,static}/{name} now pin the twin they
target — _save_deployment/_delete_deployment take an explicit kind so a
patch to a 'backup' service can't bleed into a 'backup' job/tool. Add
/tools and /static endpoints; keep /deployments/{name} kind-agnostic.
New test_kind_twins proves per-kind save/delete isolation on disk.
UI: ConfigPanel and CreateDeploymentForm write to the kind-scoped
resource; GatewayPanel/NodeDetail/DeploymentsSection link via a shared
detailPath(name, kind) helper instead of the ambiguous /deployment/:name.
Includes incidental ruff-format reflow of untouched api files.
This commit is contained in:
@@ -279,10 +279,15 @@ async def delete_program(name: str, cascade: bool = False) -> dict:
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return {"ok": True, "program": name, "action": "deleted", "removed_deployments": removed}
|
||||
return {
|
||||
"ok": True,
|
||||
"program": name,
|
||||
"action": "deleted",
|
||||
"removed_deployments": removed,
|
||||
}
|
||||
|
||||
|
||||
def _save_deployment(name: str, config_dict: dict) -> dict:
|
||||
def _save_deployment(name: str, config_dict: dict, kind: str | None = None) -> dict:
|
||||
"""Create/update a deployment (any manager) with PATCH semantics.
|
||||
|
||||
The incoming config is shallow-merged over the existing spec, so a save can
|
||||
@@ -290,18 +295,26 @@ def _save_deployment(name: str, config_dict: dict) -> dict:
|
||||
a present key replaces wholesale, an **omitted** key is preserved, and an
|
||||
explicit ``null`` clears the key (back to its default). On CREATE there's no
|
||||
base, so the incoming config stands alone.
|
||||
|
||||
``kind`` pins the twin this save targets — a kind-scoped endpoint
|
||||
(``/services|/jobs|/tools|/static``) passes it so a partial patch to a
|
||||
``backup`` service can never bleed into a ``backup`` job/tool sharing the
|
||||
name. The kind-agnostic ``/deployments/{name}`` leaves it None and infers.
|
||||
"""
|
||||
_require_repo()
|
||||
config = get_config()
|
||||
incoming = dict(config_dict)
|
||||
|
||||
# Resolve the (name, kind) this save targets. A partial patch (e.g. just
|
||||
# Resolve the (name, kind) this save targets. An explicit kind is
|
||||
# authoritative (kind-scoped endpoint). Otherwise: a partial patch (e.g. just
|
||||
# {reach: off}) has no manager, so we can't derive kind from it — prefer the
|
||||
# existing same-named deployment when there's exactly one; otherwise derive the
|
||||
# existing same-named deployment when there's exactly one; else derive the
|
||||
# kind from the incoming spec (a create, or disambiguating a shared name).
|
||||
named = config.deployments_named(name)
|
||||
existing = None
|
||||
if len(named) == 1:
|
||||
if kind is not None:
|
||||
existing = config.deployment(kind, name)
|
||||
elif len(named) == 1:
|
||||
existing = named[0][1]
|
||||
else:
|
||||
try:
|
||||
@@ -332,17 +345,25 @@ def _save_deployment(name: str, config_dict: dict) -> dict:
|
||||
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
|
||||
detail=f"Invalid deployment config: {e}",
|
||||
)
|
||||
config.store_for(kind_for(dep))[name] = dep
|
||||
target_kind = kind_for(dep)
|
||||
# A field edit that changes the derived kind (e.g. adds a schedule) moves the
|
||||
# spec to the new store; drop the stale entry under the requested kind.
|
||||
if kind is not None and target_kind != kind:
|
||||
config.store_for(kind).pop(name, None)
|
||||
config.store_for(target_kind)[name] = dep
|
||||
save_config(config)
|
||||
return {"ok": True, "deployment": name}
|
||||
|
||||
|
||||
def _delete_deployment(name: str) -> dict:
|
||||
def _delete_deployment(name: str, kind: str | None = None) -> dict:
|
||||
"""Remove a deployment. A kind-scoped delete drops only that twin; the
|
||||
kind-agnostic path removes every kind sharing the name."""
|
||||
config = get_config()
|
||||
removed = False
|
||||
for kind in KINDS:
|
||||
if name in config.store_for(kind):
|
||||
del config.store_for(kind)[name]
|
||||
kinds = (kind,) if kind is not None else KINDS
|
||||
for k in kinds:
|
||||
if name in config.store_for(k):
|
||||
del config.store_for(k)[name]
|
||||
removed = True
|
||||
if not removed:
|
||||
raise HTTPException(
|
||||
@@ -391,26 +412,50 @@ def set_deployment_enabled(name: str, request: EnabledRequest) -> dict:
|
||||
return {"ok": True, "deployment": name, "enabled": request.enabled}
|
||||
|
||||
|
||||
# Kind-scoped endpoints — pin the twin so a save/delete can't hit a same-named
|
||||
# deployment of another kind (a `backup` service vs job vs tool).
|
||||
@router.put("/services/{name}")
|
||||
def save_service(name: str, request: ServiceConfigRequest) -> dict:
|
||||
"""Alias of PUT /deployments/{name} (kept for the existing dashboard)."""
|
||||
return _save_deployment(name, request.config)
|
||||
"""Create/update the *service* named `name`."""
|
||||
return _save_deployment(name, request.config, kind="service")
|
||||
|
||||
|
||||
@router.delete("/services/{name}")
|
||||
def delete_service(name: str) -> dict:
|
||||
return _delete_deployment(name)
|
||||
return _delete_deployment(name, kind="service")
|
||||
|
||||
|
||||
@router.put("/jobs/{name}")
|
||||
def save_job(name: str, request: JobConfigRequest) -> dict:
|
||||
"""Alias of PUT /deployments/{name} (kept for the existing dashboard)."""
|
||||
return _save_deployment(name, request.config)
|
||||
"""Create/update the *job* named `name`."""
|
||||
return _save_deployment(name, request.config, kind="job")
|
||||
|
||||
|
||||
@router.delete("/jobs/{name}")
|
||||
def delete_job(name: str) -> dict:
|
||||
return _delete_deployment(name)
|
||||
return _delete_deployment(name, kind="job")
|
||||
|
||||
|
||||
@router.put("/tools/{name}")
|
||||
def save_tool(name: str, request: ServiceConfigRequest) -> dict:
|
||||
"""Create/update the *tool* named `name`."""
|
||||
return _save_deployment(name, request.config, kind="tool")
|
||||
|
||||
|
||||
@router.delete("/tools/{name}")
|
||||
def delete_tool(name: str) -> dict:
|
||||
return _delete_deployment(name, kind="tool")
|
||||
|
||||
|
||||
@router.put("/static/{name}")
|
||||
def save_static(name: str, request: ServiceConfigRequest) -> dict:
|
||||
"""Create/update the *static* frontend named `name`."""
|
||||
return _save_deployment(name, request.config, kind="static")
|
||||
|
||||
|
||||
@router.delete("/static/{name}")
|
||||
def delete_static(name: str) -> dict:
|
||||
return _delete_deployment(name, kind="static")
|
||||
|
||||
|
||||
@router.post("/apply", response_model=ApplyResponse)
|
||||
|
||||
@@ -60,7 +60,10 @@ async def _check_systemd(name: str) -> HealthStatus:
|
||||
"""Check a managed service's health via its systemd unit state."""
|
||||
unit = f"castle-{name}.service"
|
||||
proc = await asyncio.create_subprocess_exec(
|
||||
"systemctl", "--user", "is-active", unit,
|
||||
"systemctl",
|
||||
"--user",
|
||||
"is-active",
|
||||
unit,
|
||||
stdout=asyncio.subprocess.PIPE,
|
||||
stderr=asyncio.subprocess.PIPE,
|
||||
)
|
||||
|
||||
@@ -31,7 +31,11 @@ async def get_logs(
|
||||
config = load_config(root)
|
||||
# A name may span kinds — the managed (systemd) one owns the journal.
|
||||
dep_kind = next(
|
||||
((k, s) for k, s in config.deployments_named(name) if getattr(s, "manage", None)),
|
||||
(
|
||||
(k, s)
|
||||
for k, s in config.deployments_named(name)
|
||||
if getattr(s, "manage", None)
|
||||
),
|
||||
None,
|
||||
)
|
||||
if dep_kind is None:
|
||||
|
||||
@@ -35,7 +35,9 @@ class CastleMDNS:
|
||||
self._service_info: ServiceInfo | None = None
|
||||
|
||||
# Discovered state
|
||||
self.peers: dict[str, dict] = {} # hostname -> {gateway_port, api_port, addresses}
|
||||
self.peers: dict[
|
||||
str, dict
|
||||
] = {} # hostname -> {gateway_port, api_port, addresses}
|
||||
self.mqtt_broker: dict | None = None # {host, port} or None
|
||||
|
||||
def _on_service_state_change(
|
||||
@@ -67,14 +69,18 @@ class CastleMDNS:
|
||||
def _handle_castle_peer(self, info: ServiceInfo) -> None:
|
||||
"""Process a discovered castle peer."""
|
||||
props = {
|
||||
k.decode() if isinstance(k, bytes) else k: v.decode() if isinstance(v, bytes) else v
|
||||
k.decode() if isinstance(k, bytes) else k: v.decode()
|
||||
if isinstance(v, bytes)
|
||||
else v
|
||||
for k, v in info.properties.items()
|
||||
}
|
||||
peer_hostname = props.get("hostname", "")
|
||||
if not peer_hostname or peer_hostname == self._hostname:
|
||||
return
|
||||
|
||||
addresses = [socket.inet_ntoa(addr) for addr in info.addresses if len(addr) == 4]
|
||||
addresses = [
|
||||
socket.inet_ntoa(addr) for addr in info.addresses if len(addr) == 4
|
||||
]
|
||||
|
||||
self.peers[peer_hostname] = {
|
||||
"gateway_port": int(props.get("gateway_port", 9000)),
|
||||
@@ -85,13 +91,17 @@ class CastleMDNS:
|
||||
|
||||
def _handle_mqtt_broker(self, info: ServiceInfo) -> None:
|
||||
"""Process a discovered MQTT broker."""
|
||||
addresses = [socket.inet_ntoa(addr) for addr in info.addresses if len(addr) == 4]
|
||||
addresses = [
|
||||
socket.inet_ntoa(addr) for addr in info.addresses if len(addr) == 4
|
||||
]
|
||||
if addresses:
|
||||
self.mqtt_broker = {
|
||||
"host": addresses[0],
|
||||
"port": info.port,
|
||||
}
|
||||
logger.info("mDNS: discovered MQTT broker at %s:%d", addresses[0], info.port)
|
||||
logger.info(
|
||||
"mDNS: discovered MQTT broker at %s:%d", addresses[0], info.port
|
||||
)
|
||||
|
||||
def start(self) -> None:
|
||||
"""Start advertising and browsing."""
|
||||
@@ -109,14 +119,24 @@ class CastleMDNS:
|
||||
},
|
||||
)
|
||||
self._zeroconf.register_service(self._service_info)
|
||||
logger.info("mDNS: advertising %s on port %d", self._hostname, self._gateway_port)
|
||||
logger.info(
|
||||
"mDNS: advertising %s on port %d", self._hostname, self._gateway_port
|
||||
)
|
||||
|
||||
# Browse for peers and MQTT broker
|
||||
self._browsers.append(
|
||||
ServiceBrowser(self._zeroconf, CASTLE_SERVICE_TYPE, handlers=[self._on_service_state_change])
|
||||
ServiceBrowser(
|
||||
self._zeroconf,
|
||||
CASTLE_SERVICE_TYPE,
|
||||
handlers=[self._on_service_state_change],
|
||||
)
|
||||
)
|
||||
self._browsers.append(
|
||||
ServiceBrowser(self._zeroconf, MQTT_SERVICE_TYPE, handlers=[self._on_service_state_change])
|
||||
ServiceBrowser(
|
||||
self._zeroconf,
|
||||
MQTT_SERVICE_TYPE,
|
||||
handlers=[self._on_service_state_change],
|
||||
)
|
||||
)
|
||||
|
||||
def stop(self) -> None:
|
||||
|
||||
@@ -40,7 +40,9 @@ class MeshStateManager:
|
||||
def update_node(self, hostname: str, registry: NodeRegistry) -> None:
|
||||
"""Add or update a remote node's registry."""
|
||||
self._nodes[hostname] = RemoteNode(registry=registry)
|
||||
logger.info("Mesh: updated node %s (%d deployed)", hostname, len(registry.deployed))
|
||||
logger.info(
|
||||
"Mesh: updated node %s (%d deployed)", hostname, len(registry.deployed)
|
||||
)
|
||||
|
||||
def set_offline(self, hostname: str) -> None:
|
||||
"""Mark a node as offline (LWT received)."""
|
||||
|
||||
@@ -171,7 +171,9 @@ class CastleMQTTClient:
|
||||
return
|
||||
|
||||
self._connected = True
|
||||
logger.info("Connected to MQTT broker at %s:%d", self._broker_host, self._broker_port)
|
||||
logger.info(
|
||||
"Connected to MQTT broker at %s:%d", self._broker_host, self._broker_port
|
||||
)
|
||||
|
||||
# Publish our status as online (retained)
|
||||
client.publish(
|
||||
@@ -222,7 +224,9 @@ class CastleMQTTClient:
|
||||
if payload == "offline":
|
||||
mesh_state.set_offline(hostname)
|
||||
asyncio.run_coroutine_threadsafe(
|
||||
broadcast("mesh", {"event": "node_offline", "hostname": hostname}),
|
||||
broadcast(
|
||||
"mesh", {"event": "node_offline", "hostname": hostname}
|
||||
),
|
||||
self._loop,
|
||||
)
|
||||
|
||||
|
||||
@@ -154,7 +154,9 @@ def _summary_from_service(
|
||||
)
|
||||
|
||||
|
||||
def _summary_from_job(name: str, job: SystemdDeployment, config: object) -> DeploymentSummary:
|
||||
def _summary_from_job(
|
||||
name: str, job: SystemdDeployment, config: object
|
||||
) -> DeploymentSummary:
|
||||
"""Build a DeploymentSummary from a systemd deployment (job, non-deployed)."""
|
||||
managed = bool(job.manage and job.manage.systemd and job.manage.systemd.enable)
|
||||
|
||||
@@ -280,7 +282,9 @@ def _service_from_deployed(name: str, deployed: object) -> ServiceSummary:
|
||||
)
|
||||
|
||||
|
||||
def _service_from_spec(name: str, svc: SystemdDeployment, config: object) -> ServiceSummary:
|
||||
def _service_from_spec(
|
||||
name: str, svc: SystemdDeployment, config: object
|
||||
) -> ServiceSummary:
|
||||
"""Build a ServiceSummary from a systemd deployment."""
|
||||
port = None
|
||||
health_path = None
|
||||
@@ -911,7 +915,8 @@ def get_gateway() -> GatewayInfo:
|
||||
# is not a service, so filtering to services dropped its public_url (calculator).
|
||||
public_domain = registry.node.public_domain
|
||||
public_names = {
|
||||
name for _k, name, dep in (config.all_deployments() if config else [])
|
||||
name
|
||||
for _k, name, dep in (config.all_deployments() if config else [])
|
||||
if getattr(dep, "public", False)
|
||||
}
|
||||
|
||||
@@ -937,7 +942,8 @@ def get_gateway() -> GatewayInfo:
|
||||
tunnel_connected = (
|
||||
subprocess.run(
|
||||
["systemctl", "--user", "is-active", "castle-castle-tunnel.service"],
|
||||
capture_output=True, text=True,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
).stdout.strip()
|
||||
== "active"
|
||||
)
|
||||
@@ -970,7 +976,7 @@ def save_gateway_config(request: GatewayConfigRequest) -> dict[str, str]:
|
||||
from castle_core.config import load_config, save_config
|
||||
|
||||
config = load_config(root)
|
||||
norm = lambda v: (v or None) # noqa: E731 — empty string clears
|
||||
norm = lambda v: v or None # noqa: E731 — empty string clears
|
||||
config.gateway.tls = norm(request.tls)
|
||||
config.gateway.domain = norm(request.domain)
|
||||
config.gateway.public_domain = norm(request.public_domain)
|
||||
|
||||
@@ -107,7 +107,9 @@ async def _deferred_systemctl(action: str, unit: str, delay: float = 0.5) -> Non
|
||||
async def _do_action(name: str, action: str) -> JSONResponse:
|
||||
"""Execute a systemctl action and broadcast updated health."""
|
||||
deployed = _managed(name)
|
||||
unit = unit_name(name, deployed.kind) if deployed else f"{UNIT_PREFIX}{name}.service"
|
||||
unit = (
|
||||
unit_name(name, deployed.kind) if deployed else f"{UNIT_PREFIX}{name}.service"
|
||||
)
|
||||
|
||||
# Self-restart: defer the systemctl call so the response can be sent first
|
||||
if name == SELF_NAME and action in ("restart", "stop"):
|
||||
|
||||
Reference in New Issue
Block a user