Repository navigation
Expand file tree
/
Copy pathfull_stack.py
More file actions
1772 lines (1578 loc) · 76.5 KB
/
Copy pathfull_stack.py
File metadata and controls
1772 lines (1578 loc) · 76.5 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
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
"""
Build the whole stack, bare machines to a served model, in one run.
This is the composition the per-component examples only hint at:
RKE2 a cluster on bare machines
-> Storage block storage (Longhorn), which RKE2 ships without
-> MinIO object storage, on top of that block storage
-> Inference a model served from that object storage (vLLM)
Each layer's output is the next one's input. The kubeconfig threads
through all four; Longhorn provides the StorageClass MinIO's PVCs bind
to; MinIO's endpoint is where vLLM reads its weights — so the inference
pod never needs internet access at all.
pip install -e . # from the repo root
python3 examples/full_stack.py
DESTRUCTIVE. Stage 1 installs RKE2 on every machine in NODES, so point
it at hosts you are willing to have reformatted. Needs `ssh`, `helm` and
`kubectl` on PATH, passwordless key-based SSH to every node, and a user
that is root or has non-interactive `sudo -n`.
Re-running is safe: every stage is idempotent. Trim STAGES to re-run
just part of it — the later stages read what the earlier ones left
behind rather than depending on values passed in memory.
Not everything here is SDK surface. Bucket and object management, node
labels and Secrets are deliberately outside the components, so this
file does them with `kubectl` — from
`multistack.kube`, which is where the mechanics live, rather than
reimplemented here. Those parts are marked `-- glue --`, and they are the
honest cost of the composition: the SDK gets you four layers, and the
decisions between them are yours.
"""
import base64
import json
import os
import secrets
import sys
from functools import lru_cache, partial
from urllib.parse import quote
from multistack import (
Billing,
BillingBackend,
Cache,
CacheBackend,
ControlPlane,
ControlPlaneBackend,
CNPGOptions,
Database,
DatabaseBackend,
DatabaseConfig,
Enricher,
EnricherBackend,
EnricherOptions,
Gateway,
GatewayBackend,
Inference,
InferenceBackend,
MinIOTenant,
NatsQueueOptions,
KubePrometheusStackOptions,
Observability,
ObservabilityBackend,
Policy,
PolicyBackend,
Portal,
PortalBackend,
Queue,
QueueBackend,
RKE2Cluster,
Route,
RouteBackend,
RKE2Node,
Stack,
Storage,
StorageBackend,
Tokenizer,
TokenizerBackend,
prompt_secret,
)
from multistack.ingress_gateway import IngressGateway, IngressGatewayBackend
from multistack.ingress_gateway.spec import MetalLBIstioOptions
from multistack.accelerator import (
Accelerator,
AcceleratorBackend,
NODE_FEATURE_LABEL,
)
from multistack.kube import apply, node_name_for, wait_for_job
from multistack.kube import kubectl as _kubectl
from multistack.gateway import ModelGatewayOptions
from multistack.cache import ValkeyOptions
from multistack.policy.spec import RateLimits
from multistack.storage import LonghornOptions
from multistack.backends.minio_client import MinIOBackend
from multistack.backends.rke2_client import RKE2Backend
# ---------------------------------------------------------------- config --
# Your lab, from the environment, so this file needs no editing. The
# defaults are RFC 5737 documentation addresses and a documentation user:
# that range is reserved for exactly this purpose, so nothing here can be
# mistaken for a real host.
SSH_USER = os.environ.get("MULTISTACK_SSH_USER", "ubuntu")
SSH_KEY = os.environ.get("MULTISTACK_SSH_KEY", "~/.ssh/id_ed25519")
# Not /tmp: that does not survive a reboot, and losing the kubeconfig
# leaves you with a cluster you cannot reach and no obvious reason why.
KUBECONFIG = os.path.expanduser(
os.environ.get("MULTISTACK_KUBECONFIG", "~/.multistack/kubeconfig")
)
# The machines. One server, any number of agents, comma-separated.
SERVER = os.environ.get("MULTISTACK_SERVER", "192.0.2.10")
AGENTS = [a.strip() for a in os.environ.get(
"MULTISTACK_AGENTS", "192.0.2.11,192.0.2.12,192.0.2.13").split(",") if a.strip()]
# The machine with the accelerator, by IP like every other node here.
# Unset skips the GPU stage entirely, which is the right default: most
# labs have no card, and a cluster without one should not fail to build.
# Not INFERENCE_NODE's default either — on this lab they are different
# machines, because the GPU is a compute capability 6.1 card that vLLM's
# PyTorch build cannot use (see stage_inference's own note below).
GPU_NODE = os.environ.get("MULTISTACK_GPU_NODE", "")
# The GPU node's own SSH credential, when it differs from every other
# node's. Unset falls back to SSH_USER/SSH_KEY, so a lab where every node
# shares one credential -- the common case -- needs neither of these set.
# This lab's own GPU node needs both: a different user and a different
# key than the other three, on a different subnet
# (scratch/rebuild/README.md, where this was worked out the first time
# this cluster was built from bare machines).
GPU_SSH_USER = os.environ.get("MULTISTACK_GPU_SSH_USER", SSH_USER)
GPU_SSH_KEY = os.environ.get("MULTISTACK_GPU_SSH_KEY", SSH_KEY)
def _node(address: str, role: str) -> RKE2Node:
"""-- glue -- the one address that may need GPU_SSH_USER/GPU_SSH_KEY
instead of the uniform SSH_USER/SSH_KEY every other node gets. Not a
general per-node override table: one address needing one different
(user, key) pair is the shape this lab actually has. RKE2Node itself
already carries a `user`/`ssh_key` per node -- reach for that
directly if a lab needs more than this.
"""
if address and address == GPU_NODE:
return RKE2Node(address=address, user=GPU_SSH_USER, role=role,
ssh_key=GPU_SSH_KEY)
return RKE2Node(address=address, user=SSH_USER, role=role, ssh_key=SSH_KEY)
NODES = [
_node(SERVER, "server"),
*[_node(a, "agent") for a in AGENTS],
]
# The node inference is pinned to. Labelled and tainted below so nothing
# else schedules there and vLLM gets the whole machine's CPU. Not the GPU
# node: its card is compute capability 6.1, below the 7.0 floor of every
# PyTorch build vLLM ships, so vLLM there fails with "no kernel image is
# available for execution on the device". CPU on a plain worker is the
# working configuration.
INFERENCE_NODE = os.environ.get("MULTISTACK_INFERENCE_NODE", AGENTS[0])
# The first-party images live in a private GHCR namespace, so every
# namespace that runs one needs its own copy of the pull credential --
# an imagePullSecret is namespace-scoped. Unset skips it entirely, which
# is right for a public package or a mirror that needs no auth.
#
# From the environment rather than a flag, and never printed: a token on
# a command line is readable from /proc, and one in shell history is
# readable forever. A file read into the environment is better still.
GHCR_HOST = os.environ.get("MULTISTACK_GHCR_HOST", "ghcr.io")
GHCR_USER = os.environ.get("MULTISTACK_GHCR_USER", "")
GHCR_TOKEN = os.environ.get("MULTISTACK_GHCR_TOKEN", "")
PULL_SECRET = os.environ.get("MULTISTACK_PULL_SECRET", "ghcr-creds")
# The LAN range MetalLB may hand out, and the node it announces from.
# Unset skips the front door entirely, same as the GPU stage: a lab with
# no spare addresses on its subnet should still build.
#
# Layer-2 announcement means the node answers ARP for an address that is
# not configured on any of its interfaces. On a cloud network that is
# usually filtered — on OpenStack the pool has to be added to
# `allowed_address_pairs` on this node's port, or Neutron drops the ARP
# and the address is assigned, healthy, and unreachable.
ADDRESS_POOL = [a for a in os.environ.get(
"MULTISTACK_ADDRESS_POOL", "").split(",") if a.strip()]
INGRESS_NODE = os.environ.get("MULTISTACK_INGRESS_NODE", AGENTS[-1] if AGENTS else "")
MODEL = "Qwen/Qwen2.5-0.5B-Instruct"
BUCKET = "models"
MODEL_KEY = MODEL.split("/")[-1] # Qwen2.5-0.5B-Instruct
CREDS_SECRET = "minio-creds"
INFERENCE_NS = "inference"
# The platform layers each land in the namespace their own spec already
# declares — Tokenizer and Policy in `policy`, Gateway in `gateway`, both
# control planes in `control-plane`, both portals in `frontend`, the
# database in `postgres`. So none of those is named here and no stage below passes
# `namespace=`: an example that restates a default teaches a layout the
# SDK does not use, and what this file taught until now was to put every
# service into one flat `platform` namespace.
#
# Where a Secret has to be applied into the same namespace as the release
# that reads it, the namespace is read back off the spec class rather than
# typed out again, so there is still exactly one source for each.
#
# The cache is the deliberate exception. `Cache.DEFAULT_NAMESPACES` still
# says `valkey` -- the implementation's name -- but the platform groups it
# by function alongside everything else, and the VALKEY_URL in both
# control-plane Secrets says `cache`. So this overrides the default rather
# than restating it, and the two want reconciling.
CACHE_NS = "cache"
# One database per service, each its own PostgreSQL cluster: every one of
# these owns its schema and migrates it independently, so none can share
# one. Owner names are part of the contract — they appear in the
# DATABASE_URL each service authenticates with — so they are stated here
# rather than derived. Not just the two control planes: billing is
# database-per-service too (Billing.REQUIRES names "database" the same
# way ControlPlane does), so its row lives here rather than being
# special-cased in stage_billing.
DATABASES = (
# component, cluster, database, owner
("admin", "admin-pg", "admin_control_plane", "admin"),
("organization", "org-pg", "org_control_plane", "orguser"),
("billing", "billing-pg", "billing", "billing"),
)
# The Secret billing's chart reads DATABASE_URL, SERVICE_API_KEY and
# (optionally) the two Stripe credentials from.
BILLING_SECRET = "billing-secrets"
# There is no in-cluster registry, so the portal images are built locally
# and sideloaded into one node's containerd. A pod scheduled anywhere else
# stays ImagePullBackOff, which is why this pin is not optional.
IMAGE_NODE = os.environ.get("MULTISTACK_IMAGE_NODE", AGENTS[-1] if AGENTS else "")
# Gateway API is a set of CRDs, and neither RKE2 nor any chart this SDK
# deploys installs them -- multistack/route/README.md says so, and the
# route driver's check_prerequisites reports it. So the route stage
# applies them, pinned: `standard` is the channel holding Gateway and
# HTTPRoute as v1, and the version is the one the working cluster runs
# rather than whatever `latest` resolves to on the day of a rebuild.
#
# Fetched by the kubectl on *this* machine, not from inside the cluster,
# so the nodes need no internet for it.
GATEWAY_API_VERSION = os.environ.get("MULTISTACK_GATEWAY_API_VERSION", "v1.2.1")
GATEWAY_API_MANIFEST = (
"https://github.com/kubernetes-sigs/gateway-api/releases/download/"
f"{GATEWAY_API_VERSION}/standard-install.yaml"
)
# The hostnames each service answers on, behind the one front-door
# address. Hostname routing rather than path prefixes for the portals:
# each is served at `/` and its built asset paths assume it, so a prefix
# would return the SPA's own HTML with a 200 for every asset. The model
# gateway takes both -- a hostname *and* /v1 -- because an API client
# sends an explicit path and a browser does not.
#
# These need DNS (or /etc/hosts) pointing at the front door to be
# reachable; the routes attach either way, which is what makes them worth
# creating before the DNS exists.
ADMIN_HOST = os.environ.get("MULTISTACK_ADMIN_HOST", "admin.example.com")
ORG_HOST = os.environ.get("MULTISTACK_ORG_HOST", "org.example.com")
API_HOST = os.environ.get("MULTISTACK_API_HOST", "api.example.com")
# The queue stage below deploys the event backbone (NATS/JetStream) and
# publishes this automatically when "queue" is in STAGES.
# Set here only to point at an externally-managed one instead, or as the
# value main() falls back to when "queue" is trimmed out of a re-run.
# Left empty either way, the policy stage deploys an always-allow limiter
# deliberately rather than one that looks enforcing and is not.
EVENT_BACKBONE = os.environ.get("MULTISTACK_EVENT_BACKBONE", "")
# Which stages to run. Every stage after the first can run alone — the
# specs below are rebuilt either way, and each backend's create() adopts
# whatever is already live rather than replacing it. Overridable so a
# resumed run needs no edit to this file:
#
# MULTISTACK_STAGES=weights,inference python3 examples/full_stack.py
STAGES = os.environ.get(
"MULTISTACK_STAGES",
"cluster,accelerator,ingress,storage,objectstore,weights,inference,"
"cache,database,registry,tokenizer,queue,enricher,policy,gateway,"
"billing,controlplane,portal,route,observability"
).split(",")
# The order stages run in, and what to call them. One list rather than a
# number typed into each `print` below: inserting a stage used to mean
# renumbering every banner after it, and twice now the docstrings and the
# banners disagreed about which number a stage was.
STAGE_LABELS = [
("cluster", "RKE2 cluster"),
("accelerator", "GPU scheduling"),
("ingress", "Front door (MetalLB + Istio)"),
("storage", "Longhorn block storage"),
("objectstore", "MinIO object storage"),
("inference", "vLLM inference"),
("cache", "Valkey cache"),
("database", "CloudNativePG"),
("registry", "Image pull credentials"),
("tokenizer", "Tokenizer"),
("queue", "Event backbone / queue (JetStream)"),
("enricher", "Event enricher (token-count guarantee)"),
("policy", "Rate-limit policies"),
("gateway", "Gateway"),
("billing", "Usage metering / billing"),
("controlplane", "Control planes"),
("portal", "Web portals"),
("route", "Routes through the front door (API + both portals)"),
("observability", "Observability"),
]
# Which namespaces pull a first-party image, read off the specs rather
# than listed here -- the namespace regrouping already moved these once.
PULL_SECRET_NAMESPACES = sorted({
Tokenizer.DEFAULT_NAMESPACES["tiktoken"],
Gateway.DEFAULT_NAMESPACES["modelgateway"],
*Policy.DEFAULT_NAMESPACES.values(),
*Enricher.DEFAULT_NAMESPACES.values(),
*Billing.DEFAULT_NAMESPACES.values(),
*ControlPlane.DEFAULT_NAMESPACES.values(),
*Portal.DEFAULT_NAMESPACES.values(),
})
def _banner(stage: str) -> None:
"""`=== 3/15 Front door ===`, with the number derived rather than typed.
The stage docstrings below carry no number for the same reason:
STAGE_LABELS is the one place the order is stated, so inserting a
stage cannot leave a banner and a docstring disagreeing -- which is
exactly what happened twice before this existed.
`weights` has no banner of its own -- it is a Job the objectstore
stage's tenant feeds, not a layer -- so it is absent from the list
above and does not shift anything.
"""
names = [name for name, _ in STAGE_LABELS]
label = dict(STAGE_LABELS)[stage]
print(f"\n=== {names.index(stage) + 1}/{len(names)} {label} ===")
# The composition root. Which cluster is stated once, here, and threaded
# into every spec below — so `kubeconfig_path` is not repeated five times,
# and the StorageClass MinIO binds against comes from the Storage spec
# that produced it rather than a hardcoded "longhorn" that goes quietly
# wrong the day the implementation changes.
#
# Not ambient: nothing is read from the environment, and stack.outputs
# shows exactly what will be filled in. A spec built directly still works
# — every other example does that.
stack = Stack(kubeconfig_path=KUBECONFIG)
# ------------------------------------------------------------------ glue --
# Bound to this cluster once, so no call below can forget the kubeconfig.
# These are the SDK's own: running kubectl, applying a manifest, waiting on
# a Job and mapping a node IP to its Kubernetes name are mechanism every
# capability needs, so they live in multistack.kube rather than being
# rewritten per script. What stays glue here is the policy — which node to
# reserve, which bucket, what to call the Secret.
kubectl = partial(_kubectl, KUBECONFIG)
_apply = partial(apply, KUBECONFIG)
node_name = partial(node_name_for, KUBECONFIG)
def _authenticated(url: str, owner: str, password: str) -> str:
"""-- glue -- turn a published `database_url` into one a chart can use.
Two things that URL deliberately is not.
`+asyncpg`, because each control plane hands this value straight to
SQLAlchemy's create_async_engine() and has no scheme fix-up of its
own — a bare `postgresql://` is refused there as a sync driver, at
startup, before the first query.
And authenticated: `database_url` omits the password on purpose, so
that a spec is a file people can commit. Percent-encoded, because a
prompted password containing `@` or `/` would otherwise be read as
part of the host.
"""
return url.replace(
"postgresql://", "postgresql+asyncpg://", 1
).replace(f"{owner}@", f"{owner}:{quote(password, safe='')}@", 1)
CACHE_AUTH_SECRET = "valkey-auth"
CACHE_AUTH_KEY = "valkey-password"
@lru_cache(maxsize=None)
def _cache_password() -> str:
"""-- glue -- the Valkey password, asked for once per run.
Generated by default: nobody needs to choose this one, and the only
things that read it are the release itself and the four consumers
wired below. MULTISTACK_CACHE_PASSWORD keeps it stable across
separate runs, which is what re-running one stage on its own needs --
without it, a second run of `cache` would rotate the password and
leave the consumers from the first run authenticating with the old
one.
Cached for the same reason `_org_verify_key` is: three stages need
this value and a second `generate=True` prompt would hand them
different ones.
"""
return prompt_secret(
"Valkey password (the cache every rate-limit policy counts in)",
env_var="MULTISTACK_CACHE_PASSWORD",
generate=True,
)
def _cache_authenticated(url: str, password: str) -> str:
"""-- glue -- put the password into a published `cache_url`.
The same shape as `_authenticated` above and for the same reason:
`Cache.endpoint` omits the credential on purpose, so that a spec is a
file people can commit, and whatever must authenticate assembles it
here instead.
`redis://:password@host` -- Valkey has a password but no username in
its default ACL, so the user half is empty and the colon stays.
Percent-encoded, because a generated password containing `@` or `/`
would otherwise be read as part of the host.
"""
return url.replace("redis://", f"redis://:{quote(password, safe='')}@", 1)
@lru_cache(maxsize=None)
def _org_verify_key() -> str:
"""-- glue -- the one bearer the gateway presents to the org plane.
Two releases have to carry this exact same value: it is the
gateway's `ORG_VERIFY_API_KEY` and the organization control plane's
`MG_SERVICE_API_KEY`. ADR-026 removed the gateway's static-key auth
mode, so every /v1 request now asks that plane to resolve the
caller's bearer into a real per-org key — and if the two values
disagree the gateway fails closed, which reaches a caller as a 503
on every single request.
Deliberately not SERVICE_API_KEY, which the planes use on each
other: that bearer also opens the admin and org-to-org routes, so
sharing it would let a compromised gateway escalate into them. This
one resolves API keys and nothing else.
Cached because both stages below need it, and a second
`generate=True` prompt would hand them two different values — the
exact mismatch described above. MULTISTACK_ORG_VERIFY_API_KEY keeps
it stable across separate runs, which is what re-running one stage
on its own needs.
"""
return prompt_secret(
"Gateway's bearer against the organization control plane "
"(its ORG_VERIFY_API_KEY = that plane's MG_SERVICE_API_KEY)",
env_var="MULTISTACK_ORG_VERIFY_API_KEY",
generate=True,
)
@lru_cache(maxsize=None)
def _admin_service_api_key() -> str:
"""-- glue -- admin-control-plane's own SERVICE_API_KEY, reused by
organization-control-plane and billing.
One shared value, not per-service: admin-cp's own outbound calls
to organization-cp carry `Bearer {admin's SERVICE_API_KEY}`
(admin-control-plane/main.py's org_client), and organization-cp
checks incoming calls against its *own* configured SERVICE_API_KEY
(require_service_key) — so the two must be issued the same value,
or every admin<->org call 401s despite both planes installing
cleanly and passing their own health probes (neither exercises
this path). Billing's own `SERVICE_API_KEY` must equal it too: it
is the outbound bearer admin-cp reuses to notify billing of
org/plan changes (POST /internal/v1/subscriptions, ADR-032 — see
examples/billing/install.py). A mismatch there doesn't fail closed
the way the admin<->org one does; it just makes that best-effort
notify 401 silently on every plan change.
Generated, not prompted, the same way each control plane's own
SERVICE_API_KEY already is — this is a service-to-service credential
no person ever types. Cached because stage_controlplane (for both
planes) and stage_billing all need this exact value, and a second
`secrets.token_urlsafe` call would hand them different ones.
"""
return secrets.token_urlsafe(24)
def _ensure_namespace(namespace: str) -> None:
"""-- glue -- create a namespace a Secret is about to be applied into.
Helm creates the namespace it installs into, but these Secrets are
applied *before* the chart that reads them, so at that point nothing
has created it yet and the apply fails.
"""
_apply({
"apiVersion": "v1", "kind": "Namespace",
"metadata": {"name": namespace},
})
def _ensure_gateway_api() -> None:
"""-- glue -- install the Gateway API CRDs if they are not there.
Not a capability: these are cluster-scoped CRDs that several things
could own, and installing them is a decision about the cluster rather
than about a route. The SDK deliberately does not install them, and
the route driver reports their absence -- but reporting it at the
sixteenth stage of a rebuild, forty minutes in, is too late to be
useful, and this file is where the cluster-level decisions live.
Idempotent, and quiet when they already exist: `kubectl apply` on the
same bundle is a no-op, but skipping it entirely keeps a rebuild from
needing the internet for a cluster that is already complete.
"""
existing = kubectl(
"get", "crd", "httproutes.gateway.networking.k8s.io",
"--ignore-not-found", "-o", "name", check=False,
).strip()
if existing:
print("[stack] Gateway API CRDs already installed")
return
print(f"[stack] installing Gateway API {GATEWAY_API_VERSION} "
"(standard channel)")
kubectl("apply", "-f", GATEWAY_API_MANIFEST)
# The CRDs have to be established before a route can be applied
# against them, and apply returns as soon as the objects are created.
kubectl("wait", "--for=condition=Established", "--timeout=60s",
"crd/gateways.gateway.networking.k8s.io",
"crd/httproutes.gateway.networking.k8s.io")
print("[stack] Gateway API CRDs established")
# ----------------------------------------------------------------- stages --
def stage_cluster() -> str:
"""RKE2 on bare machines."""
cluster = stack.build(
RKE2Cluster,
name="ai-cluster",
version="v1.32.5+rke2r1",
nodes=NODES,
cni="cilium",
disable_kube_proxy=True,
# Multi-homed hosts otherwise register on whichever NIC the
# kubelet picks, which may not be the address declared above.
pin_node_ip=True,
# Cilium defaults to 1450. Where the real path MTU is lower, pod
# traffic to the affected node fails in ways that look like
# anything but MTU. Measured path MTU in this lab is 1428, so
# 1428 - 50 (VXLAN overhead) = 1378; 1350 is what the working
# cluster runs, kept here rather than re-deriving it.
cilium_mtu=1350,
)
RKE2Backend().create(cluster)
return cluster.kubeconfig_path
def stage_accelerator() -> str | None:
"""GPU scheduling — so that a pod can ask for a card at all.
Kubernetes does not know a node has a GPU. Until something
advertises one as an extended resource, `nvidia.com/gpu` is not a
thing a pod can request, so a GPU workload stays Pending on a
machine with a perfectly healthy card in it.
Returns the resource name a pod can now ask for, or None when no GPU
node was named.
"""
if not GPU_NODE:
print("[stack] skipped: no GPU node (set MULTISTACK_GPU_NODE to "
"the address of the machine with the card)")
return None
# -- glue -- the chart's own affinity requires a label that Node
# Feature Discovery would apply, and this cluster does not run NFD.
# Labelling a node is a change to the machine rather than to a
# release, which is why it sits here next to the inference stage's
# label and taint rather than inside the capability. The driver
# refuses to install without it, instead of leaving a DaemonSet that
# wants zero pods and reports no error.
name = node_name(GPU_NODE)
kubectl("label", "node", name, f"{NODE_FEATURE_LABEL}=true", "--overwrite")
print(f"[stack] labelled {name} {NODE_FEATURE_LABEL}=true")
accelerator = stack.build(
Accelerator,
# Pinned: on a cluster where one machine has the card, an
# unpinned DaemonSet is a pod on every node for the sake of one.
node_selector={"kubernetes.io/hostname": name},
)
resource = AcceleratorBackend().create(accelerator)
# Deliberately not recorded: this capability publishes nothing. It
# advertises a property of a node rather than an address, so there is
# no value for a later spec to fill a field from — see
# multistack/accelerator/spec.py.
print(f"[stack] {name} can now be asked for {resource}")
return resource
def stage_ingress() -> str | None:
"""The front door — one external address for the whole platform.
MetalLB hands a real LAN address to the Istio ingress gateway's
Service, because on bare metal no cloud provider does it. Everything
reachable from outside the cluster goes through this one address
rather than taking a second one each.
Returns the external address, or None when no pool was named.
"""
if not ADDRESS_POOL:
print("[stack] skipped: no address pool (set "
"MULTISTACK_ADDRESS_POOL=192.0.2.240-192.0.2.250)")
return None
ingress = stack.build(
IngressGateway,
address_pool=ADDRESS_POOL,
# Which nodes may announce the pool. Empty would mean every node,
# which is only right when every node sits on the pool's subnet.
node_selector=({"kubernetes.io/hostname": node_name(INGRESS_NODE)}
if INGRESS_NODE else {}),
# Same local-chart escape hatch storage's LonghornOptions.chart
# uses -- these repos' index.yaml lives on GitHub Pages/GCS, but
# the tarballs themselves are fetched from
# release-assets.githubusercontent.com, which is the flaky leg.
# Unset envs fall back to the chart name, i.e. today's default
# repo-based resolution.
options=MetalLBIstioOptions(
metallb_chart=os.environ.get("MULTISTACK_METALLB_CHART", "metallb"),
istio_base_chart=os.environ.get("MULTISTACK_ISTIO_BASE_CHART", "base"),
istiod_chart=os.environ.get("MULTISTACK_ISTIOD_CHART", "istiod"),
ingress_chart=os.environ.get("MULTISTACK_ISTIO_GATEWAY_CHART", "gateway"),
),
)
backend = IngressGatewayBackend()
for warning in backend.check_prerequisites(ingress):
print(f"[ingress] warning: {warning}")
endpoint = backend.create(ingress)
# publishes ingress_gateway_endpoint -- which works because the
# driver writes the assigned address back onto the spec
# (drivers/metallb_istio.py). record() skips a None, so the route
# stage's own guard depends on that write having happened.
stack.record(ingress)
print(f"[stack] front door at {endpoint}")
return endpoint
def stage_storage() -> str:
"""Block storage — the StorageClass RKE2 doesn't ship."""
storage = stack.build(
Storage,
# The implementation. Everything below `type` is the same question
# asked of any of them; `options` is Longhorn's own.
type="longhorn",
# Needs this many schedulable nodes, or volumes stay Degraded.
replica_count=min(3, len(NODES)),
default_storage_class=True,
options=LonghornOptions(
# Longhorn's own detection reads kubelet's cmdline via pod
# logs and fails whenever the API server can't reach a node,
# taking the CSI driver with it. Pinned rather than detected.
kubelet_root_dir="/var/lib/kubelet",
# A local `helm pull longhorn/longhorn -d <dir>` .tgz path
# here, instead of the chart name, points the install at that
# file directly (HelmRunner's own _local_chart resolution) --
# useful when the chart repo's release-asset CDN is being
# slow/flaky and a chart already sitting on disk is faster
# and more reliable than re-fetching it.
chart=os.environ.get("MULTISTACK_LONGHORN_CHART", "longhorn"),
),
)
backend = StorageBackend()
# Mutates the nodes: open-iscsi + nfs-common. Safe to re-run. Takes the
# spec as well as the nodes — the spec is what says whose packages.
backend.install_prerequisites(storage, NODES)
# EVERY node, not a subset: Longhorn's node components are a
# DaemonSet and land on nodes whether or not they're listed here.
storage_class = backend.create(storage, nodes=NODES)
# After create(), not before: the StorageClass name is knowable from
# the spec, but publishing it early would let the next layer bind
# against a class that isn't installed yet.
stack.record(storage)
return storage_class
def stage_objectstore() -> MinIOTenant:
"""A MinIO tenant on Longhorn.
Re-runnable on its own, which the header promises: if the storage
stage did not run in this process, nothing published a StorageClass,
so name the one that is already on the cluster. The backend still
verifies it exists.
"""
if "storage_class" not in stack:
stack.provide(storage_class="longhorn")
tenant = stack.build(
MinIOTenant,
name="minio-test",
namespace="minio-test",
# 2 x 2 = 4 drives, the minimum for erasure coding.
servers=2,
volumes_per_server=2,
volume_size="20Gi",
# storage_class comes from the Storage spec recorded above.
# root_user/root_password unset: generated, and returned below.
)
info = MinIOBackend().create(tenant)
stack.record(tenant) # publishes the endpoint vLLM reads from
# -- glue -- vLLM's init container reads the weights as this user, so
# the credentials have to exist as a Secret in the inference
# namespace. They never go in a spec.
#
# Built as a manifest and piped to stdin, not assembled with
# `kubectl create secret --from-literal`: that renders the Secret
# locally, but the keys are in the process arguments while it runs,
# where any local user can read them out of /proc. Every other
# credential in this file already went through `apply`; this one was
# the exception.
_apply({
"apiVersion": "v1", "kind": "Namespace",
"metadata": {"name": INFERENCE_NS},
})
_apply({
"apiVersion": "v1", "kind": "Secret",
"metadata": {"name": CREDS_SECRET, "namespace": INFERENCE_NS},
"type": "Opaque",
"data": {
"AWS_ACCESS_KEY_ID":
base64.b64encode(info.root_user.encode()).decode(),
"AWS_SECRET_ACCESS_KEY":
base64.b64encode(info.root_password.encode()).decode(),
},
})
print(f"[stack] wrote {CREDS_SECRET} to namespace {INFERENCE_NS}")
# Generated credentials are worth nothing if nobody ever sees them.
# Printed once, here, because the only other copies are inside two
# Secrets and there is no way to recover them from the spec.
print(f"[stack] MinIO root user: {info.root_user}")
print(f"[stack] MinIO root password: {info.root_password}")
print("[stack] store these — they are generated per tenant and not "
"recoverable from this file")
# -- glue -- the Secret name is a choice made here, not something the
# tenant spec knows, so it is published explicitly.
stack.provide(s3_secret_name=CREDS_SECRET)
return tenant
def stage_weights(tenant: MinIOTenant) -> None:
"""3b. Put the model in the bucket.
-- glue -- Object management isn't SDK surface. This runs one Job on
the cluster that downloads the model and uploads it to MinIO. It uses
hostNetwork because the pod network here has no route to the
internet; the whole point of the exercise is that after this runs,
nothing needs one again.
"""
script = f"""
set -e
pip install --quiet --no-cache-dir huggingface_hub boto3
# -u, because this script's stdout is a pipe and Python block-buffers a
# pipe. Without it `kubectl logs` shows the download finishing and then
# nothing at all -- a stalled upload and a working one look identical for
# as long as it takes to fill 8KB. That cost a nine-minute wait on
# 2026-09-15 to notice MinIO had no write quorum.
python -u - <<'PY'
import glob, os, urllib3
import boto3
from botocore.config import Config
from huggingface_hub import snapshot_download
import os
urllib3.disable_warnings()
path = snapshot_download(
"{MODEL}",
allow_patterns=["*.json", "*.safetensors", "*.txt", "*.model"],
)
s3 = boto3.client(
"s3",
endpoint_url="{tenant.endpoint()}",
aws_access_key_id=os.environ["AWS_ACCESS_KEY_ID"],
aws_secret_access_key=os.environ["AWS_SECRET_ACCESS_KEY"],
# MinIO is path-style, not the virtual-host addressing boto3 assumes.
config=Config(s3={{"addressing_style": "path"}}),
verify=False, # operator's own CA
)
try:
s3.create_bucket(Bucket="{BUCKET}")
except Exception as exc:
print("bucket:", exc)
for f in sorted(glob.glob(os.path.join(path, "*"))):
if not os.path.isfile(f):
continue
key = "{MODEL_KEY}/" + os.path.basename(f)
size = os.path.getsize(f)
try:
# Same name and size already there: a re-run resumes rather than
# re-uploading a gigabyte.
if s3.head_object(Bucket="{BUCKET}", Key=key)["ContentLength"] == size:
print("present", key)
continue
except Exception:
pass
print("uploading", key, size)
s3.upload_file(f, "{BUCKET}", key)
print("done")
PY
"""
job = {
"apiVersion": "batch/v1",
"kind": "Job",
"metadata": {"name": "stage-weights", "namespace": INFERENCE_NS},
"spec": {
"backoffLimit": 1,
# Clean itself up an hour after finishing. Without this a
# completed Job and its pod sit in the namespace forever --
# `kubectl get pods` then shows a Completed pod next to the
# running ones, which reads as something that failed to exit.
"ttlSecondsAfterFinished": 3600,
"template": {"spec": {
"restartPolicy": "Never",
# The pod network has no internet egress; the host does.
"hostNetwork": True,
"dnsPolicy": "ClusterFirstWithHostNet",
"containers": [{
"name": "stage",
"image": "python:3.11-slim",
"command": ["sh", "-c", script],
"envFrom": [{"secretRef": {"name": CREDS_SECRET}}],
}],
}},
},
}
kubectl("delete", "job", "stage-weights", "-n", INFERENCE_NS,
"--ignore-not-found")
_apply(job)
print(f"[stack] staging {MODEL} into s3://{BUCKET}/{MODEL_KEY} (a few minutes)")
wait_for_job(KUBECONFIG, "stage-weights", INFERENCE_NS)
logs = kubectl("logs", "job/stage-weights", "-n", INFERENCE_NS,
check=False).strip().splitlines()
print(f"[stack] staging finished: {logs[-1] if logs else 'no output'}")
def stage_inference(tenant: MinIOTenant) -> str:
"""Serve the model, reading weights from the tenant above."""
# -- glue -- keep everything else off this node, so vLLM gets its CPU.
name = node_name(INFERENCE_NODE)
kubectl("label", "node", name, "multistack.io/workload=inference", "--overwrite")
kubectl("taint", "node", name, "dedicated=inference:NoSchedule", "--overwrite")
print(f"[stack] reserved node {name} for inference")
service = stack.build(
Inference,
# The composition: weights come from the tenant created above, so
# this pod needs no internet access. s3_endpoint_url and
# s3_secret_name are both filled in from the stack.
model=f"s3://{BUCKET}/{MODEL_KEY}",
s3_insecure_tls=tenant.request_auto_cert,
name="vllm-qwen",
namespace=INFERENCE_NS,
device="cpu",
# None: chosen from the target node's CPU flags, so a bfloat16
# model isn't emulated in software.
dtype=None,
max_model_len=4096,
# Leaves cores for kubelet/CNI/DaemonSets — asking for all of
# them leaves the pod Pending.
cpu_cores=6,
memory_gb=12,
kv_cache_gb=4,
node_selector={"multistack.io/workload": "inference"},
tolerations=[{
"key": "dedicated", "operator": "Equal",
"value": "inference", "effect": "NoSchedule",
}],
# No ingress here, so this is what makes the endpoint reachable
# from outside the cluster and lets the warmup request land.
host_network=True,
)
target = next(n for n in NODES if n.address == INFERENCE_NODE)
endpoint = InferenceBackend().create(service, nodes=[target])
# Publishes inference_endpoint. Nothing's FROM_STACK reads it — the
# gateway's upstream_url is not wired, deliberately — but the stack
# should still know what it built.
stack.record(service)
return endpoint
def stage_cache() -> str:
"""Valkey — the counter store every rate-limit policy reads.
Returns the *authenticated* URL. Four things read this cache — both
rate limiters and both control planes — and the URL the Stack
publishes is deliberately credential-free, so the authenticated one
is threaded through the return value the same way the database
stage's is.
The chart enables auth whether or not it is asked to: `auth.enabled`
defaults true and `auth.password` empty means "generate a random
ten-character password", into a Secret nobody here reads. That is the
worst of the three possible outcomes -- Helm succeeds, the pods go
Ready, and every consumer gets NOAUTH at its first request, which for
a rate limiter means failing open. So the password originates here and
the release is pointed at a Secret holding it.
"""
# -- glue -- the credential the release reads, applied before the
# chart that mounts it. Through stdin like every other Secret in this
# file: `kubectl create secret --from-literal` would put the password
# in the process arguments, readable from /proc while it runs.
password = _cache_password()
_ensure_namespace(CACHE_NS)
_apply({
"apiVersion": "v1", "kind": "Secret",
"metadata": {"name": CACHE_AUTH_SECRET, "namespace": CACHE_NS},
"type": "Opaque",
"data": {CACHE_AUTH_KEY:
base64.b64encode(password.encode()).decode()},
})
print(f"[stack] wrote {CACHE_AUTH_SECRET} to namespace {CACHE_NS}")
valkey = stack.build(
Cache,
name="valkey", namespace=CACHE_NS,
# No password in the spec. `existingSecret` is a *reference* to
# the Secret applied above, which is the difference between a spec
# people can commit and one they cannot.
options=ValkeyOptions(values={
"auth": {
"enabled": True,
"existingSecret": CACHE_AUTH_SECRET,
"existingSecretPasswordKey": CACHE_AUTH_KEY,
},
# One node, not the chart's default `replication`. These are
# rate-limit counters: they are rebuilt from the event stream
# within a window, so a replica set costs three PVCs and a
# failover story to protect data that is cheap to lose. State
# it, rather than inheriting a topology nobody chose.
"architecture": "standalone",
# AOF, so a restart does not silently reset every counter to
# zero -- which reads as "nobody has made a request yet" and
# grants everyone a fresh allowance.
"primary": {"persistence": {"enabled": True, "size": "8Gi"}},
}),
)
CacheBackend().create(valkey)
stack.record(valkey) # publishes cache_url, credential-free
print(f"[stack] cache at {valkey.endpoint}")
return _cache_authenticated(valkey.endpoint, password)
def stage_database() -> dict:
"""CloudNativePG — one Postgres cluster per database-owning service
(both control planes, and billing).
Returns {component: authenticated DATABASE_URL}. The password is the
one credential that has to originate here: CloudNativePG reads
bootstrap.initdb only at bootstrap, so this value is what the role
is created with. Everything afterwards reads it from the Secret,
which is why `database_url` carries no password — and why the
authenticated URL is returned to the caller rather than published.
Asked for at run time rather than generated, so that whoever runs
this knows the passwords afterwards. MULTISTACK_ADMIN_DB_PASSWORD,
MULTISTACK_ORGANIZATION_DB_PASSWORD and MULTISTACK_BILLING_DB_PASSWORD
keep the script runnable unattended -- ORGANIZATION, not ORG: the
name is derived from `component` below, and DATABASES spells that
out in full.
"""
backend = DatabaseBackend()
urls = {}
for component, cluster_name, db_name, owner in DATABASES:
password = prompt_secret(
f"PostgreSQL password for role {owner!r} ({cluster_name})",
env_var=f"MULTISTACK_{component.upper()}_DB_PASSWORD",
confirm=True,
)
database = stack.build(
Database,
# namespace omitted: the cnpg implementation's own default
# is already `postgres`.
name=cluster_name,
database=DatabaseConfig(name=db_name, owner=owner,
password=password),
# Same local-chart escape hatch as storage's LonghornOptions.chart
# -- the cnpg repo's index.yaml is on GitHub Pages, but the
# tarball itself comes from release-assets.githubusercontent.com,
# the same flaky leg that hit Longhorn. A local path already
# pins a version by definition, and HelmRunner's own
# _local_chart resolution refuses to combine one with a
# chart_version -- so drop the operator's default pin
# (0.29.0) whenever a local path is in play.
options=CNPGOptions(
operator_chart=os.environ.get(