From 07c22b994634ee08b014705e7d11fddd195d1b27 Mon Sep 17 00:00:00 2001 From: Paolo Date: Wed, 30 Sep 2026 22:00:50 +0000 Subject: [PATCH] feat: cluster-wide Swarm discovery behind DOCKTAIL_DISCOVERY=swarm Container discovery is node-local: /containers returns only the containers running on the Docker node that serves it, on a manager exactly as on a worker. An agent therefore advertises only the services it is colocated with, which in a Swarm cluster means one agent per node. Add a second discovery source that reads the whole cluster. /services is manager-only and returns every service with its labels, and each service carries Endpoint.VirtualIPs: the VIP the routing mesh keeps alive and load-balanced across that service's tasks, on the overlay its tasks are attached to. A VIP is reachable from any node on that overlay, so one agent covers the cluster as long as the overlay carries the traffic. Labels are read from Spec.TaskTemplate.ContainerSpec.Labels, where both `labels:` and `deploy.labels:` land, so a Swarm service is labelled exactly like the equivalent standalone container and every existing label applies unchanged. Parsing is shared: the label translation is factored out of the inspect-and-parse path so both sources run the same code. The destination is the service VIP rather than a container IP, so the mode requires the agent and the exposed services to share an overlay. When a service sits on several networks, DockTail prefers one it is also attached to, and DOCKTAIL_SWARM_NETWORK pins the choice. A VIP only exists in vip mode, so a dnsrr service or one with no network is reported and skipped. Event watching stays local, so a change on another node lands on the next RECONCILE_INTERVAL rather than immediately. --- docker/client.go | 25 +++- docker/swarm.go | 265 ++++++++++++++++++++++++++++++++++++++ docker/swarm_test.go | 271 +++++++++++++++++++++++++++++++++++++++ docs/07-reference.md | 46 +++++++ main.go | 19 +++ reconciler/reconciler.go | 65 +++++++++- 6 files changed, 685 insertions(+), 6 deletions(-) create mode 100644 docker/swarm.go create mode 100644 docker/swarm_test.go diff --git a/docker/client.go b/docker/client.go index 44c56b4..a3ac133 100644 --- a/docker/client.go +++ b/docker/client.go @@ -54,6 +54,12 @@ type Client struct { // own network IP (see sharesNetworkWith). Resolved at most once. selfNetOnce sync.Once selfNetIDs map[string]struct{} + + // netNames caches network ID -> network name, used by the swarm discovery + // path to resolve the operator's DOCKTAIL_SWARM_NETWORK preference. Network + // names do not change under a running cluster, and the lookup is otherwise + // repeated for every candidate VIP on every reconcile. + netNames sync.Map } // NewClient creates a new Docker client @@ -792,9 +798,7 @@ func (c *Client) parseContainer(ctx context.Context, containerID string, labels } func (c *Client) parseContainerServices(ctx context.Context, containerID string, labels map[string]string) ([]*apptypes.ContainerService, error) { - serviceEnabled := isServiceEnabled(labels) - funnelEnabled := isFunnelEnabled(labels) - if !serviceEnabled && !funnelEnabled { + if !isManagedContainer(labels) { return nil, nil } @@ -804,6 +808,21 @@ func (c *Client) parseContainerServices(ctx context.Context, containerID string, return nil, fmt.Errorf("failed to inspect container: %w", err) } + return c.parseFromInspect(ctx, containerID, inspect, labels) +} + +// parseFromInspect runs the label -> service translation against an already +// resolved inspect response. Keeping it separate from parseContainerServices +// lets the swarm source (docker/swarm.go) reuse the whole label surface for +// Swarm services, whose "container" is a service VIP rather than a locally +// inspectable container. +func (c *Client) parseFromInspect(ctx context.Context, containerID string, inspect container.InspectResponse, labels map[string]string) ([]*apptypes.ContainerService, error) { + serviceEnabled := isServiceEnabled(labels) + funnelEnabled := isFunnelEnabled(labels) + if !serviceEnabled && !funnelEnabled { + return nil, nil + } + containerName := strings.TrimPrefix(inspect.Name, "/") cctx := &containerCtx{ diff --git a/docker/swarm.go b/docker/swarm.go new file mode 100644 index 0000000..7d34dac --- /dev/null +++ b/docker/swarm.go @@ -0,0 +1,265 @@ +package docker + +import ( + "bytes" + "context" + "fmt" + "net" + "os" + "sort" + "strings" + + "github.com/docker/docker/api/types/container" + network "github.com/docker/docker/api/types/network" + "github.com/docker/docker/api/types/swarm" + "github.com/rs/zerolog/log" + + apptypes "github.com/marvinvr/docktail/types" +) + +// Swarm discovery mode. +// +// GetEnabledContainers only sees containers running on the agent's own Docker +// node — the /containers endpoint is node-local in a Swarm cluster, on manager +// nodes as much as workers. An agent can therefore only advertise services whose +// containers are colocated with it, which forces one agent per node for a +// cluster-wide setup. +// +// The Swarm API does offer a cluster view, but at service granularity: +// /services returns every service in the cluster together with its labels, and +// each service carries Endpoint.VirtualIPs — the VIP the cluster's routing mesh +// already keeps alive and load-balanced across that service's tasks, on the +// overlay network the tasks are attached to. A VIP is reachable from any node +// attached to that overlay, so one agent can advertise a whole cluster as long +// as the overlay carries the traffic. +// +// Labels come from Spec.TaskTemplate.ContainerSpec.Labels, which is where both +// `labels:` and `deploy.labels:` from a compose file land, so a service is +// labelled exactly like the equivalent standalone container would be. + +// swarmNetworkEnv names the Docker network whose VIP should be advertised when +// a service is attached to several. Optional: without it, the lowest VIP +// (deterministically, by address) wins. +const swarmNetworkEnv = "DOCKTAIL_SWARM_NETWORK" + +// GetEnabledSwarmServices returns the Swarm services managed by DockTail, +// discovered across the whole cluster rather than only on this node. +// +// It requires a manager endpoint: /services is manager-only. On a worker-only +// agent it fails, and the caller should treat that as "cannot run in swarm +// mode" rather than retrying. +func (c *Client) GetEnabledSwarmServices(ctx context.Context) ([]*apptypes.ContainerService, error) { + services, err := c.cli.ServiceList(ctx, swarm.ServiceListOptions{}) + if err != nil { + return nil, fmt.Errorf("failed to list swarm services (is this Docker endpoint a swarm manager?): %w", err) + } + + var managed []*apptypes.ContainerService + for _, svc := range services { + labels := swarmServiceLabels(svc) + if !isManagedContainer(labels) { + continue + } + + parsed, err := c.parseSwarmService(ctx, svc, labels) + if err != nil { + log.Warn(). + Err(err). + Str("service", svc.Spec.Name). + Msg("Failed to parse swarm service, skipping") + continue + } + managed = append(managed, parsed...) + } + + return managed, nil +} + +// swarmServiceLabels returns the labels a Swarm service should be configured +// from: the task template's container labels, which is where compose `labels:` +// and `deploy.labels:` both land. Returns nil for a service with no task +// template (an externally-managed service). +func swarmServiceLabels(svc swarm.Service) map[string]string { + if svc.Spec.TaskTemplate.ContainerSpec == nil { + return nil + } + return svc.Spec.TaskTemplate.ContainerSpec.Labels +} + +// parseSwarmService turns one Swarm service into its Tailscale services. The +// service VIP stands in for the container IP: it is the cluster-stable address +// for those tasks and the one an agent on any node can dial. +func (c *Client) parseSwarmService(ctx context.Context, svc swarm.Service, labels map[string]string) ([]*apptypes.ContainerService, error) { + networkID, vip, err := c.selectServiceVIP(ctx, svc) + if err != nil { + return nil, err + } + + // A Swarm service is network_mode: none only when it has no network at all, + // in which case there is no VIP to advertise and selectServiceVIP has + // already failed. host-mode tasks are reachable through the service's + // published ports rather than a VIP; that is the one case where direct mode + // cannot work, so it is rejected here with the same guidance the container + // path gives. + inspect := container.InspectResponse{ + ContainerJSONBase: &container.ContainerJSONBase{ + ID: svc.ID, + Name: "/" + svc.Spec.Name, + HostConfig: &container.HostConfig{ + NetworkMode: container.NetworkMode(swarmNetworkMode(svc)), + }, + }, + Config: &container.Config{Labels: labels}, + NetworkSettings: &container.NetworkSettings{ + Networks: map[string]*network.EndpointSettings{ + networkNameResolver(c, networkID): { + NetworkID: networkID, + IPAddress: vip, + }, + }, + }, + } + + services, err := c.parseFromInspect(ctx, svc.ID, inspect, labels) + if err != nil { + return nil, err + } + return withCloudLabels(services, labels), nil +} + +// selectServiceVIP picks the overlay VIP to advertise for a service. A service +// attached to several networks gets one VIP per network; only the one on a +// network the agent is also attached to is dialable, so the agent's own +// networks win, then an operator-pinned network, then a deterministic choice. +func (c *Client) selectServiceVIP(ctx context.Context, svc swarm.Service) (networkID, vip string, err error) { + // Resolve the networks the agent is attached to so the chosen VIP is one + // this agent can actually dial, then defer to the pure selection logic. + c.ensureSelfNetIDs(ctx) + return c.selectServiceVIPFrom(svc, c.selfNetIDs) +} + +// selectServiceVIPFrom is selectServiceVIP with the set of network IDs the +// agent is attached to passed in, so the choice logic can be exercised without +// a live Docker daemon. +func (c *Client) selectServiceVIPFrom(svc swarm.Service, sharedNetIDs map[string]struct{}) (networkID, vip string, err error) { + if len(svc.Endpoint.VirtualIPs) == 0 { + return "", "", fmt.Errorf("service %s has no virtual IP: a service only gets one in vip mode "+ + "(replicated or global). Publish the port or use a network the service is attached to", svc.Spec.Name) + } + + candidates := make([]swarm.EndpointVirtualIP, 0, len(svc.Endpoint.VirtualIPs)) + candidates = append(candidates, svc.Endpoint.VirtualIPs...) + sort.Slice(candidates, func(i, j int) bool { + return lessIP(candidates[i].Addr, candidates[j].Addr) + }) + + if preferred := strings.TrimSpace(getSwarmNetworkEnv()); preferred != "" { + for _, v := range candidates { + if networkNameResolver(c, v.NetworkID) == preferred { + return v.NetworkID, stripCIDR(v.Addr), nil + } + } + return "", "", fmt.Errorf("service %s is not attached to network %q (attached to: %v)", + svc.Spec.Name, preferred, c.swarmNetworkNames(svc)) + } + + if len(sharedNetIDs) > 0 { + for _, v := range candidates { + if _, ok := sharedNetIDs[v.NetworkID]; ok { + return v.NetworkID, stripCIDR(v.Addr), nil + } + } + } + + return candidates[0].NetworkID, stripCIDR(candidates[0].Addr), nil +} + +// swarmNetworkEnvVar is a package-level indirection so tests can pin the +// operator's network preference without touching the process environment. +var swarmNetworkEnvVar = func() string { return os.Getenv(swarmNetworkEnv) } + +func getSwarmNetworkEnv() string { return swarmNetworkEnvVar() } + +// swarmNetworkNames resolves the Docker network names a service is attached to. +// The alias in a service's NetworkAttachmentConfig is the DNS name the *service* +// gets on that network, not the network's own name, so it cannot answer +// DOCKTAIL_SWARM_NETWORK. Names come from inspecting each attached network. +// A network that cannot be inspected (one the agent has no permission for, or +// that has since been removed) falls back to its ID so the error message names +// something the operator can look up. +func (c *Client) swarmNetworkNames(svc swarm.Service) []string { + names := make([]string, 0, len(svc.Spec.TaskTemplate.Networks)) + for _, n := range svc.Spec.TaskTemplate.Networks { + names = append(names, networkNameResolver(c, n.Target)) + } + sort.Strings(names) + return names +} + +// networkNameResolver maps a Docker network ID to its network name. Overridden +// in tests so the selection logic can be exercised without a live daemon. +var networkNameResolver = func(c *Client, networkID string) string { return c.networkName(networkID) } + +// networkName resolves one network ID to its Docker network name, caching the +// answer for the life of the agent. Names rarely change, and the lookup costs +// an API call per candidate per reconcile otherwise. +func (c *Client) networkName(networkID string) string { + if networkID == "host" { + return "host" + } + if name, ok := c.netNames.Load(networkID); ok { + return name.(string) //nolint:forcetypeassert // only strings are ever stored + } + + name := networkID + if inspect, _, err := c.cli.NetworkInspectWithRaw(context.Background(), networkID, network.InspectOptions{}); err == nil && inspect.Name != "" { + name = inspect.Name + } else { + log.Debug(). + Err(err). + Str("network_id", networkID). + Msg("Could not resolve network name; using its ID") + } + + c.netNames.Store(networkID, name) + return name +} + +// swarmNetworkMode reports the effective network mode of a Swarm service's +// tasks: "host" when every attached network is the host network, "" otherwise. +func swarmNetworkMode(svc swarm.Service) string { + networks := svc.Spec.TaskTemplate.Networks + if len(networks) == 0 { + return "" + } + for _, n := range networks { + if n.Target != "host" { + return "" + } + } + return "host" +} + +// stripCIDR turns "10.0.4.7/24" into "10.0.4.7". Swarm reports VirtualIPs in +// CIDR form; everything downstream (tailscale serve, the health check) wants a +// bare address. +func stripCIDR(addr string) string { + if ip, _, err := net.ParseCIDR(addr); err == nil { + return ip.String() + } + if host, _, err := net.SplitHostPort(addr); err == nil { + return host + } + return addr +} + +// lessIP orders two CIDR addresses numerically when possible, falling back to +// string order so the choice stays deterministic either way. +func lessIP(a, b string) bool { + ipa := net.ParseIP(stripCIDR(a)) + ipb := net.ParseIP(stripCIDR(b)) + if ipa == nil || ipb == nil { + return a < b + } + return bytes.Compare(ipa, ipb) < 0 +} diff --git a/docker/swarm_test.go b/docker/swarm_test.go new file mode 100644 index 0000000..2e7c517 --- /dev/null +++ b/docker/swarm_test.go @@ -0,0 +1,271 @@ +package docker + +import ( + "testing" + + "github.com/docker/docker/api/types/swarm" +) + +// stubNetworkNames makes network-ID -> name resolution deterministic without a +// live Docker daemon. IDs in the tests are netA/netB, so name them by a lookup +// table the test supplies. +func stubNetworkNames(t *testing.T, names map[string]string) { + t.Helper() + prev := networkNameResolver + networkNameResolver = func(_ *Client, networkID string) string { + if name, ok := names[networkID]; ok { + return name + } + return networkID + } + t.Cleanup(func() { networkNameResolver = prev }) +} + +func newTestService(networks []swarm.NetworkAttachmentConfig, vips []swarm.EndpointVirtualIP) swarm.Service { + return swarm.Service{ + Spec: swarm.ServiceSpec{ + Annotations: swarm.Annotations{Name: "entertainment_plex"}, + TaskTemplate: swarm.TaskSpec{Networks: networks}, + }, + Endpoint: swarm.Endpoint{VirtualIPs: vips}, + } +} + +func TestStripCIDR(t *testing.T) { + tests := []struct { + name string + in string + want string + }{ + {name: "strips the prefix length", in: "10.0.4.7/24", want: "10.0.4.7"}, + {name: "passes a bare address through", in: "10.0.4.7", want: "10.0.4.7"}, + {name: "keeps a /128 host address", in: "fd7a::1/128", want: "fd7a::1"}, + {name: "drops a port", in: "10.0.4.7:8080", want: "10.0.4.7"}, + {name: "returns garbage unchanged", in: "not-an-ip", want: "not-an-ip"}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := stripCIDR(tt.in); got != tt.want { + t.Errorf("stripCIDR(%q) = %q, want %q", tt.in, got, tt.want) + } + }) + } +} + +func TestLessIP(t *testing.T) { + // Numeric, not lexicographic: 10.0.4.9 must sort before 10.0.4.10. + if !lessIP("10.0.4.9/24", "10.0.4.10/24") { + t.Error("10.0.4.9 should sort before 10.0.4.10 (numeric, not string order)") + } + if lessIP("10.0.4.10/24", "10.0.4.9/24") { + t.Error("10.0.4.10 should not sort before 10.0.4.9") + } + if !lessIP("10.0.4.2/24", "10.0.10.2/24") { + t.Error("10.0.4.2 should sort before 10.0.10.2") + } + // Unparseable addresses fall back to string order, which must still be a + // strict ordering so sort.Slice does not misbehave. + if lessIP("aaa", "bbb") == lessIP("bbb", "aaa") { + t.Error("string fallback should order aaa before bbb") + } +} + +func TestSwarmServiceLabels(t *testing.T) { + tests := []struct { + name string + svc swarm.Service + want map[string]string + }{ + { + name: "reads the task template container labels", + svc: swarm.Service{ + Spec: swarm.ServiceSpec{ + TaskTemplate: swarm.TaskSpec{ + ContainerSpec: &swarm.ContainerSpec{ + Labels: map[string]string{ + "docktail.service.enable": "true", + "docktail.service.name": "plex", + }, + }, + }, + }, + }, + want: map[string]string{ + "docktail.service.enable": "true", + "docktail.service.name": "plex", + }, + }, + { + name: "returns nil when there is no container spec", + svc: swarm.Service{Spec: swarm.ServiceSpec{TaskTemplate: swarm.TaskSpec{}}}, + want: nil, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := swarmServiceLabels(tt.svc) + if len(got) != len(tt.want) { + t.Fatalf("swarmServiceLabels() = %v, want %v", got, tt.want) + } + for k, v := range tt.want { + if got[k] != v { + t.Errorf("label %q = %q, want %q", k, got[k], v) + } + } + }) + } +} + +func TestSwarmNetworkNames(t *testing.T) { + stubNetworkNames(t, map[string]string{"netA": "tailscale"}) + + svc := swarm.Service{Spec: swarm.ServiceSpec{TaskTemplate: swarm.TaskSpec{ + Networks: []swarm.NetworkAttachmentConfig{ + {Target: "netA", Aliases: []string{"entertainment_plex"}}, + {Target: "netZ", Aliases: []string{"entertainment_plex"}}, + }, + }}} + + got := (&Client{}).swarmNetworkNames(svc) + // Sorted, so the unresolved ID ("netZ") precedes the resolved name. + want := []string{"netZ", "tailscale"} + + if len(got) != len(want) { + t.Fatalf("swarmNetworkNames() = %v, want %v", got, want) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("name[%d] = %q, want %q", i, got[i], want[i]) + } + } +} + +func TestSwarmNetworkMode(t *testing.T) { + tests := []struct { + name string + svc swarm.Service + want string + }{ + { + name: "overlay attachment is not host mode", + svc: swarm.Service{Spec: swarm.ServiceSpec{TaskTemplate: swarm.TaskSpec{ + Networks: []swarm.NetworkAttachmentConfig{{Target: "netid", Aliases: []string{"tailscale"}}}, + }}}, + want: "", + }, + { + name: "every attachment on host means host mode", + svc: swarm.Service{Spec: swarm.ServiceSpec{TaskTemplate: swarm.TaskSpec{ + Networks: []swarm.NetworkAttachmentConfig{{Target: "host"}}, + }}}, + want: "host", + }, + { + name: "no attachment is not host mode", + svc: swarm.Service{Spec: swarm.ServiceSpec{TaskTemplate: swarm.TaskSpec{}}}, + want: "", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := swarmNetworkMode(tt.svc); got != tt.want { + t.Errorf("swarmNetworkMode() = %q, want %q", got, tt.want) + } + }) + } +} + +func TestSelectServiceVIP(t *testing.T) { + c := &Client{} + + t.Run("no virtual IP is an error, not a silent skip", func(t *testing.T) { + _, _, err := c.selectServiceVIPFrom(newTestService(nil, nil), nil) + if err == nil { + t.Fatal("expected an error for a service without a VIP") + } + }) + + t.Run("picks the lowest VIP when the agent shares no network", func(t *testing.T) { + stubNetworkNames(t, map[string]string{"netA": "appnet"}) + svc := newTestService( + []swarm.NetworkAttachmentConfig{{Target: "netA", Aliases: []string{"plex"}}}, + []swarm.EndpointVirtualIP{{NetworkID: "netA", Addr: "10.0.9.20/24"}}, + ) + id, vip, err := c.selectServiceVIPFrom(svc, nil) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if id != "netA" || vip != "10.0.9.20" { + t.Errorf("got (%q, %q), want (\"netA\", \"10.0.9.20\")", id, vip) + } + }) + + t.Run("prefers a network the agent is attached to", func(t *testing.T) { + stubNetworkNames(t, map[string]string{"netA": "appnet", "netB": "tailscale"}) + svc := newTestService( + []swarm.NetworkAttachmentConfig{ + {Target: "netA", Aliases: []string{"plex"}}, + {Target: "netB", Aliases: []string{"plex"}}, + }, + []swarm.EndpointVirtualIP{ + {NetworkID: "netA", Addr: "10.0.4.2/24"}, + {NetworkID: "netB", Addr: "10.0.10.9/24"}, + }, + ) + id, vip, err := c.selectServiceVIPFrom(svc, map[string]struct{}{"netB": {}}) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if id != "netB" || vip != "10.0.10.9" { + t.Errorf("got (%q, %q), want (\"netB\", \"10.0.10.9\") — the shared network must win", id, vip) + } + }) + + t.Run("operator-pinned network overrides the shared-network preference", func(t *testing.T) { + stubNetworkNames(t, map[string]string{"netA": "appnet", "netB": "tailscale"}) + t.Setenv(swarmNetworkEnv, "appnet") + svc := newTestService( + []swarm.NetworkAttachmentConfig{ + {Target: "netA", Aliases: []string{"plex"}}, + {Target: "netB", Aliases: []string{"plex"}}, + }, + []swarm.EndpointVirtualIP{ + {NetworkID: "netA", Addr: "10.0.4.2/24"}, + {NetworkID: "netB", Addr: "10.0.10.9/24"}, + }, + ) + id, vip, err := c.selectServiceVIPFrom(svc, map[string]struct{}{"netB": {}}) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if id != "netA" || vip != "10.0.4.2" { + t.Errorf("got (%q, %q), want (\"netA\", \"10.0.4.2\")", id, vip) + } + }) + + t.Run("pinned network the service is not on is an error", func(t *testing.T) { + stubNetworkNames(t, map[string]string{"netA": "appnet"}) + t.Setenv(swarmNetworkEnv, "absent") + svc := newTestService( + []swarm.NetworkAttachmentConfig{{Target: "netA", Aliases: []string{"plex"}}}, + []swarm.EndpointVirtualIP{{NetworkID: "netA", Addr: "10.0.4.2/24"}}, + ) + if _, _, err := c.selectServiceVIPFrom(svc, nil); err == nil { + t.Fatal("expected an error when the pinned network is not attached") + } + }) + + t.Run("host network has no VIP and is reported as such", func(t *testing.T) { + stubNetworkNames(t, nil) + svc := newTestService( + []swarm.NetworkAttachmentConfig{{Target: "host"}}, + nil, + ) + if _, _, err := c.selectServiceVIPFrom(svc, nil); err == nil { + t.Fatal("expected an error for a host-networked service with no VIP") + } + }) +} diff --git a/docs/07-reference.md b/docs/07-reference.md index 5d6b943..fa42993 100644 --- a/docs/07-reference.md +++ b/docs/07-reference.md @@ -16,6 +16,8 @@ Use this section when checking exact configuration names, defaults, and supporte | `SKIP_SHUTDOWN_CLEANUP` | `false` | When `true`, DockTail leaves its services and Funnels advertised on shutdown instead of draining and clearing them. This can keep ports exposed on the tailnet beyond what your current labels define; see [Cleanup Behavior](#cleanup-behavior). | | `LOG_LEVEL` | `info` | Logging level for all output, including the DockTail Cloud module: `debug`, `info`, `warn`, or `error`. Any other value means `info`. | | `RECONCILE_INTERVAL` | `60s` | State reconciliation interval. | +| `DOCKTAIL_DISCOVERY` | `containers` | Where DockTail looks for things to expose. `containers` exposes the containers running on **this Docker node**. `swarm` exposes every labelled service in the **whole cluster**, and requires a manager endpoint. See [Docker Swarm](#docker-swarm). | +| `DOCKTAIL_SWARM_NETWORK` | - | With `DOCKTAIL_DISCOVERY=swarm`: pin which Docker network a service's virtual IP is taken from when the service is attached to several. Defaults to a network DockTail is also attached to, then to the lowest IP. | | `DOCKER_HOST` | `unix:///var/run/docker.sock` | Docker daemon socket. Rootless Docker typically uses `unix:///run/user//docker.sock`. | | `TAILSCALE_SOCKET` | `/var/run/tailscale/tailscaled.sock` | The `tailscaled` socket DockTail checks: the startup missing-socket hint, the [socket-loss check](#tailscale-socket-loss), and the image's health check. DockTail does not pass it to the bundled `tailscale` CLI, which does the serve and Funnel work at its own default of `/var/run/tailscale/tailscaled.sock`, so mount the daemon's socket directory at `/var/run/tailscale` either way. | | `EXIT_ON_SOCKET_LOSS` | `true` | When `true`, DockTail exits if the Tailscale socket stays unreachable past the grace period, so the container's restart policy can re-establish the mount. See [Tailscale Socket Loss](#tailscale-socket-loss). | @@ -69,6 +71,50 @@ endpoint. `ws://` is allowed for loopback endpoints; non-loopback plaintext requires `DOCKTAIL_CLOUD_ALLOW_INSECURE=true` and must never be used in production. Non-loopback production endpoints must use `wss://`. +### Docker Swarm + +By default DockTail exposes the containers running on **its own Docker node**. In a Swarm cluster that is a narrow slice of the cluster: the `/containers` endpoint is node-local, on manager nodes as much as workers, so one agent sees only the containers it happens to be colocated with. + +`DOCKTAIL_DISCOVERY=swarm` switches to the cluster view. DockTail then lists every labelled **service** in the cluster through the Swarm manager API and advertises each one at its service's virtual IP. + +```yaml +services: + docktail: + image: ghcr.io/marvinvr/docktail:latest + environment: + - DOCKTAIL_DISCOVERY=swarm + - TAILSCALE_API_KEY=${TAILSCALE_API_KEY} + volumes: + - /var/run/docker.sock:/var/run/docker.sock:ro + - /var/run/tailscale:/var/run/tailscale +``` + +What you need to know: + +- **A manager endpoint is required.** `/services` is manager-only. On a worker-only agent the reconcile fails and nothing is advertised. +- **One agent is enough.** Services on every node are covered by the single agent, because the service VIP is cluster-wide. +- **The overlay carries the traffic.** A service VIP is only dialable from a node attached to the same overlay. Attach DockTail — and therefore the `tailscaled` it configures — to the same network the exposed services use, and prefer that network when picking a VIP. When a service is attached to several networks, DockTail picks one DockTail is also attached to; set `DOCKTAIL_SWARM_NETWORK` to pin it explicitly. +- **A VIP only exists in `vip` mode.** A `dnsrr` service, or one with no network, has none and is skipped with a warning. +- **Changes on other nodes arrive on the next reconcile tick.** Docker event watching is still local, so a service created or relabelled on another node is picked up within `RECONCILE_INTERVAL`, not immediately. + +Labels come from the service's task template, which is where both `labels:` and `deploy.labels:` land in a compose file — so a Swarm service is labelled exactly like the equivalent standalone container, and every label in [Labels](04-labels.md) applies unchanged. + +```yaml +services: + plex: + image: lscr.io/linuxserver/plex:latest + networks: + - tailscale + labels: + - "docktail.service.enable=true" + - "docktail.service.name=plex" + - "docktail.service.port=32400" + +networks: + tailscale: + external: true +``` + ### Supported Protocols Tailscale-facing `docktail.service.service-protocol` values: diff --git a/main.go b/main.go index 909cf32..480c2bb 100644 --- a/main.go +++ b/main.go @@ -44,6 +44,21 @@ func main() { reconcileInterval := getEnvDuration("RECONCILE_INTERVAL", 60*time.Second) tailscaleSocket := getEnv("TAILSCALE_SOCKET", "/var/run/tailscale/tailscaled.sock") + // Discovery mode: "containers" (default) reads the containers running on + // this Docker node; "swarm" reads every labelled service in the cluster + // through the Swarm manager API, so one agent covers the whole cluster. + discoveryMode := getEnv("DOCKTAIL_DISCOVERY", "containers") + switch discoveryMode { + case "containers", "swarm": + default: + log.Warn(). + Str("key", "DOCKTAIL_DISCOVERY"). + Str("value", discoveryMode). + Str("default", "containers"). + Msg("Unknown discovery mode, using containers") + discoveryMode = "containers" + } + // Control Plane Configuration tailscaleAPIKey := getSecretEnv("TAILSCALE_API_KEY", "") tailscaleOAuthClientID := getSecretEnv("TAILSCALE_OAUTH_CLIENT_ID", "") @@ -106,6 +121,7 @@ func main() { log.Info(). Dur("reconcile_interval", reconcileInterval). + Str("discovery", discoveryMode). Str("tailscale_socket", tailscaleSocket). Str("api_sync_method", apiSyncMethod). Str("tailnet", tailscaleTailnet). @@ -167,6 +183,9 @@ func main() { // Create reconciler rec := reconciler.NewReconciler(dockerClient, tailscaleClient, reconcileInterval) + if discoveryMode == "swarm" { + rec = reconciler.NewSwarmReconciler(dockerClient, tailscaleClient, reconcileInterval) + } // Setup signal handling ctx, cancel := context.WithCancel(context.Background()) diff --git a/reconciler/reconciler.go b/reconciler/reconciler.go index 286e391..d38c1b1 100644 --- a/reconciler/reconciler.go +++ b/reconciler/reconciler.go @@ -33,15 +33,49 @@ type Observer interface { OnEvent(ctx context.Context, event events.Message) } +// Discovery supplies the set of Tailscale services to reconcile. The default +// implementation reads the containers running on this Docker node; the Swarm +// implementation (see ../docker/swarm.go) reads every labelled service in the +// cluster instead, so one agent can cover a multi-node cluster. +type Discovery interface { + // Name identifies the discovery mode in logs. + Name() string + // List returns every DockTail-managed service/funnel to reconcile. + List(ctx context.Context) ([]*apptypes.ContainerService, error) +} + // Reconciler manages the reconciliation loop type Reconciler struct { dockerClient *docker.Client + discovery Discovery tailscaleClient *tailscale.Client interval time.Duration observer Observer // optional; nil unless the cloud module is enabled onResult func(err error) // optional; told the outcome of every reconcile } +// containerDiscovery is the node-local default: labelled containers running on +// this Docker node, watched through the local docker socket. +type containerDiscovery struct{ client *docker.Client } + +func (containerDiscovery) Name() string { return "containers" } + +func (d containerDiscovery) List(ctx context.Context) ([]*apptypes.ContainerService, error) { + return d.client.GetEnabledContainers(ctx) +} + +// swarmDiscovery reads labelled services across the whole cluster through the +// Swarm manager API. Service event watching still comes from the local socket, +// so a service change on another node is picked up on the next reconcile tick +// rather than immediately. +type swarmDiscovery struct{ client *docker.Client } + +func (swarmDiscovery) Name() string { return "swarm" } + +func (d swarmDiscovery) List(ctx context.Context) ([]*apptypes.ContainerService, error) { + return d.client.GetEnabledSwarmServices(ctx) +} + // SetResultHook installs a function told the outcome of every reconcile cycle // (nil on success) — the local health status. Safe to call once before Run. func (r *Reconciler) SetResultHook(fn func(err error)) { @@ -67,10 +101,34 @@ func triggersReconcile(a events.Action) bool { } } -// NewReconciler creates a new reconciler +// NewReconciler creates a new reconciler using node-local container discovery. func NewReconciler(dockerClient *docker.Client, tailscaleClient *tailscale.Client, interval time.Duration) *Reconciler { + return NewReconcilerWithDiscovery( + dockerClient, + tailscaleClient, + interval, + containerDiscovery{client: dockerClient}, + ) +} + +// NewSwarmReconciler creates a reconciler that discovers labelled services +// across the whole Swarm cluster instead of only this node's containers. The +// docker client must point at a manager endpoint. +func NewSwarmReconciler(dockerClient *docker.Client, tailscaleClient *tailscale.Client, interval time.Duration) *Reconciler { + return NewReconcilerWithDiscovery( + dockerClient, + tailscaleClient, + interval, + swarmDiscovery{client: dockerClient}, + ) +} + +// NewReconcilerWithDiscovery creates a reconciler with an explicit discovery +// source. +func NewReconcilerWithDiscovery(dockerClient *docker.Client, tailscaleClient *tailscale.Client, interval time.Duration, discovery Discovery) *Reconciler { return &Reconciler{ dockerClient: dockerClient, + discovery: discovery, tailscaleClient: tailscaleClient, interval: interval, } @@ -143,14 +201,15 @@ func (r *Reconciler) Reconcile(ctx context.Context) error { func (r *Reconciler) reconcile(ctx context.Context) error { log.Info().Msg("Starting reconciliation") - // Get all enabled containers from Docker - containers, err := r.dockerClient.GetEnabledContainers(ctx) + // Get all enabled containers (or, in swarm mode, services) from Docker + containers, err := r.discovery.List(ctx) if err != nil { return fmt.Errorf("%w: %w", ErrListContainers, err) } log.Info(). Int("count", len(containers)). + Str("source", r.discovery.Name()). Msg("Found enabled containers") for _, container := range containers {