-
Notifications
You must be signed in to change notification settings - Fork 27
Expand file tree
/
Copy pathdeployment_guard.py
More file actions
114 lines (99 loc) · 3.83 KB
/
Copy pathdeployment_guard.py
File metadata and controls
114 lines (99 loc) · 3.83 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
"""Host-side admission adapter shared by CLI, boot and node-agent operations."""
import os
from contextlib import contextmanager
from contextvars import ContextVar
from pathlib import Path
import requests
import yaml
from dotenv import dotenv_values
from service_deployments import PlacementError, Plan
_held = ContextVar("deployment_activations", default=frozenset())
def placement_config(root: Path) -> dict:
path = root / "config" / "config.yml"
if not path.exists():
return {}
config = yaml.safe_load(path.read_text()) or {}
return config.get("service_placement") or {}
def control_token(root: Path) -> str:
return (
os.getenv("SERVICE_PLACEMENT_TOKEN")
or dotenv_values(root / "backend" / ".env").get("SERVICE_PLACEMENT_TOKEN")
or ""
)
def control_headers(root: Path) -> dict:
token = control_token(root)
return {"Authorization": f"Bearer {token}"} if token else {}
def local_services(root: Path, services: list[str]) -> list[str]:
"""Select bulk activation targets; admission still fences each activation.
Groups absent from the plan remain node-local. Both primary and warm standby
instances are eligible. An unreadable plan must never become permission to
start every enabled service.
"""
cfg = placement_config(root)
if not cfg:
return services
url = cfg.get("coordinator_url", "").rstrip("/")
node = cfg.get("node_id")
if not url or not node:
raise PlacementError("service_placement requires coordinator_url and node_id")
try:
response = requests.get(
f"{url}/deployments", headers=control_headers(root), timeout=15
)
response.raise_for_status()
plan = Plan.model_validate(response.json()["plan"])
except (requests.RequestException, ValueError, KeyError, TypeError) as exc:
raise PlacementError(
f"Cannot read deployment plan; refusing bulk activation: {exc}"
) from exc
if node not in plan.nodes:
raise PlacementError(
f"Register node {node} in the deployment plan before starting services"
)
return [
name
for name in services
if name not in plan.deployments
or any(instance.node == node for instance in plan.deployments[name].instances)
]
@contextmanager
def activation(root: Path, service: str, *, stop_revision: int | None = None):
cfg = placement_config(root)
if not cfg or service in _held.get():
yield
return
url = cfg.get("coordinator_url", "").rstrip("/")
node = cfg.get("node_id")
if not url or not node:
raise PlacementError("service_placement requires coordinator_url and node_id")
try:
response = requests.post(
f"{url}/deployments/activate",
json={"service": service, "node": node, "stop_revision": stop_revision},
headers=control_headers(root),
timeout=60,
)
if not response.ok:
raise PlacementError(response.json().get("detail", response.text))
token = response.json().get("token")
except requests.RequestException as exc:
raise PlacementError(
f"Deployment authority unavailable; refusing activation: {exc}"
) from exc
marker = _held.set(_held.get() | {service})
try:
yield
finally:
_held.reset(marker)
if token:
try:
response = requests.delete(
f"{url}/deployments/activations/{token}",
headers=control_headers(root),
timeout=15,
)
response.raise_for_status()
except requests.RequestException as exc:
raise PlacementError(
f"Activation completed but reservation {token} could not be released: {exc}"
) from exc