Monorepo for Tangled
0

Configure Feed

Select the types of activity you want to include in your feed.

spindle/microvm: add resource budget limits and optional cgroup enforcement

Signed-off-by: dawn <dawn@tangled.org>

dawn (Jun 11, 2026, 5:14 PM +0300) ff106eb7 62565c57

+741 -56
+8 -1
go.mod
··· 29 29 github.com/charmbracelet/ssh v0.0.0-20250128164007-98fd5ae11894 30 30 github.com/charmbracelet/wish v1.4.7 31 31 github.com/cloudflare/cloudflare-go/v6 v6.7.0 32 + github.com/containerd/cgroups/v3 v3.1.3 32 33 github.com/cyphar/filepath-securejoin v0.4.1 33 34 github.com/dgraph-io/ristretto v0.2.0 34 35 github.com/did-method-plc/go-didplc v0.2.2 ··· 57 58 github.com/openbao/openbao/api/v2 v2.3.0 58 59 github.com/posthog/posthog-go v1.5.5 59 60 github.com/prometheus/client_golang v1.23.2 61 + github.com/prometheus/procfs v0.19.2 60 62 github.com/redis/go-redis/v9 v9.7.3 61 63 github.com/resend/resend-go/v3 v3.5.0 62 64 github.com/sethvargo/go-envconfig v1.1.0 ··· 139 141 github.com/charmbracelet/x/term v0.2.2 // indirect 140 142 github.com/charmbracelet/x/termios v0.1.0 // indirect 141 143 github.com/charmbracelet/x/windows v0.2.0 // indirect 144 + github.com/cilium/ebpf v0.16.0 // indirect 142 145 github.com/clipperhouse/displaywidth v0.9.0 // indirect 143 146 github.com/clipperhouse/stringish v0.1.1 // indirect 144 147 github.com/clipperhouse/uax29/v2 v2.5.0 // indirect ··· 146 149 github.com/containerd/errdefs v1.0.0 // indirect 147 150 github.com/containerd/errdefs/pkg v0.3.0 // indirect 148 151 github.com/containerd/log v0.1.0 // indirect 152 + github.com/coreos/go-systemd/v22 v22.5.0 // indirect 149 153 github.com/creack/pty v1.1.21 // indirect 150 154 github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect 151 155 github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect ··· 169 173 github.com/go-redis/cache/v9 v9.0.0 // indirect 170 174 github.com/go-test/deep v1.1.1 // indirect 171 175 github.com/goccy/go-json v0.10.5 // indirect 176 + github.com/godbus/dbus/v5 v5.1.0 // indirect 172 177 github.com/gogo/protobuf v1.3.2 // indirect 173 178 github.com/golang-jwt/jwt v3.2.2+incompatible // indirect 174 179 github.com/golang-jwt/jwt/v5 v5.3.0 // indirect ··· 231 236 github.com/mitchellh/mapstructure v1.5.0 // indirect 232 237 github.com/moby/docker-image-spec v1.3.1 // indirect 233 238 github.com/moby/sys/atomicwriter v0.1.0 // indirect 239 + github.com/moby/sys/userns v0.1.0 // indirect 234 240 github.com/moby/term v0.5.2 // indirect 235 241 github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect 236 242 github.com/modern-go/reflect2 v1.0.2 // indirect ··· 248 254 github.com/onsi/gomega v1.37.0 // indirect 249 255 github.com/opencontainers/go-digest v1.0.0 // indirect 250 256 github.com/opencontainers/image-spec v1.1.1 // indirect 257 + github.com/opencontainers/runtime-spec v1.3.0 // indirect 251 258 github.com/opentracing/opentracing-go v1.2.1-0.20220228012449-10b1cf09e00b // indirect 252 259 github.com/pjbgf/sha1cd v0.3.2 // indirect 253 260 github.com/pkg/errors v0.9.1 // indirect ··· 255 262 github.com/polydawn/refmt v0.89.1-0.20221221234430-40501e09de1f // indirect 256 263 github.com/prometheus/client_model v0.6.2 // indirect 257 264 github.com/prometheus/common v0.67.5 // indirect 258 - github.com/prometheus/procfs v0.19.2 // indirect 259 265 github.com/puzpuzpuz/xsync/v4 v4.2.0 // indirect 260 266 github.com/rivo/uniseg v0.4.7 // indirect 261 267 github.com/ryanuber/go-glob v1.0.0 // indirect 262 268 github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 // indirect 269 + github.com/sirupsen/logrus v1.9.3 // indirect 263 270 github.com/spaolacci/murmur3 v1.1.0 // indirect 264 271 github.com/tidwall/gjson v1.18.0 // indirect 265 272 github.com/tidwall/match v1.2.0 // indirect
+14
go.sum
··· 193 193 github.com/chzyer/logex v1.1.10/go.mod h1:+Ywpsq7O8HXn0nuIou7OrIPyXbp3wmkHB+jjWRnGsAI= 194 194 github.com/chzyer/readline v0.0.0-20180603132655-2972be24d48e/go.mod h1:nSuG5e5PlCu98SY8svDHJxuZscDgtXS6KTTbou5AhLI= 195 195 github.com/chzyer/test v0.0.0-20180213035817-a1ea475d72b1/go.mod h1:Q3SI9o4m/ZMnBNeIyt5eFwwo7qiLfzFZmjNmxjkiQlU= 196 + github.com/cilium/ebpf v0.16.0 h1:+BiEnHL6Z7lXnlGUsXQPPAE7+kenAd4ES8MQ5min0Ok= 197 + github.com/cilium/ebpf v0.16.0/go.mod h1:L7u2Blt2jMM/vLAVgjxluxtBKlz3/GWjB0dMOEngfwE= 196 198 github.com/clipperhouse/displaywidth v0.9.0 h1:Qb4KOhYwRiN3viMv1v/3cTBlz3AcAZX3+y9OLhMtAtA= 197 199 github.com/clipperhouse/displaywidth v0.9.0/go.mod h1:aCAAqTlh4GIVkhQnJpbL0T/WfcrJXHcj8C0yjYcjOZA= 198 200 github.com/clipperhouse/stringish v0.1.1 h1:+NSqMOr3GR6k1FdRhhnXrLfztGzuG+VuFDfatpWHKCs= ··· 203 205 github.com/cloudflare/circl v1.6.2-0.20250618153321-aa837fd1539d/go.mod h1:uddAzsPgqdMAYatqJ0lsjX1oECcQLIlRpzZh3pJrofs= 204 206 github.com/cloudflare/cloudflare-go/v6 v6.7.0 h1:MP6Xy5WmsyrxgTxoLeq/vraqR0nbTtXoHhW4vAYc4SY= 205 207 github.com/cloudflare/cloudflare-go/v6 v6.7.0/go.mod h1:Lj3MUqjvKctXRpdRhLQxZYRrNZHuRs0XYuH8JtQGyoI= 208 + github.com/containerd/cgroups/v3 v3.1.3 h1:eUNflyMddm18+yrDmZPn3jI7C5hJ9ahABE5q6dyLYXQ= 209 + github.com/containerd/cgroups/v3 v3.1.3/go.mod h1:PKZ2AcWmSBsY/tJUVhtS/rluX0b1uq1GmPO1ElCmbOw= 206 210 github.com/containerd/errdefs v1.0.0 h1:tg5yIfIlQIrxYtu9ajqY42W3lpS19XqdxRQeEwYG8PI= 207 211 github.com/containerd/errdefs v1.0.0/go.mod h1:+YBYIdtsnF4Iw6nWZhJcqGSg/dwvV7tyJ/kCkyJ2k+M= 208 212 github.com/containerd/errdefs/pkg v0.3.0 h1:9IKJ06FvyNlexW690DXuQNx2KA2cUJXx151Xdx3ZPPE= 209 213 github.com/containerd/errdefs/pkg v0.3.0/go.mod h1:NJw6s9HwNuRhnjJhM7pylWwMyAkmCQvQ4GpJHEqRLVk= 210 214 github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= 211 215 github.com/containerd/log v0.1.0/go.mod h1:VRRf09a7mHDIRezVKTRCrOq78v577GXq3bSa3EhrzVo= 216 + github.com/coreos/go-systemd/v22 v22.5.0 h1:RrqgGjYQKalulkV8NGVIfkXQf6YYmOyiJKk8iXXhfZs= 217 + github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc= 212 218 github.com/cpuguy83/go-md2man/v2 v2.0.0-20190314233015-f79a8a8ca69d/go.mod h1:maD7wRr/U5Z6m/iR4s+kqSMx2CaBsrgA7czyZG/E6dU= 213 219 github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= 214 220 github.com/creack/pty v1.1.21 h1:1/QdRyBaHHJP61QkWMXlOIBfsgdDeeKfK8SYVUWJKf0= ··· 315 321 github.com/gobwas/ws v1.4.0/go.mod h1:G3gNqMNtPppf5XUz7O4shetPpcZ1VJ7zt18dlUeakrc= 316 322 github.com/goccy/go-json v0.10.5 h1:Fq85nIqj+gXn/S5ahsiTlK3TmC85qgirsdTP/+DeaC4= 317 323 github.com/goccy/go-json v0.10.5/go.mod h1:oq7eo15ShAhp70Anwd5lgX2pLfOS3QCiwU/PULtXL6M= 324 + github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= 325 + github.com/godbus/dbus/v5 v5.1.0 h1:4KLkAxT3aOY8Li4FRJe/KvhoNFFxo0m6fNuFUO8QJUk= 326 + github.com/godbus/dbus/v5 v5.1.0/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= 318 327 github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= 319 328 github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= 320 329 github.com/golang-jwt/jwt v3.2.2+incompatible h1:IfV12K8xAKAnZqdXVzCZ+TOjboZ2keLg81eXfW3O+oY= ··· 552 561 github.com/moby/sys/atomicwriter v0.1.0/go.mod h1:Ul8oqv2ZMNHOceF643P6FKPXeCmYtlQMvpizfsSoaWs= 553 562 github.com/moby/sys/sequential v0.6.0 h1:qrx7XFUd/5DxtqcoH1h438hF5TmOvzC/lspjy7zgvCU= 554 563 github.com/moby/sys/sequential v0.6.0/go.mod h1:uyv8EUTrca5PnDsdMGXhZe6CCe8U/UiTWd+lL+7b/Ko= 564 + github.com/moby/sys/userns v0.1.0 h1:tVLXkFOxVu9A64/yh59slHVv9ahO9UIev4JZusOLG/g= 565 + github.com/moby/sys/userns v0.1.0/go.mod h1:IHUYgu/kao6N8YZlp9Cf444ySSvCmDlmzUcYfDHOl28= 555 566 github.com/moby/term v0.5.2 h1:6qk3FJAFDs6i/q3W/pQ97SX192qKfZgGjCQqfCJkgzQ= 556 567 github.com/moby/term v0.5.2/go.mod h1:d3djjFCrjnB+fl8NJux+EJzu0msscUP+f8it8hPkFLc= 557 568 github.com/modern-go/concurrent v0.0.0-20180228061459-e0a39a4cb421/go.mod h1:6dJC0mAP4ikYIbvyc7fijjWJddQyLn8Ig3JB5CqoB9Q= ··· 624 635 github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= 625 636 github.com/opencontainers/image-spec v1.1.1 h1:y0fUlFfIZhPF1W537XOLg0/fcx6zcHCJwooC2xJA040= 626 637 github.com/opencontainers/image-spec v1.1.1/go.mod h1:qpqAh3Dmcf36wStyyWU+kCeDgrGnAve2nCC8+7h8Q0M= 638 + github.com/opencontainers/runtime-spec v1.3.0 h1:YZupQUdctfhpZy3TM39nN9Ika5CBWT5diQ8ibYCRkxg= 639 + github.com/opencontainers/runtime-spec v1.3.0/go.mod h1:jwyrGlmzljRJv/Fgzds9SsS/C5hL+LL3ko9hs6T5lQ0= 627 640 github.com/opentracing/opentracing-go v1.2.0/go.mod h1:GxEUsuufX4nBwe+T+Wl9TAgYrxe9dPLANfrWvHYVTgc= 628 641 github.com/opentracing/opentracing-go v1.2.1-0.20220228012449-10b1cf09e00b h1:FfH+VrHHk6Lxt9HdVS0PXzSXFyS2NbZKXv33FYPol0A= 629 642 github.com/opentracing/opentracing-go v1.2.1-0.20220228012449-10b1cf09e00b/go.mod h1:AC62GU6hc0BrNm+9RK9VSiwa/EUe1bkIeFORAMcHvJU= ··· 867 880 golang.org/x/sys v0.0.0-20220319134239-a9b59b0215f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= 868 881 golang.org/x/sys v0.0.0-20220422013727-9388b58f7150/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= 869 882 golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= 883 + golang.org/x/sys v0.0.0-20220715151400-c0bba94af5f8/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= 870 884 golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= 871 885 golang.org/x/sys v0.0.0-20220908164124-27713097b956/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= 872 886 golang.org/x/sys v0.1.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
+21
nix/gomod2nix.toml
··· 256 256 [mod."github.com/charmbracelet/x/windows"] 257 257 version = "v0.2.0" 258 258 hash = "sha256-pDAd1E5w66E/d3vuTyzgnW+W/KegZ2sxQQMfoEn7S1A=" 259 + [mod."github.com/cilium/ebpf"] 260 + version = "v0.16.0" 261 + hash = "sha256-xACuieGmiUUjoTT/9MpvPBNexp98S/AZbLxm5f9nqDk=" 259 262 [mod."github.com/clipperhouse/displaywidth"] 260 263 version = "v0.9.0" 261 264 hash = "sha256-9CNyTZPSncKQ7Y0my9DR4WYXDjtDHYNL512D691WDAM=" ··· 271 274 [mod."github.com/cloudflare/cloudflare-go/v6"] 272 275 version = "v6.7.0" 273 276 hash = "sha256-ycQpx1II/JgBgrCRwY5qiVKStGv5wuCANy1091sJ5Zw=" 277 + [mod."github.com/containerd/cgroups/v3"] 278 + version = "v3.1.3" 279 + hash = "sha256-1a5heWXIzME7iMu2L35OBiAOi2Z/gnpg2fjvP6On9sM=" 274 280 [mod."github.com/containerd/errdefs"] 275 281 version = "v1.0.0" 276 282 hash = "sha256-wMZGoeqvRhuovYCJx0Js4P3qFCNTZ/6Atea/kNYoPMI=" ··· 280 286 [mod."github.com/containerd/log"] 281 287 version = "v0.1.0" 282 288 hash = "sha256-vuE6Mie2gSxiN3jTKTZovjcbdBd1YEExb7IBe3GM+9s=" 289 + [mod."github.com/coreos/go-systemd/v22"] 290 + version = "v22.5.0" 291 + hash = "sha256-E2zXikbmIQImghstLUWuey1YgA0Folu3F+fi5k4hCxA=" 283 292 [mod."github.com/creack/pty"] 284 293 version = "v1.1.21" 285 294 hash = "sha256-pjGw6wQlrVhN65XaIxZueNJqnXThGu00u24rKOLzxS0=" ··· 386 395 [mod."github.com/goccy/go-json"] 387 396 version = "v0.10.5" 388 397 hash = "sha256-/EtlGihP0/7oInzMC5E0InZ4b5Ad3s4xOpqotloi3xw=" 398 + [mod."github.com/godbus/dbus/v5"] 399 + version = "v5.1.0" 400 + hash = "sha256-xOCMJpQK3KTmHTPn/CdqI4j0eENCtMmJDgAIoYqYOEY=" 389 401 [mod."github.com/gogo/protobuf"] 390 402 version = "v1.3.2" 391 403 hash = "sha256-pogILFrrk+cAtb0ulqn9+gRZJ7sGnnLLdtqITvxvG6c=" ··· 608 620 [mod."github.com/moby/sys/atomicwriter"] 609 621 version = "v0.1.0" 610 622 hash = "sha256-i46GNrsICnJ0AYkN+ocbVZ2GNTQVEsrVX5WcjKzjtBM=" 623 + [mod."github.com/moby/sys/userns"] 624 + version = "v0.1.0" 625 + hash = "sha256-zwXKyEZIH/FZjSVuSGmtwThDxPutj1pY+N6Ae6oVPuc=" 611 626 [mod."github.com/moby/term"] 612 627 version = "v0.5.2" 613 628 hash = "sha256-/G20jUZKx36ktmPU/nEw/gX7kRTl1Dbu7zvNBYNt4xU=" ··· 665 680 [mod."github.com/opencontainers/image-spec"] 666 681 version = "v1.1.1" 667 682 hash = "sha256-bxBjtl+6846Ed3QHwdssOrNvlHV6b+Dn17zPISSQGP8=" 683 + [mod."github.com/opencontainers/runtime-spec"] 684 + version = "v1.3.0" 685 + hash = "sha256-B2QF7FlUYZDL9eNA0+JD7WasnBryMXNIDbdSGS4MMG4=" 668 686 [mod."github.com/opentracing/opentracing-go"] 669 687 version = "v1.2.1-0.20220228012449-10b1cf09e00b" 670 688 hash = "sha256-77oWcDviIoGWHVAotbgmGRpLGpH5AUy+pM15pl3vRrw=" ··· 717 735 [mod."github.com/sethvargo/go-envconfig"] 718 736 version = "v1.1.0" 719 737 hash = "sha256-WelRHfyZG9hrA4fbQcfBawb2ZXBQNT1ourEYHzQdZ4w=" 738 + [mod."github.com/sirupsen/logrus"] 739 + version = "v1.9.3" 740 + hash = "sha256-EnxsWdEUPYid+aZ9H4/iMTs1XMvCLbXZRDyvj89Ebms=" 720 741 [mod."github.com/spaolacci/murmur3"] 721 742 version = "v1.1.0" 722 743 hash = "sha256-RWD4PPrlAsZZ8Xy356MBxpj+/NZI7w2XOU14Ob7/Y9M="
+13
spindle/config/config.go
··· 65 65 DefaultImage string `env:"DEFAULT_IMAGE, default=spindle-base"` 66 66 EnableKVM bool `env:"ENABLE_KVM, default=true"` 67 67 WorkflowTimeout string `env:"WORKFLOW_TIMEOUT, default=5m"` 68 + 69 + MaxTotalMemoryMiB int64 `env:"BUDGET_MEMORY_MIB, default=0"` 70 + MaxTotalVCPUs int64 `env:"BUDGET_VCPUS, default=0"` 71 + MaxTotalDiskMiB int64 `env:"BUDGET_DISK_MIB, default=0"` 72 + 73 + MaxWorkflowMemoryMiB int64 `env:"MAX_WORKFLOW_MEMORY_MIB, default=0"` 74 + MaxWorkflowVCPUs int64 `env:"MAX_WORKFLOW_VCPUS, default=0"` 75 + MaxWorkflowDiskMiB int64 `env:"MAX_WORKFLOW_DISK_MIB, default=0"` 76 + 77 + EnableCgroups bool `env:"ENABLE_CGROUPS, default=false"` 78 + CgroupParent string `env:"CGROUP_PARENT, default=self"` 79 + CgroupPidsMax int64 `env:"CGROUP_PIDS_MAX, default=4096"` 80 + CgroupSwapMaxMiB *int64 `env:"CGROUP_SWAP_MAX_MIB"` 68 81 } 69 82 70 83 type NixCache struct {
+14 -5
spindle/engine/engine.go
··· 24 24 FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) error 25 25 } 26 26 27 - func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, db *db.DB, n *notifier.Notifier, workflowSem chan struct{}, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { 27 + func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, db *db.DB, n *notifier.Notifier, globalSlotter WorkflowSlotter, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { 28 28 l.Info("starting all workflows in parallel", "pipeline", pipelineId) 29 29 30 30 // extract secrets ··· 78 78 defer wfLogger.Close() 79 79 } 80 80 81 + l.Info("waiting for slot", "wid", wid) 82 + s, _ := eng.(WorkflowSlotter) 83 + slot, err := ChainSlotter{globalSlotter, s}.AcquireWorkflowSlot(ctx, wid, &w) 84 + if err != nil { 85 + l.Error("failed to acquire slot", "wid", wid, "err", err) 86 + dbErr := db.StatusFailed(wid, err.Error(), -1, n) 87 + if dbErr != nil { 88 + l.Error("failed to set workflow status to failed", "wid", wid, "err", dbErr) 89 + } 90 + return 91 + } 92 + defer slot.Release() 93 + 81 94 err = db.StatusRunning(wid, n) 82 95 if err != nil { 83 96 l.Error("failed to set workflow status to running", "wid", wid, "err", err) 84 97 return 85 98 } 86 - 87 - // acquire semaphore slot before starting the container 88 - workflowSem <- struct{}{} 89 - defer func() { <-workflowSem }() 90 99 91 100 err = eng.SetupWorkflow(ctx, wid, &w, wfLogger) 92 101 if err != nil {
+79
spindle/engine/slot.go
··· 1 + package engine 2 + 3 + import ( 4 + "context" 5 + "errors" 6 + 7 + "tangled.org/core/spindle/models" 8 + ) 9 + 10 + var ErrNoWorkflowSlots = errors.New("no workflow slots available") 11 + 12 + type WorkflowSlot interface { 13 + Release() 14 + } 15 + 16 + type WorkflowSlotter interface { 17 + AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow) (WorkflowSlot, error) 18 + } 19 + 20 + type releaseFunc func() 21 + 22 + func (f releaseFunc) Release() { 23 + if f != nil { 24 + f() 25 + } 26 + } 27 + 28 + type NoopSlot struct{} 29 + 30 + func (NoopSlot) Release() {} 31 + 32 + // limit by concurrent workflow count 33 + type SemaphoreSlotter struct { 34 + slots chan struct{} 35 + } 36 + 37 + func NewSemaphoreSlotter(maxConcurrent int) *SemaphoreSlotter { 38 + if maxConcurrent <= 0 { 39 + return &SemaphoreSlotter{} 40 + } 41 + return &SemaphoreSlotter{slots: make(chan struct{}, maxConcurrent)} 42 + } 43 + 44 + func (a *SemaphoreSlotter) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow) (WorkflowSlot, error) { 45 + if a == nil || a.slots == nil { 46 + return NoopSlot{}, nil 47 + } 48 + select { 49 + case a.slots <- struct{}{}: 50 + return releaseFunc(func() { <-a.slots }), nil 51 + case <-ctx.Done(): 52 + return nil, ctx.Err() 53 + } 54 + } 55 + 56 + // compose multiple slotters, all must be satisfied before the slot is acquired 57 + type ChainSlotter []WorkflowSlotter 58 + 59 + func (c ChainSlotter) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow) (WorkflowSlot, error) { 60 + acquired := make([]WorkflowSlot, 0, len(c)) 61 + for _, s := range c { 62 + if s == nil { 63 + continue 64 + } 65 + slot, err := s.AcquireWorkflowSlot(ctx, wid, wf) 66 + if err != nil { 67 + for _, s := range acquired { 68 + s.Release() 69 + } 70 + return nil, err 71 + } 72 + acquired = append(acquired, slot) 73 + } 74 + return releaseFunc(func() { 75 + for _, s := range acquired { 76 + s.Release() 77 + } 78 + }), nil 79 + }
+243
spindle/engines/microvm/budget.go
··· 1 + package microvm 2 + 3 + import ( 4 + "context" 5 + "fmt" 6 + "runtime" 7 + "sync" 8 + 9 + "github.com/prometheus/procfs" 10 + 11 + "tangled.org/core/spindle/config" 12 + "tangled.org/core/spindle/engine" 13 + "tangled.org/core/spindle/models" 14 + ) 15 + 16 + type vmResources struct { 17 + MemoryMiB int64 18 + VCPUs int64 19 + DiskMiB int64 20 + } 21 + 22 + type vmScheduler struct { 23 + mu sync.Mutex 24 + budget vmResources 25 + max vmResources 26 + used vmResources 27 + queue []*vmWaiter 28 + } 29 + 30 + type vmWaiter struct { 31 + req vmResources 32 + ready chan struct{} 33 + } 34 + 35 + type vmLease struct { 36 + scheduler *vmScheduler 37 + req vmResources 38 + once sync.Once 39 + } 40 + 41 + func newVMScheduler(budget, max vmResources) *vmScheduler { 42 + return &vmScheduler{budget: budget, max: max} 43 + } 44 + 45 + func newVMBudgetConfig(cfg config.MicroVMPipelines) (vmResources, vmResources, error) { 46 + budget := vmResources{ 47 + MemoryMiB: cfg.MaxTotalMemoryMiB, 48 + VCPUs: cfg.MaxTotalVCPUs, 49 + DiskMiB: cfg.MaxTotalDiskMiB, 50 + } 51 + if budget.MemoryMiB <= 0 { 52 + hostMemory := hostMemoryMiB() 53 + if hostMemory <= 0 { 54 + return vmResources{}, vmResources{}, fmt.Errorf("detect host memory: set SPINDLE_MICROVM_PIPELINES_BUDGET_MEMORY_MIB explicitly") 55 + } 56 + budget.MemoryMiB = hostMemory - 1024 57 + if budget.MemoryMiB <= 0 { 58 + return vmResources{}, vmResources{}, fmt.Errorf("microVM budget memory is non-positive: hostMemoryMiB=%d reservedMemoryMiB=1024", hostMemory) 59 + } 60 + } 61 + if budget.VCPUs <= 0 { 62 + budget.VCPUs = int64(runtime.NumCPU()) 63 + } 64 + 65 + maxReq := vmResources{ 66 + MemoryMiB: cfg.MaxWorkflowMemoryMiB, 67 + VCPUs: cfg.MaxWorkflowVCPUs, 68 + DiskMiB: cfg.MaxWorkflowDiskMiB, 69 + } 70 + if maxReq.MemoryMiB <= 0 { 71 + maxReq.MemoryMiB = budget.MemoryMiB 72 + } 73 + if maxReq.VCPUs <= 0 { 74 + maxReq.VCPUs = budget.VCPUs 75 + } 76 + if maxReq.DiskMiB <= 0 { 77 + maxReq.DiskMiB = budget.DiskMiB 78 + } 79 + 80 + return budget, maxReq, nil 81 + } 82 + 83 + func hostMemoryMiB() int64 { 84 + fs, err := procfs.NewFS("/proc") 85 + if err != nil { 86 + return 0 87 + } 88 + meminfo, err := fs.Meminfo() 89 + if err != nil || meminfo.MemTotalBytes == nil { 90 + return 0 91 + } 92 + return int64(*meminfo.MemTotalBytes / 1024 / 1024) 93 + } 94 + 95 + func (e *Engine) AcquireWorkflowSlot(ctx context.Context, wid models.WorkflowId, wf *models.Workflow) (engine.WorkflowSlot, error) { 96 + state, ok := wf.Data.(*workflowState) 97 + if !ok || state == nil { 98 + return nil, fmt.Errorf("microVM workflow state is not initialized") 99 + } 100 + if e.scheduler == nil { 101 + return engine.NoopSlot{}, nil 102 + } 103 + return e.scheduler.Acquire(ctx, resourcesForImage(state.ImageSpec)) 104 + } 105 + 106 + func resourcesForImage(spec ImageSpec) vmResources { 107 + var diskMiB int64 108 + for _, volume := range spec.Volumes { 109 + diskMiB += volume.SizeMiB 110 + } 111 + return vmResources{ 112 + MemoryMiB: int64(spec.MemoryMiB), 113 + VCPUs: int64(spec.VCPUs), 114 + DiskMiB: diskMiB, 115 + } 116 + } 117 + 118 + func (s *vmScheduler) Acquire(ctx context.Context, req vmResources) (engine.WorkflowSlot, error) { 119 + if s == nil { 120 + return engine.NoopSlot{}, nil 121 + } 122 + if err := validateRequest(req); err != nil { 123 + return nil, err 124 + } 125 + 126 + s.mu.Lock() 127 + if !req.fits(s.budget) || !req.fits(s.max) { 128 + s.mu.Unlock() 129 + return nil, fmt.Errorf("%w: request=%s budget=%s max=%s", engine.ErrNoWorkflowSlots, req, s.budget, s.max) 130 + } 131 + if len(s.queue) == 0 && s.availableLocked(req) { 132 + s.used = s.used.add(req) 133 + s.mu.Unlock() 134 + return &vmLease{scheduler: s, req: req}, nil 135 + } 136 + 137 + waiter := &vmWaiter{req: req, ready: make(chan struct{})} 138 + s.queue = append(s.queue, waiter) 139 + s.scheduleLocked() 140 + s.mu.Unlock() 141 + 142 + select { 143 + case <-waiter.ready: 144 + return &vmLease{scheduler: s, req: req}, nil 145 + case <-ctx.Done(): 146 + s.mu.Lock() 147 + select { 148 + case <-waiter.ready: 149 + // undo committed resources, drainLocked already did that 150 + s.used = s.used.sub(req) 151 + default: 152 + // still in queue, just remove 153 + s.removeLocked(waiter) 154 + } 155 + s.scheduleLocked() 156 + s.mu.Unlock() 157 + return nil, ctx.Err() 158 + } 159 + } 160 + 161 + func (l *vmLease) Release() { 162 + if l == nil || l.scheduler == nil { 163 + return 164 + } 165 + l.once.Do(func() { 166 + l.scheduler.release(l.req) 167 + }) 168 + } 169 + 170 + func (s *vmScheduler) release(req vmResources) { 171 + s.mu.Lock() 172 + defer s.mu.Unlock() 173 + s.used = s.used.sub(req) 174 + s.scheduleLocked() 175 + } 176 + 177 + func (s *vmScheduler) scheduleLocked() { 178 + for len(s.queue) > 0 { 179 + waiter := s.queue[0] 180 + if !s.availableLocked(waiter.req) { 181 + return 182 + } 183 + s.queue = s.queue[1:] 184 + s.used = s.used.add(waiter.req) 185 + close(waiter.ready) 186 + } 187 + } 188 + 189 + func (s *vmScheduler) removeLocked(waiter *vmWaiter) { 190 + for i, candidate := range s.queue { 191 + if candidate != waiter { 192 + continue 193 + } 194 + copy(s.queue[i:], s.queue[i+1:]) 195 + s.queue[len(s.queue)-1] = nil 196 + s.queue = s.queue[:len(s.queue)-1] 197 + return 198 + } 199 + } 200 + 201 + func (s *vmScheduler) availableLocked(req vmResources) bool { 202 + return s.used.add(req).fits(s.budget) 203 + } 204 + 205 + func validateRequest(req vmResources) error { 206 + if req.MemoryMiB < 0 || req.VCPUs < 0 || req.DiskMiB < 0 { 207 + return fmt.Errorf("microVM resource request must not be negative: %s", req) 208 + } 209 + return nil 210 + } 211 + 212 + func (r vmResources) fits(limit vmResources) bool { 213 + if limit.MemoryMiB > 0 && r.MemoryMiB > limit.MemoryMiB { 214 + return false 215 + } 216 + if limit.VCPUs > 0 && r.VCPUs > limit.VCPUs { 217 + return false 218 + } 219 + if limit.DiskMiB > 0 && r.DiskMiB > limit.DiskMiB { 220 + return false 221 + } 222 + return true 223 + } 224 + 225 + func (r vmResources) add(other vmResources) vmResources { 226 + return vmResources{ 227 + MemoryMiB: r.MemoryMiB + other.MemoryMiB, 228 + VCPUs: r.VCPUs + other.VCPUs, 229 + DiskMiB: r.DiskMiB + other.DiskMiB, 230 + } 231 + } 232 + 233 + func (r vmResources) sub(other vmResources) vmResources { 234 + return vmResources{ 235 + MemoryMiB: max(0, r.MemoryMiB-other.MemoryMiB), 236 + VCPUs: max(0, r.VCPUs-other.VCPUs), 237 + DiskMiB: max(0, r.DiskMiB-other.DiskMiB), 238 + } 239 + } 240 + 241 + func (r vmResources) String() string { 242 + return fmt.Sprintf("memory=%dMiB vcpus=%d disk=%dMiB", r.MemoryMiB, r.VCPUs, r.DiskMiB) 243 + }
+246
spindle/engines/microvm/cgroup.go
··· 1 + package microvm 2 + 3 + import ( 4 + "fmt" 5 + "log/slog" 6 + "os" 7 + "path/filepath" 8 + "regexp" 9 + "strings" 10 + 11 + cgroups "github.com/containerd/cgroups/v3" 12 + "github.com/containerd/cgroups/v3/cgroup2" 13 + "github.com/prometheus/procfs" 14 + ) 15 + 16 + var ( 17 + cgroupInvalidChar = regexp.MustCompile(`[^a-zA-Z0-9\-_.]`) 18 + cgroupConsecutiveSep = regexp.MustCompile(`[-_.]{2,}`) 19 + ) 20 + 21 + const ( 22 + cgroupParentSelf = "self" 23 + supervisorCgroupName = "supervisor" 24 + ) 25 + 26 + type CgroupLimits struct { 27 + Enabled bool 28 + Parent *CgroupParent 29 + Name string 30 + MemoryMaxMiB int64 31 + SwapMaxMiB *int64 32 + PidsMax int64 33 + } 34 + 35 + type CgroupParent struct { 36 + root *cgroup2.Manager 37 + mountpoint string 38 + group string 39 + } 40 + 41 + type CgroupHandle struct { 42 + manager *cgroup2.Manager 43 + } 44 + 45 + func initCgroupParent(parent string, logger *slog.Logger) (*CgroupParent, error) { 46 + if parent == "" { 47 + parent = cgroupParentSelf 48 + } 49 + if cgroups.Mode() != cgroups.Unified { 50 + return nil, fmt.Errorf("microVM cgroups require cgroup v2 unified mode") 51 + } 52 + 53 + mountpoint, group, err := resolveCgroupParent(parent) 54 + if err != nil { 55 + return nil, err 56 + } 57 + if _, err := os.Stat(filepath.Join(mountpoint, strings.TrimPrefix(group, "/"))); err != nil { 58 + return nil, fmt.Errorf("stat cgroup parent %q:%q: %w", mountpoint, group, err) 59 + } 60 + 61 + root, err := cgroup2.Load(group, cgroup2.WithMountpoint(mountpoint)) 62 + if err != nil { 63 + return nil, fmt.Errorf("load cgroup parent %q:%q: %w", mountpoint, group, err) 64 + } 65 + 66 + if group != "/" { 67 + if err := moveParentProcesses(root, logger); err != nil { 68 + return nil, err 69 + } 70 + } 71 + 72 + if logger != nil { 73 + logger.Info("initialized microVM cgroup parent", "mountpoint", mountpoint, "group", group) 74 + } 75 + return &CgroupParent{root: root, mountpoint: mountpoint, group: group}, nil 76 + } 77 + 78 + func prepareCgroup(limits CgroupLimits, logger *slog.Logger) (*CgroupHandle, error) { 79 + if !limits.Enabled { 80 + return nil, nil 81 + } 82 + if limits.Parent == nil || limits.Parent.root == nil { 83 + return nil, fmt.Errorf("cgroup parent is not initialized") 84 + } 85 + name := sanitizeCgroupName(limits.Name) 86 + if name == "" { 87 + return nil, fmt.Errorf("cgroup name is empty") 88 + } 89 + 90 + manager, err := limits.Parent.root.NewChild(name, cgroupResources(limits)) 91 + if err != nil { 92 + return nil, fmt.Errorf("create cgroup %q: %w", name, err) 93 + } 94 + 95 + if logger != nil { 96 + logger.Info("created microVM cgroup", "name", name, "parentGroup", limits.Parent.group) 97 + } 98 + return &CgroupHandle{manager: manager}, nil 99 + } 100 + 101 + func cgroupResources(limits CgroupLimits) *cgroup2.Resources { 102 + resources := &cgroup2.Resources{} 103 + if limits.MemoryMaxMiB > 0 || limits.SwapMaxMiB != nil { 104 + memory := &cgroup2.Memory{} 105 + if limits.MemoryMaxMiB > 0 { 106 + maxBytes := limits.MemoryMaxMiB * 1024 * 1024 107 + memory.Max = &maxBytes 108 + } 109 + if limits.SwapMaxMiB != nil { 110 + swapBytes := *limits.SwapMaxMiB * 1024 * 1024 111 + memory.Swap = &swapBytes 112 + } 113 + oomGroup := true 114 + memory.OOMGroup = &oomGroup 115 + resources.Memory = memory 116 + } 117 + if limits.PidsMax > 0 { 118 + resources.Pids = &cgroup2.Pids{Max: limits.PidsMax} 119 + } 120 + return resources 121 + } 122 + 123 + func (h *CgroupHandle) AddProcess(pid int, logger *slog.Logger) error { 124 + if h == nil || h.manager == nil { 125 + return nil 126 + } 127 + if pid <= 0 { 128 + return fmt.Errorf("invalid pid %d", pid) 129 + } 130 + if err := h.manager.AddProc(uint64(pid)); err != nil { 131 + return fmt.Errorf("add pid %d to cgroup: %w", pid, err) 132 + } 133 + if logger != nil { 134 + logger.Info("added process to microVM cgroup", "pid", pid) 135 + } 136 + return nil 137 + } 138 + 139 + func (h *CgroupHandle) Close() error { 140 + if h == nil || h.manager == nil { 141 + return nil 142 + } 143 + return h.manager.Delete() 144 + } 145 + 146 + func resolveCgroupParent(parent string) (string, string, error) { 147 + mountpoint, err := cgroup2Mountpoint() 148 + if err != nil { 149 + return "", "", err 150 + } 151 + 152 + if parent == "" || parent == cgroupParentSelf { 153 + group, err := selfCgroupV2Path() 154 + if err != nil { 155 + return "", "", err 156 + } 157 + return mountpoint, group, nil 158 + } 159 + if !filepath.IsAbs(parent) { 160 + return "", "", fmt.Errorf("cgroup parent must be %q or an absolute delegated cgroupfs path: %q", cgroupParentSelf, parent) 161 + } 162 + 163 + cleanParent := filepath.Clean(parent) 164 + rel, err := filepath.Rel(mountpoint, cleanParent) 165 + if err != nil { 166 + return "", "", fmt.Errorf("resolve cgroup parent %q relative to cgroup2 mount %q: %w", cleanParent, mountpoint, err) 167 + } 168 + if rel == ".." || strings.HasPrefix(rel, "../") { 169 + return "", "", fmt.Errorf("cgroup parent %q is outside cgroup2 mount %q", cleanParent, mountpoint) 170 + } 171 + if rel == "." { 172 + return mountpoint, "/", nil 173 + } 174 + 175 + group := "/" + filepath.ToSlash(rel) 176 + if err := cgroup2.VerifyGroupPath(group); err != nil { 177 + return "", "", fmt.Errorf("invalid cgroup parent path %q: %w", group, err) 178 + } 179 + return mountpoint, group, nil 180 + } 181 + 182 + func cgroup2Mountpoint() (string, error) { 183 + mounts, err := procfs.GetMounts() 184 + if err != nil { 185 + return "", fmt.Errorf("read procfs mountinfo: %w", err) 186 + } 187 + for _, mount := range mounts { 188 + if mount.FSType == "cgroup2" { 189 + return mount.MountPoint, nil 190 + } 191 + } 192 + return "", fmt.Errorf("cgroup v2 mountpoint not found") 193 + } 194 + 195 + func selfCgroupV2Path() (string, error) { 196 + self, err := procfs.Self() 197 + if err != nil { 198 + return "", fmt.Errorf("open procfs self: %w", err) 199 + } 200 + groups, err := self.Cgroups() 201 + if err != nil { 202 + return "", fmt.Errorf("read procfs self cgroups: %w", err) 203 + } 204 + for _, group := range groups { 205 + if group.HierarchyID != 0 { 206 + continue 207 + } 208 + path := group.Path 209 + if path == "" { 210 + path = "/" 211 + } 212 + if err := cgroup2.VerifyGroupPath(path); err != nil { 213 + return "", fmt.Errorf("invalid self cgroup path %q: %w", path, err) 214 + } 215 + return path, nil 216 + } 217 + return "", fmt.Errorf("current process has no cgroup v2 hierarchy entry") 218 + } 219 + 220 + func moveParentProcesses(parent *cgroup2.Manager, logger *slog.Logger) error { 221 + supervisor, err := parent.NewChild(supervisorCgroupName, nil) 222 + if err != nil { 223 + return fmt.Errorf("create supervisor cgroup: %w", err) 224 + } 225 + 226 + procs, err := parent.Procs(false) 227 + if err != nil { 228 + return fmt.Errorf("list parent cgroup processes: %w", err) 229 + } 230 + for _, pid := range procs { 231 + if err := supervisor.AddProc(pid); err != nil { 232 + return fmt.Errorf("move pid %d to supervisor cgroup: %w", pid, err) 233 + } 234 + } 235 + 236 + if logger != nil && len(procs) > 0 { 237 + logger.Info("moved spindle processes to supervisor cgroup", "processes", len(procs)) 238 + } 239 + return nil 240 + } 241 + 242 + func sanitizeCgroupName(name string) string { 243 + name = cgroupInvalidChar.ReplaceAllLiteralString(name, "-") 244 + name = cgroupConsecutiveSep.ReplaceAllLiteralString(name, "-") 245 + return strings.Trim(name, "-_.") 246 + }
+40 -9
spindle/engines/microvm/engine.go
··· 36 36 type cleanupFunc func(context.Context) error 37 37 38 38 type Engine struct { 39 - l *slog.Logger 40 - cfg *config.Config 41 - db *db.DB 42 - agent *agentHub 39 + l *slog.Logger 40 + cfg *config.Config 41 + db *db.DB 42 + agent *agentHub 43 + scheduler *vmScheduler 44 + cgroupParent *CgroupParent 43 45 44 46 cleanupMu sync.Mutex 45 47 cleanup map[string][]cleanupFunc ··· 65 67 if err != nil { 66 68 return nil, err 67 69 } 70 + budget, max, err := newVMBudgetConfig(cfg.MicroVMPipelines) 71 + if err != nil { 72 + return nil, err 73 + } 74 + l.Info("initialized microVM workflow budget", "budget", budget.String(), "maxWorkflow", max.String()) 75 + 76 + var cgroupParent *CgroupParent 77 + if cfg.MicroVMPipelines.EnableCgroups { 78 + cgroupParent, err = initCgroupParent(cfg.MicroVMPipelines.CgroupParent, l) 79 + if err != nil { 80 + return nil, err 81 + } 82 + } 83 + 68 84 return &Engine{ 69 - l: l, 70 - cfg: cfg, 71 - db: d, 72 - agent: agent, 73 - cleanup: make(map[string][]cleanupFunc), 85 + l: l, 86 + cfg: cfg, 87 + db: d, 88 + agent: agent, 89 + scheduler: newVMScheduler(budget, max), 90 + cgroupParent: cgroupParent, 91 + cleanup: make(map[string][]cleanupFunc), 74 92 }, nil 75 93 } 76 94 ··· 208 226 EnableKVM: e.cfg.MicroVMPipelines.EnableKVM, 209 227 KeepWorkDir: true, 210 228 WorkDir: workDir, 229 + Cgroup: e.cgroupLimits(wid, state.ImageSpec), 211 230 }, l) 212 231 if err != nil { 213 232 return err ··· 407 426 delete(e.cleanup, key) 408 427 return fns 409 428 } 429 + 430 + func (e *Engine) cgroupLimits(wid models.WorkflowId, spec ImageSpec) CgroupLimits { 431 + cfg := e.cfg.MicroVMPipelines 432 + return CgroupLimits{ 433 + Enabled: cfg.EnableCgroups, 434 + Parent: e.cgroupParent, 435 + Name: "workflow-" + wid.String(), 436 + MemoryMaxMiB: resourcesForImage(spec).MemoryMiB, 437 + SwapMaxMiB: cfg.CgroupSwapMaxMiB, 438 + PidsMax: cfg.CgroupPidsMax, 439 + } 440 + }
+28 -8
spindle/engines/microvm/qemu.go
··· 44 44 SerialLogPath string 45 45 WorkDir string 46 46 VolumeBaseName string 47 + Cgroup CgroupLimits 47 48 } 48 49 49 50 type QEMUVMHandle struct { ··· 60 61 keepWorkDir bool 61 62 ownsWorkDir bool 62 63 qemuLogFile *os.File 64 + cgroup *CgroupHandle 63 65 slirpCmd *exec.Cmd 64 66 slirpExit *os.File 65 67 waitErr error ··· 191 193 cmd.Stdout = qemuLogFile 192 194 cmd.Stderr = qemuLogFile 193 195 196 + cgroup, err := prepareCgroup(cfg.Cgroup, logger) 197 + if err != nil { 198 + return nil, err 199 + } 200 + handle.cgroup = cgroup 201 + 194 202 logger.Info("starting qemu microvm", "cid", cid, "workDir", workDir, "serialLog", serialLogPath, "qmp", qmpPath) 195 203 if err := cmd.Start(); err != nil { 196 204 return nil, fmt.Errorf("starting qemu: %w", err) 197 205 } 198 206 handle.cmd = cmd 199 207 handle.Process = cmd.Process 200 - 201 - if slirpNet != nil { 202 - handle.slirpCmd, handle.slirpExit, err = slirpNet.Start(ctx, qemuLogFile, logger) 203 - if err != nil { 204 - return nil, err 205 - } 206 - } 207 - 208 208 handle.done = make(chan struct{}) 209 209 go func() { 210 210 err := cmd.Wait() ··· 214 214 close(handle.done) 215 215 }() 216 216 217 + if err := cgroup.AddProcess(cmd.Process.Pid, logger); err != nil { 218 + return nil, err 219 + } 220 + 221 + if slirpNet != nil { 222 + handle.slirpCmd, handle.slirpExit, err = slirpNet.Start(ctx, qemuLogFile, logger) 223 + if err != nil { 224 + return nil, err 225 + } 226 + if handle.slirpCmd != nil && handle.slirpCmd.Process != nil { 227 + if err := cgroup.AddProcess(handle.slirpCmd.Process.Pid, logger); err != nil { 228 + return nil, err 229 + } 230 + } 231 + } 232 + 217 233 qmpTimeout := cfg.BootTimeout 218 234 if qmpTimeout == 0 { 219 235 qmpTimeout = defaultQMPTimeout ··· 314 330 if h.qemuLogFile != nil { 315 331 closeErr = errors.Join(closeErr, h.qemuLogFile.Close()) 316 332 h.qemuLogFile = nil 333 + } 334 + if h.cgroup != nil { 335 + closeErr = errors.Join(closeErr, h.cgroup.Close()) 336 + h.cgroup = nil 317 337 } 318 338 if h.ownsWorkDir && !h.keepWorkDir && h.workDir != "" { 319 339 closeErr = errors.Join(closeErr, os.RemoveAll(h.workDir))
+2
spindle/engines/microvm/vm.go
··· 160 160 EnableKVM bool 161 161 KeepWorkDir bool 162 162 WorkDir string 163 + Cgroup CgroupLimits 163 164 164 165 QEMU QEMUVMConfig 165 166 ··· 348 349 SerialLogPath: serialLogPath, 349 350 WorkDir: cfg.WorkDir, 350 351 MkfsExt4: cfg.MkfsExt4, 352 + Cgroup: cfg.Cgroup, 351 353 }, logger) 352 354 case "firecracker": 353 355 return nil, fmt.Errorf("runner type %q not implemented yet", runnerType)
+33 -33
spindle/server.go
··· 41 41 ) 42 42 43 43 type Spindle struct { 44 - jc *jetstream.JetstreamClient 45 - tap *Tap 46 - embedTap *embeddedTap 47 - db *db.DB 48 - e *rbac.Enforcer 49 - l *slog.Logger 50 - n *notifier.Notifier 51 - engs map[string]models.Engine 52 - jq *queue.Queue 53 - cfg *config.Config 54 - ks *eventconsumer.Consumer 55 - res *idresolver.Resolver 56 - vault secrets.Manager 57 - motd []byte 58 - motdMu sync.RWMutex 59 - workflowSem chan struct{} 60 - rootCtx context.Context 44 + jc *jetstream.JetstreamClient 45 + tap *Tap 46 + embedTap *embeddedTap 47 + db *db.DB 48 + e *rbac.Enforcer 49 + l *slog.Logger 50 + n *notifier.Notifier 51 + engs map[string]models.Engine 52 + jq *queue.Queue 53 + cfg *config.Config 54 + ks *eventconsumer.Consumer 55 + res *idresolver.Resolver 56 + vault secrets.Manager 57 + motd []byte 58 + motdMu sync.RWMutex 59 + workflowSlotter engine.WorkflowSlotter 60 + rootCtx context.Context 61 61 } 62 62 63 63 // New creates a new Spindle server with the provided configuration and engines. ··· 104 104 jq := queue.NewQueue(cfg.Server.QueueSize, cfg.Server.MaxJobCount) 105 105 logger.Info("initialized queue", "queueSize", cfg.Server.QueueSize, "numWorkers", cfg.Server.MaxJobCount) 106 106 107 - workflowSem := make(chan struct{}, cfg.Server.MaxConcurrentWorkflows) 108 - logger.Info("initialized workflow semaphore", "maxConcurrentWorkflows", cfg.Server.MaxConcurrentWorkflows) 107 + workflowSlotter := engine.NewSemaphoreSlotter(cfg.Server.MaxConcurrentWorkflows) 108 + logger.Info("initialized default workflow slotter", "maxConcurrentWorkflows", cfg.Server.MaxConcurrentWorkflows) 109 109 110 110 collections := []string{ 111 111 tangled.SpindleMemberNSID, ··· 140 140 resolver := idresolver.DefaultResolver(cfg.Server.PlcUrl) 141 141 142 142 spindle := &Spindle{ 143 - jc: jc, 144 - e: e, 145 - db: d, 146 - l: logger, 147 - n: &n, 148 - engs: engines, 149 - jq: jq, 150 - cfg: cfg, 151 - res: resolver, 152 - vault: vault, 153 - motd: defaultMotd, 154 - workflowSem: workflowSem, 155 - rootCtx: ctx, 143 + jc: jc, 144 + e: e, 145 + db: d, 146 + l: logger, 147 + n: &n, 148 + engs: engines, 149 + jq: jq, 150 + cfg: cfg, 151 + res: resolver, 152 + vault: vault, 153 + motd: defaultMotd, 154 + workflowSlotter: workflowSlotter, 155 + rootCtx: ctx, 156 156 } 157 157 158 158 err = e.AddSpindle(rbacDomain) ··· 479 479 480 480 ok := s.jq.Enqueue(queue.Job{ 481 481 Run: func() error { 482 - engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.workflowSem, ctx, &models.Pipeline{ 482 + engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.workflowSlotter, ctx, &models.Pipeline{ 483 483 RepoDid: repoDid, 484 484 Workflows: workflows, 485 485 }, pipelineId)