From 01c08baab5d28892aee242d98c19d4001a851c51 Mon Sep 17 00:00:00 2001 From: dawn Date: Fri, 18 Sep 2026 14:03:05 +0300 Subject: [PATCH] spindle/engines/nixery: give each workflow its own seat network Signed-off-by: dawn --- spindle/engines/nixery/engine.go | 27 +++-- spindle/engines/nixery/seats.go | 129 +++++++++++++++++++++++ spindle/engines/nixery/seats_test.go | 147 +++++++++++++++++++++++++++ 3 files changed, 289 insertions(+), 14 deletions(-) create mode 100644 spindle/engines/nixery/seats.go create mode 100644 spindle/engines/nixery/seats_test.go diff --git a/spindle/engines/nixery/engine.go b/spindle/engines/nixery/engine.go index 484ce6ba1..78dc0b554 100644 --- a/spindle/engines/nixery/engine.go +++ b/spindle/engines/nixery/engine.go @@ -16,7 +16,6 @@ import ( "github.com/docker/docker/api/types/container" "github.com/docker/docker/api/types/image" "github.com/docker/docker/api/types/mount" - "github.com/docker/docker/api/types/network" "github.com/docker/docker/client" "github.com/docker/docker/pkg/stdcopy" "gopkg.in/yaml.v3" @@ -43,6 +42,7 @@ type Engine struct { cfg *config.Config slotter engine.WorkflowSlotter + seats *seatNetworks cleanupMu sync.Mutex cleanup map[string][]cleanupFunc @@ -206,6 +206,11 @@ func New(ctx context.Context, cfg *config.Config) (*Engine, error) { slotter: engine.NewSemaphoreSlotter(cfg.NixeryPipelines.MaxConcurrentWorkflows), } + namespace := networkNamespace(cfg) + e.seats = newSeatNetworks(namespace, func(ctx context.Context, name string) error { + return e.ensureSeatNetwork(ctx, name, namespace) + }) + e.cleanup = make(map[string][]cleanupFunc) return e, nil @@ -264,18 +269,14 @@ func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *m return err } - /// -------------------------NETWORK CREATION--------------------------------------- - _, err = e.docker.NetworkCreate(ctx, networkName(wid), network.CreateOptions{ - Driver: "bridge", - }) + /// -------------------------WORKFLOW NETWORK--------------------------------------- + seat, releaseSeat, err := e.seats.take(ctx) if err != nil { - return err + return fmt.Errorf("claiming workflow network: %w", err) } - e.registerCleanup(wid, func(ctx context.Context) error { - if err := e.docker.NetworkRemove(ctx, networkName(wid)); err != nil { - return fmt.Errorf("removing network: %w", err) - } + e.registerCleanup(wid, func(context.Context) error { + releaseSeat() return nil }) @@ -304,7 +305,7 @@ func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *m } /// -------------------------CONTAINER CREATION------------------------------------- - l.Info("creating container") + l.Info("creating container", "network", seat) wfLogger.DataWriter(setupStepIdx, "stdout").Write([]byte("Creating container...")) extraHosts := []string{"host.docker.internal:host-gateway"} @@ -331,6 +332,7 @@ func (e *Engine) SetupWorkflow(ctx context.Context, wid models.WorkflowId, wf *m // TODO(winter): investigate whether environment variables passed here // get propagated to ContainerExec processes }, &container.HostConfig{ + NetworkMode: container.NetworkMode(seat), Mounts: append([]mount.Mount{ { Type: mount.TypeTmpfs, @@ -590,9 +592,6 @@ func (e *Engine) drainCleanups(wid models.WorkflowId) []cleanupFunc { return fns } -func networkName(wid models.WorkflowId) string { - return fmt.Sprintf("workflow-network-%s", wid) -} func (e *Engine) QuotaResources(wf *models.Workflow) quota.Resources { return quota.Resources{ quota.ResourceWorkflows: 1, diff --git a/spindle/engines/nixery/seats.go b/spindle/engines/nixery/seats.go new file mode 100644 index 000000000..7aebb591c --- /dev/null +++ b/spindle/engines/nixery/seats.go @@ -0,0 +1,129 @@ +package nixery + +import ( + "context" + "fmt" + "strings" + "sync" + + "github.com/docker/docker/api/types/network" + "github.com/docker/docker/client" + "tangled.org/core/spindle/config" +) + +const ( + // one bridge per running workflow, so concurrent workflows can't reach + // each other + seatNetworkPrefix = "spindle-wf" + // seats are host-global, so a seat records which engine made it and we only + // adopt our own + seatOwnerLabel = "sh.tangled.pipeline/engine" +) + +// seat networks are recycled, so an engine owns at most as many as its peak +// concurrency, not one per workflow it has run +type seatNetworks struct { + owner string + ensure func(ctx context.Context, name string) error + + mu sync.Mutex + busy []bool + free []int +} + +func newSeatNetworks(owner string, ensure func(context.Context, string) error) *seatNetworks { + return &seatNetworks{owner: owner, ensure: ensure} +} + +func (s *seatNetworks) take(ctx context.Context) (string, func(), error) { + seat := s.claim() + name := s.name(seat) + if err := s.ensure(ctx, name); err != nil { + s.release(seat) + return "", nil, err + } + return name, func() { s.release(seat) }, nil +} + +func (s *seatNetworks) claim() int { + s.mu.Lock() + defer s.mu.Unlock() + + if n := len(s.free); n > 0 { + seat := s.free[n-1] + s.free = s.free[:n-1] + s.busy[seat] = true + return seat + } + s.busy = append(s.busy, true) + return len(s.busy) - 1 +} + +// release only touches memory, so a canceled workflow still frees its seat, and +// a second release can't hand the seat to two workflows +func (s *seatNetworks) release(seat int) { + s.mu.Lock() + defer s.mu.Unlock() + + if seat >= len(s.busy) || !s.busy[seat] { + return + } + s.busy[seat] = false + s.free = append(s.free, seat) +} + +func (s *seatNetworks) name(seat int) string { + return fmt.Sprintf("%s-%s-%d", seatNetworkPrefix, s.owner, seat) +} + +func (e *Engine) ensureSeatNetwork(ctx context.Context, name, owner string) error { + docker, err := e.ensureDocker() + if err != nil { + return err + } + + existing, err := docker.NetworkInspect(ctx, name, network.InspectOptions{}) + switch { + case err == nil: + if held := existing.Labels[seatOwnerLabel]; held != owner { + return fmt.Errorf("network %s is owned by %q, want %q", name, held, owner) + } + return nil + case client.IsErrNotFound(err): + if _, err := docker.NetworkCreate(ctx, name, network.CreateOptions{ + Driver: "bridge", + Labels: map[string]string{seatOwnerLabel: owner}, + }); err != nil { + return fmt.Errorf("creating network %s: %w", name, err) + } + return nil + default: + return fmt.Errorf("inspecting network %s: %w", name, err) + } +} + +// seats are host-global, so the pool is namespaced by the engine process's own +// name +func networkNamespace(cfg *config.Config) string { + name := cfg.Logging.ServiceName + if name == "" { + name = cfg.Tracing.ServiceName + } + if name = sanitizeNamespace(name); name != "" { + return name + } + return "spindle" +} + +func sanitizeNamespace(name string) string { + return strings.Trim(strings.Map(func(r rune) rune { + switch { + case r >= 'a' && r <= 'z', r >= '0' && r <= '9', r == '-', r == '_', r == '.': + return r + case r >= 'A' && r <= 'Z': + return r + ('a' - 'A') + default: + return '-' + } + }, name), "-") +} diff --git a/spindle/engines/nixery/seats_test.go b/spindle/engines/nixery/seats_test.go new file mode 100644 index 000000000..192b44f80 --- /dev/null +++ b/spindle/engines/nixery/seats_test.go @@ -0,0 +1,147 @@ +package nixery + +import ( + "context" + "errors" + "testing" + + "tangled.org/core/spindle/config" +) + +// records what the pool asked for, or refuses it +type fakeDocker struct { + ensured []string + fail error +} + +func (f *fakeDocker) ensure(_ context.Context, name string) error { + if f.fail != nil { + return f.fail + } + f.ensured = append(f.ensured, name) + return nil +} + +func (f *fakeDocker) networks() map[string]int { + counts := map[string]int{} + for _, name := range f.ensured { + counts[name]++ + } + return counts +} + +func TestSeatNetworksGiveConcurrentWorkflowsDifferentNetworks(t *testing.T) { + docker := &fakeDocker{} + seats := newSeatNetworks("hel-1", docker.ensure) + + first, _, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + second, _, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + + if first == second { + t.Fatalf("two concurrent workflows share network %q", first) + } + if want := "spindle-wf-hel-1-1"; second != want { + t.Fatalf("second seat = %q, want %q", second, want) + } +} + +func TestSeatNetworksReuseReleasedSeats(t *testing.T) { + docker := &fakeDocker{} + seats := newSeatNetworks("hel-1", docker.ensure) + + name, release, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + release() + + again, _, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + if again != name { + t.Fatalf("released seat was not reused: got %q, want %q", again, name) + } + if counts := docker.networks(); len(counts) != 1 || counts[name] != 2 { + t.Fatalf("pool networks = %v, want only %q", counts, name) + } +} + +func TestSeatNetworksReleaseIsIdempotent(t *testing.T) { + docker := &fakeDocker{} + seats := newSeatNetworks("hel-1", docker.ensure) + + _, release, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + release() + release() + + first, _, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + second, _, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + if first == second { + t.Fatalf("a double release handed seat %q to two workflows", first) + } +} + +func TestSeatNetworksReturnSeatsWhoseNetworkCannotBeCreated(t *testing.T) { + docker := &fakeDocker{fail: errors.New("all predefined address pools have been fully subnetted")} + seats := newSeatNetworks("hel-1", docker.ensure) + + if _, _, err := seats.take(context.Background()); err == nil { + t.Fatal("take succeeded although the network could not be created") + } + + docker.fail = nil + name, _, err := seats.take(context.Background()) + if err != nil { + t.Fatal(err) + } + if want := "spindle-wf-hel-1-0"; name != want { + t.Fatalf("seat = %q, want %q", name, want) + } +} + +func TestNetworkNamespaceNamesThePoolAfterTheEngineProcess(t *testing.T) { + for _, tc := range []struct { + name string + cfg config.Config + want string + }{ + {"defaults", config.Config{}, "spindle"}, + { + "logging service name", + config.Config{Logging: config.Logging{ServiceName: "spindle-executor-hel-1"}}, + "spindle-executor-hel-1", + }, + { + "tracing service name fallback", + config.Config{Tracing: config.Tracing{ServiceName: "spindle-executor-hel-1"}}, + "spindle-executor-hel-1", + }, + { + "unsanitized service name", + config.Config{Logging: config.Logging{ServiceName: "Executor Hel 1!"}}, + "executor-hel-1", + }, + } { + t.Run(tc.name, func(t *testing.T) { + if got := networkNamespace(&tc.cfg); got != tc.want { + t.Fatalf("networkNamespace() = %q, want %q", got, tc.want) + } + }) + } +} -- 2.51.2