Merge pull request #2224 from rcourtman/maintainer/20260924T055323Z

Keep provider MSP client routes online after upgrades and isolated across providers
This commit is contained in:
pulse-triage[bot] 2026-09-24 06:28:10 +00:00 • committed by GitHub
commit 8d98add802
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
5 changed files with 552 additions and 30 deletions

View file

@ -3443,10 +3443,16 @@ client's isolated tenant network. Recreating either one, which every
`upgrade.sh` run does to the control plane, drops those attachments, cutting
the client's Traefik route and, with the old stop behaviour, taking every
client offline three minutes later. A client without an isolated network is
skipped quietly. Regression coverage:
skipped quietly during recovery. Provisioning may create a missing isolated
network, but it must not adopt a same-named pre-existing network unless the
network has the exact tenant ID and provider-MSP tenant-runtime labels; an
empty tenant ID is invalid. This prevents a foreign or unlabelled network from
becoming a client's runtime route before support containers are attached.
Regression coverage:
`TestHealthMonitorReattachesSupportContainersAndRestartsInsteadOfStopping` in
`internal/cloudcp/health_monitor_test.go` and
`TestEnsureSupportContainersOnTenantNetworkSkipsMissingNetwork` in
`TestEnsureSupportContainersOnTenantNetworkSkipsMissingNetwork` and
`TestEnsureTenantNetworkRejectsUnownedExistingNetwork` in
`internal/cloudcp/docker/manager_test.go`.
### Provider-hosted MSP platforms buy and renew their own licence
@ -3500,6 +3506,27 @@ coverage:
`TestProviderMSPWorkspaceLadderRisesFromTheEvaluation` in
`pkg/licensing/features_test.go`.
### Provider MSP health monitor recovers clients across upgrades
The health monitor rejoins Traefik and the control plane to each active
client's labelled isolated network on startup and before every health check.
Recreating either support container drops its tenant-network attachments;
without reconnection, routes fail and a healthy client can appear unreachable.
The monitor now restarts an unhealthy client with Docker's restart operation,
rather than stopping it: `unless-stopped` does not revive a container stopped
through the API. Missing isolated networks are skipped; a network whose labels
do not identify that tenant is refused rather than connected. Support containers
must also have the role label and belong to the configured provider ingress
network, avoiding a same-host provider stack with the same role label. Regression
coverage is in `internal/cloudcp/health_monitor_test.go` and
`internal/cloudcp/docker/manager_test.go`.
The same boundary applies when removing a client: cleanup only selects that
managed client's derived, correctly labelled tenant network, and only
disconnects support containers on this provider ingress that are still
attached to it. A role label by itself must not select another provider's
support container for a forced disconnect.
### Provider MSP status reads client health that any caller can observe
`provider-msp status`, which `upgrade.sh` gates on before every provider

View file

@ -2451,16 +2451,26 @@ has not repeated that live-host proof.
### Provider MSP clients survive a support-container recreate
Recreating the control plane or Traefik (every `upgrade.sh` run recreates the
control plane) used to detach them from each client's isolated network. The
health monitor then saw every client as failing and stopped it, and nothing
restarted it, so an upgrade on a v6.4.1 install with clients took all of them
offline. The monitor now reattaches the support containers on startup and on
every pass, and restarts rather than stops an unhealthy client. Verified on
2026-09-23 on a v6.4.1 provider bundle with two clients: a control-plane
recreate rejoined both client networks within a second, a Traefik-only
recreate had its client routes back on the next 60-second pass, and neither
client was stopped.
Recreating the control plane or Traefik detaches it from each client's
isolated network. The health monitor now reattaches both support containers
on startup and each pass, and restarts rather than stops an unhealthy client.
The original PR reports a two-client v6.4.1 lab reproduction and recovery on
the same installation on 2026-09-23; that is not installed acceptance of this
integrated source. Reconnection selects support containers from the configured
provider ingress network. Both reconnection and new workspace provisioning
require exact tenant ownership and tenant-runtime labels on an existing
isolated network before attaching containers; provisioning refuses an
unlabelled same-named network instead of adopting it. A legacy client without
an isolated network remains a reconnection no-op.
Client removal follows the same installation boundary: the Docker manager
identifies the managed client's network by its derived name and exact tenant
and runtime labels, not by a generic tenant-network label alone. It must force
disconnect only this provider's role-labelled support containers that are
still attached to that network, selecting them through the configured provider
ingress network. A recreated or other provider's support container is not a
cleanup target. `internal/cloudcp/docker/manager_test.go` covers the network
selection and disconnect filters; source-only coverage does not replace a
two-provider installed upgrade, recreation and cleanup check.
### Provider MSP operations accept a renewed licence

View file

@ -365,12 +365,17 @@ func (m *Manager) ensureTenantNetwork(ctx context.Context, tenantID string) (str
}
return networkName, nil
}
if strings.TrimSpace(tenantID) == "" {
return "", fmt.Errorf("tenant ID is required for an isolated network")
}
networkName := m.tenantNetworkName(tenantID)
inspect, err := m.cli.NetworkInspect(ctx, networkName, client.NetworkInspectOptions{})
if err == nil {
if got := strings.TrimSpace(inspect.Network.Labels["pulse.tenant.id"]); got != "" && got != tenantID {
return "", fmt.Errorf("tenant network %q belongs to tenant %q, not %q", networkName, got, tenantID)
// A matching name is not evidence that this is our isolated network.
// In particular, never adopt a pre-existing unlabelled Docker network.
if !isOwnedTenantNetwork(inspect.Network.Labels, tenantID) {
return "", fmt.Errorf("network %q is not the isolated network for tenant %q", networkName, tenantID)
}
return networkName, nil
}
@ -392,6 +397,10 @@ func (m *Manager) ensureTenantNetwork(ctx context.Context, tenantID string) (str
return networkName, nil
}
func isOwnedTenantNetwork(labels map[string]string, tenantID string) bool {
return tenantID != "" && labels["pulse.tenant.id"] == tenantID && labels[tenantRuntimeNetworkLabel] == tenantRuntimeNetworkLabelValue
}
func (m *Manager) connectSupportContainersToTenantNetwork(ctx context.Context, networkName string) error {
if m == nil || !m.cfg.IsolateTenantNetworks {
return nil
@ -409,8 +418,15 @@ func (m *Manager) connectSupportContainersToTenantNetwork(ctx context.Context, n
}
func (m *Manager) connectSupportContainersByLabel(ctx context.Context, networkName, label string) error {
providerNetwork := strings.TrimSpace(m.cfg.Network)
if providerNetwork == "" {
return fmt.Errorf("provider ingress network is required to select support containers")
}
filters := client.Filters{}
filters = filters.Add("label", label)
// Role labels are shared by every provider MSP installation on a host.
// Never attach another installation's support container to this tenant.
filters = filters.Add("network", providerNetwork)
result, err := m.cli.ContainerList(ctx, client.ContainerListOptions{
Filters: filters,
})
@ -437,10 +453,8 @@ func (m *Manager) ensureContainerConnectedToNetwork(ctx context.Context, network
if err != nil {
return fmt.Errorf("inspect tenant network %q: %w", networkName, err)
}
for id := range inspect.Network.Containers {
if id == containerID || strings.HasPrefix(id, containerID) || strings.HasPrefix(containerID, id) {
return nil
}
if networkHasContainer(inspect.Network.Containers, containerID) {
return nil
}
_, err = m.cli.NetworkConnect(ctx, networkName, client.NetworkConnectOptions{
Container: containerID,
@ -449,6 +463,18 @@ func (m *Manager) ensureContainerConnectedToNetwork(ctx context.Context, network
return err
}
func networkHasContainer(containers map[string]network.EndpointResource, containerID string) bool {
if containerID == "" {
return false
}
for id := range containers {
if id != "" && (id == containerID || strings.HasPrefix(id, containerID) || strings.HasPrefix(containerID, id)) {
return true
}
}
return false
}
// IsNotFound reports whether Docker treated an identifier as missing.
func IsNotFound(err error) bool {
return errdefs.IsNotFound(err)
@ -887,12 +913,18 @@ func (m *Manager) EnsureSupportContainersOnTenantNetwork(ctx context.Context, te
if networkName == "" {
return nil
}
if _, err := m.cli.NetworkInspect(ctx, networkName, client.NetworkInspectOptions{}); err != nil {
inspected, err := m.cli.NetworkInspect(ctx, networkName, client.NetworkInspectOptions{})
if err != nil {
if errdefs.IsNotFound(err) {
return nil
}
return fmt.Errorf("inspect tenant network %q: %w", networkName, err)
}
// A name alone does not establish ownership. Never attach provider support
// containers to an unrelated network that happens to have the derived name.
if !isOwnedTenantNetwork(inspected.Network.Labels, tenantID) {
return fmt.Errorf("network %q is not the isolated network for tenant %q", networkName, tenantID)
}
return m.connectSupportContainersToTenantNetwork(ctx, networkName)
}
@ -927,12 +959,17 @@ func (m *Manager) tenantNetworkNamesForContainer(ctx context.Context, containerI
return nil
}
inspect := inspectResult.Container
if inspect.NetworkSettings == nil {
if inspect.Config == nil || inspect.NetworkSettings == nil {
return nil
}
tenantID := strings.TrimSpace(inspect.Config.Labels["pulse.tenant.id"])
if tenantID == "" || inspect.Config.Labels["pulse.managed"] != "true" {
return nil
}
ownedName := m.tenantNetworkName(tenantID)
var out []string
for networkName := range inspect.NetworkSettings.Networks {
if m.isTenantNetwork(ctx, networkName) {
if networkName == ownedName && m.isTenantNetwork(ctx, networkName, tenantID) {
out = append(out, networkName)
}
}
@ -940,16 +977,16 @@ func (m *Manager) tenantNetworkNamesForContainer(ctx context.Context, containerI
return out
}
func (m *Manager) isTenantNetwork(ctx context.Context, networkName string) bool {
func (m *Manager) isTenantNetwork(ctx context.Context, networkName, tenantID string) bool {
networkName = strings.TrimSpace(networkName)
if m == nil || networkName == "" {
if m == nil || networkName == "" || tenantID == "" {
return false
}
inspect, err := m.cli.NetworkInspect(ctx, networkName, client.NetworkInspectOptions{})
if err != nil {
return false
}
return strings.TrimSpace(inspect.Network.Labels[tenantRuntimeNetworkLabel]) == tenantRuntimeNetworkLabelValue
return isOwnedTenantNetwork(inspect.Network.Labels, tenantID)
}
func (m *Manager) removeTenantNetwork(ctx context.Context, networkName string) error {
@ -969,6 +1006,14 @@ func (m *Manager) disconnectSupportContainersFromTenantNetwork(ctx context.Conte
if m == nil || !m.cfg.IsolateTenantNetworks {
return nil
}
providerNetwork := strings.TrimSpace(m.cfg.Network)
if providerNetwork == "" {
return fmt.Errorf("provider ingress network is required to select support containers")
}
inspected, err := m.cli.NetworkInspect(ctx, networkName, client.NetworkInspectOptions{})
if err != nil {
return fmt.Errorf("inspect tenant network %q for disconnect: %w", networkName, err)
}
for _, label := range m.cfg.SupportContainerLabels {
label = strings.TrimSpace(label)
if label == "" {
@ -976,6 +1021,8 @@ func (m *Manager) disconnectSupportContainersFromTenantNetwork(ctx context.Conte
}
filters := client.Filters{}
filters = filters.Add("label", label)
// The role label is shared by every provider installation on this host.
filters = filters.Add("network", providerNetwork)
result, err := m.cli.ContainerList(ctx, client.ContainerListOptions{
All: true,
Filters: filters,
@ -984,7 +1031,9 @@ func (m *Manager) disconnectSupportContainersFromTenantNetwork(ctx context.Conte
return fmt.Errorf("list provider support containers for disconnect label %q: %w", label, err)
}
for _, item := range result.Items {
if item.ID == "" {
// A replacement support container has not joined this tenant network.
// Do not try to disconnect it (or another provider's container).
if !networkHasContainer(inspected.Network.Containers, item.ID) {
continue
}
_, err := m.cli.NetworkDisconnect(ctx, networkName, client.NetworkDisconnectOptions{

View file

@ -2,6 +2,7 @@ package docker
import (
"context"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
@ -641,3 +642,212 @@ func TestEnsureSupportContainersOnTenantNetworkSkipsMissingNetwork(t *testing.T)
t.Fatalf("reconcile changed the host for a missing network: %v", mutations)
}
}
func TestEnsureSupportContainersOnTenantNetworkRejectsWrongOwner(t *testing.T) {
var mutations []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Api-Version", "1.47")
if strings.HasSuffix(r.URL.Path, "/_ping") {
_, _ = w.Write([]byte("OK"))
return
}
if r.Method != http.MethodGet {
mutations = append(mutations, r.Method+" "+r.URL.Path)
}
w.Header().Set("Content-Type", "application/json")
_, _ = w.Write([]byte(`{"Name":"pulse-provider-msp-tenant-t-acme","Id":"net-1","Labels":{"pulse.tenant.id":"t-other","pulse.provider-msp.network":"tenant"}}`))
}))
t.Cleanup(srv.Close)
t.Setenv("DOCKER_HOST", "tcp://"+strings.TrimPrefix(srv.URL, "http://"))
t.Setenv("DOCKER_TLS_VERIFY", "")
t.Setenv("DOCKER_CERT_PATH", "")
mgr, err := NewManager(ManagerConfig{Image: "pulse:test", Network: "pulse-provider-msp", IsolateTenantNetworks: true})
if err != nil {
t.Fatalf("NewManager: %v", err)
}
t.Cleanup(func() { _ = mgr.Close() })
if err := mgr.EnsureSupportContainersOnTenantNetwork(context.Background(), "t-acme"); err == nil {
t.Fatal("expected wrong-owner network to be rejected")
}
if len(mutations) != 0 {
t.Fatalf("reconcile changed wrong-owner network: %v", mutations)
}
}
func TestEnsureTenantNetworkRejectsUnownedExistingNetwork(t *testing.T) {
var labels map[string]string
var mutations []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Api-Version", "1.47")
if strings.HasSuffix(r.URL.Path, "/_ping") {
_, _ = w.Write([]byte("OK"))
return
}
if r.Method != http.MethodGet {
mutations = append(mutations, r.Method+" "+r.URL.Path)
}
w.Header().Set("Content-Type", "application/json")
_ = json.NewEncoder(w).Encode(map[string]any{
"Name": "pulse-provider-msp-tenant-t-acme",
"Id": "net-1",
"Labels": labels,
})
}))
t.Cleanup(srv.Close)
t.Setenv("DOCKER_HOST", "tcp://"+strings.TrimPrefix(srv.URL, "http://"))
t.Setenv("DOCKER_TLS_VERIFY", "")
t.Setenv("DOCKER_CERT_PATH", "")
mgr, err := NewManager(ManagerConfig{Image: "pulse:test", Network: "pulse-provider-msp", IsolateTenantNetworks: true})
if err != nil {
t.Fatalf("NewManager: %v", err)
}
t.Cleanup(func() { _ = mgr.Close() })
for _, tc := range []struct {
name string
labels map[string]string
wantOK bool
}{
{name: "unlabelled"},
{name: "wrong tenant", labels: map[string]string{"pulse.tenant.id": "t-other", tenantRuntimeNetworkLabel: tenantRuntimeNetworkLabelValue}},
{name: "missing runtime label", labels: map[string]string{"pulse.tenant.id": "t-acme"}},
{name: "wrong runtime label", labels: map[string]string{"pulse.tenant.id": "t-acme", tenantRuntimeNetworkLabel: "ingress"}},
{name: "owned tenant network", labels: map[string]string{"pulse.tenant.id": "t-acme", tenantRuntimeNetworkLabel: tenantRuntimeNetworkLabelValue}, wantOK: true},
} {
t.Run(tc.name, func(t *testing.T) {
labels = tc.labels
mutations = nil
got, err := mgr.ensureTenantNetwork(context.Background(), "t-acme")
if tc.wantOK {
if err != nil || got != "pulse-provider-msp-tenant-t-acme" {
t.Fatalf("ensureTenantNetwork = (%q, %v), want owned network", got, err)
}
} else if err == nil || got != "" {
t.Fatalf("ensureTenantNetwork = (%q, %v), want fail closed", got, err)
}
if len(mutations) != 0 {
t.Fatalf("ensureTenantNetwork mutated existing network: %v", mutations)
}
})
}
if got, err := mgr.ensureTenantNetwork(context.Background(), ""); err == nil || got != "" {
t.Fatalf("ensureTenantNetwork(empty) = (%q, %v), want rejection", got, err)
}
}
func TestTenantNetworkCleanupSelectsOnlyOwnedNetwork(t *testing.T) {
ownedLabels := map[string]string{"pulse.tenant.id": "t-acme", tenantRuntimeNetworkLabel: tenantRuntimeNetworkLabelValue}
var containerLabels, targetLabels map[string]string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Api-Version", "1.47")
w.Header().Set("Content-Type", "application/json")
switch {
case strings.HasSuffix(r.URL.Path, "/containers/client-a/json"):
_ = json.NewEncoder(w).Encode(map[string]any{
"Id": "client-a", "Config": map[string]any{"Labels": containerLabels},
"NetworkSettings": map[string]any{"Networks": map[string]any{
"provider-a-tenant-t-acme": map[string]any{},
"provider-b-tenant-t-acme": map[string]any{},
}},
})
case strings.HasSuffix(r.URL.Path, "/networks/provider-a-tenant-t-acme"):
_ = json.NewEncoder(w).Encode(map[string]any{"Name": "provider-a-tenant-t-acme", "Labels": targetLabels})
case strings.HasSuffix(r.URL.Path, "/networks/provider-b-tenant-t-acme"):
_ = json.NewEncoder(w).Encode(map[string]any{"Name": "provider-b-tenant-t-acme", "Labels": ownedLabels})
default:
w.WriteHeader(http.StatusNotFound)
}
}))
t.Cleanup(srv.Close)
t.Setenv("DOCKER_HOST", "tcp://"+strings.TrimPrefix(srv.URL, "http://"))
t.Setenv("DOCKER_TLS_VERIFY", "")
t.Setenv("DOCKER_CERT_PATH", "")
mgr, err := NewManager(ManagerConfig{Network: "provider-a", IsolateTenantNetworks: true})
if err != nil {
t.Fatalf("NewManager: %v", err)
}
t.Cleanup(func() { _ = mgr.Close() })
for _, tc := range []struct {
name string
containerLabels map[string]string
targetLabels map[string]string
want []string
}{
{name: "owned", containerLabels: map[string]string{"pulse.tenant.id": "t-acme", "pulse.managed": "true"}, targetLabels: ownedLabels, want: []string{"provider-a-tenant-t-acme"}},
{name: "unlabelled network", containerLabels: map[string]string{"pulse.tenant.id": "t-acme", "pulse.managed": "true"}},
{name: "wrong network tenant", containerLabels: map[string]string{"pulse.tenant.id": "t-acme", "pulse.managed": "true"}, targetLabels: map[string]string{"pulse.tenant.id": "t-other", tenantRuntimeNetworkLabel: tenantRuntimeNetworkLabelValue}},
{name: "unmanaged container", containerLabels: map[string]string{"pulse.tenant.id": "t-acme"}, targetLabels: ownedLabels},
{name: "no container tenant", containerLabels: map[string]string{"pulse.managed": "true"}, targetLabels: ownedLabels},
} {
t.Run(tc.name, func(t *testing.T) {
containerLabels, targetLabels = tc.containerLabels, tc.targetLabels
got := mgr.tenantNetworkNamesForContainer(context.Background(), "client-a")
if len(got) != len(tc.want) || (len(got) > 0 && got[0] != tc.want[0]) {
t.Fatalf("tenantNetworkNamesForContainer = %v, want %v", got, tc.want)
}
})
}
}
func TestTenantNetworkCleanupDisconnectsOnlyAttachedLocalSupport(t *testing.T) {
const target = "provider-a-tenant-t-acme"
const localAttached = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
const localReplacement = "cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc"
const otherAttached = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
var listFilters map[string]map[string]bool
var disconnected []string
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Api-Version", "1.47")
w.Header().Set("Content-Type", "application/json")
switch {
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/networks/"+target):
_ = json.NewEncoder(w).Encode(map[string]any{
"Name": target, "Labels": map[string]string{"pulse.tenant.id": "t-acme", tenantRuntimeNetworkLabel: tenantRuntimeNetworkLabelValue},
"Containers": map[string]any{localAttached: map[string]any{}, otherAttached: map[string]any{}},
})
case r.Method == http.MethodGet && strings.HasSuffix(r.URL.Path, "/containers/json"):
if err := json.Unmarshal([]byte(r.URL.Query().Get("filters")), &listFilters); err != nil {
w.WriteHeader(http.StatusBadRequest)
return
}
items := []map[string]string{{"Id": localAttached}, {"Id": localReplacement}}
if !listFilters["network"]["provider-a"] {
items = append(items, map[string]string{"Id": otherAttached})
}
_ = json.NewEncoder(w).Encode(items)
case r.Method == http.MethodPost && strings.HasSuffix(r.URL.Path, "/networks/"+target+"/disconnect"):
var body struct{ Container string }
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
w.WriteHeader(http.StatusBadRequest)
return
}
disconnected = append(disconnected, body.Container)
_, _ = w.Write([]byte(`{}`))
default:
w.WriteHeader(http.StatusNotFound)
}
}))
t.Cleanup(srv.Close)
t.Setenv("DOCKER_HOST", "tcp://"+strings.TrimPrefix(srv.URL, "http://"))
t.Setenv("DOCKER_TLS_VERIFY", "")
t.Setenv("DOCKER_CERT_PATH", "")
mgr, err := NewManager(ManagerConfig{Network: "provider-a", IsolateTenantNetworks: true, SupportContainerLabels: []string{providerSupportTraefikLabel}})
if err != nil {
t.Fatalf("NewManager: %v", err)
}
t.Cleanup(func() { _ = mgr.Close() })
if err := mgr.disconnectSupportContainersFromTenantNetwork(context.Background(), target); err != nil {
t.Fatalf("disconnectSupportContainersFromTenantNetwork: %v", err)
}
if !listFilters["network"]["provider-a"] || !listFilters["label"][providerSupportTraefikLabel] {
t.Fatalf("support selection did not identify the local provider: %v", listFilters)
}
if len(disconnected) != 1 || disconnected[0] != localAttached {
t.Fatalf("disconnected = %v, want only attached local support %s", disconnected, localAttached)
}
}

View file

@ -9,6 +9,7 @@ import (
"strings"
"sync"
"testing"
"time"
"github.com/rcourtman/pulse-go-rewrite/internal/cloudcp/docker"
"github.com/rcourtman/pulse-go-rewrite/internal/cloudcp/registry"
@ -24,6 +25,7 @@ type fakeProviderDaemon struct {
support map[string]string
health string
calls []string
healthChecked chan struct{}
}
var dockerAPIVersionPrefix = regexp.MustCompile(`^/v[0-9.]+`)
@ -43,7 +45,10 @@ func (d *fakeProviderDaemon) ServeHTTP(w http.ResponseWriter, r *http.Request) {
for id := range d.attached {
containers[id] = map[string]any{"Name": id}
}
_ = json.NewEncoder(w).Encode(map[string]any{"Name": d.tenantNetwork, "Id": "net-1", "Containers": containers})
_ = json.NewEncoder(w).Encode(map[string]any{
"Name": d.tenantNetwork, "Id": "net-1", "Containers": containers,
"Labels": map[string]string{"pulse.tenant.id": "t-acme", "pulse.provider-msp.network": "tenant"},
})
case r.Method == http.MethodPost && path == "/networks/"+d.tenantNetwork+"/connect":
var body struct{ Container string }
_ = json.NewDecoder(r.Body).Decode(&body)
@ -51,14 +56,26 @@ func (d *fakeProviderDaemon) ServeHTTP(w http.ResponseWriter, r *http.Request) {
d.calls = append(d.calls, "connect "+body.Container)
w.WriteHeader(http.StatusOK)
case r.Method == http.MethodGet && path == "/containers/json":
filters := r.URL.Query().Get("filters")
if !strings.Contains(filters, `"network"`) || !strings.Contains(filters, "pulse-provider-msp") {
w.WriteHeader(http.StatusBadRequest)
_, _ = w.Write([]byte(`{"message":"support containers must be selected from the provider ingress network"}`))
return
}
var items []map[string]any
for label, id := range d.support {
if strings.Contains(r.URL.Query().Get("filters"), label) {
if strings.Contains(filters, label) {
items = append(items, map[string]any{"Id": id, "State": "running"})
}
}
_ = json.NewEncoder(w).Encode(items)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/containers/") && strings.HasSuffix(path, "/json"):
if d.healthChecked != nil {
select {
case d.healthChecked <- struct{}{}:
default:
}
}
_ = json.NewEncoder(w).Encode(map[string]any{
"Id": strings.TrimSuffix(strings.TrimPrefix(path, "/containers/"), "/json"),
"State": map[string]any{"Status": "running", "Running": true, "Health": map[string]any{"Status": d.health}},
@ -94,7 +111,8 @@ func TestHealthMonitorReattachesSupportContainersAndRestartsInsteadOfStopping(t
"pulse.provider-msp.role=traefik": "traefik-recreated",
"pulse.provider-msp.role=control-plane": "control-plane-recreated",
},
health: "healthy",
health: "healthy",
healthChecked: make(chan struct{}, 1),
}
srv := httptest.NewServer(daemon)
t.Cleanup(srv.Close)
@ -123,8 +141,37 @@ func TestHealthMonitorReattachesSupportContainersAndRestartsInsteadOfStopping(t
t.Fatalf("Create tenant: %v", err)
}
monitor := NewMonitor(reg, mgr, MonitorConfig{RestartOnFail: true, FailThreshold: 3})
monitor.checkAll(context.Background())
monitor := NewMonitor(reg, mgr, MonitorConfig{Interval: time.Hour, RestartOnFail: true, FailThreshold: 3})
ctx, cancel := context.WithCancel(context.Background())
done := make(chan struct{})
go func() {
monitor.Run(ctx)
close(done)
}()
select {
case <-daemon.healthChecked:
case <-time.After(5 * time.Second):
cancel()
t.Fatal("monitor did not check the tenant on startup")
}
deadline := time.Now().Add(5 * time.Second)
for {
observed, err := reg.Get("t-acme")
if err == nil && observed != nil && observed.HealthCheckOK {
break
}
if time.Now().After(deadline) {
cancel()
t.Fatalf("monitor did not persist healthy result on startup: tenant=%+v err=%v", observed, err)
}
time.Sleep(5 * time.Millisecond)
}
cancel()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("monitor did not stop after cancellation")
}
got := strings.Join(daemon.recorded(), ", ")
for _, want := range []string{"connect traefik-recreated", "connect control-plane-recreated"} {
@ -157,3 +204,182 @@ func TestHealthMonitorReattachesSupportContainersAndRestartsInsteadOfStopping(t
t.Fatalf("unhealthy client was stopped, which unless-stopped never undoes; daemon saw %q", got)
}
}
// Two provider installations can share a Docker daemon and the same support
// role labels. A recreated container must only join the tenant network of the
// installation whose ingress network already contains it.
type multiProviderDaemon struct {
mu sync.Mutex
networks map[string]*multiProviderNetwork
support map[string]multiProviderSupport
health map[string]string
calls []string
}
type multiProviderNetwork struct {
tenantID string
attached map[string]bool
}
type multiProviderSupport struct {
role string
ingress string
}
func (d *multiProviderDaemon) ServeHTTP(w http.ResponseWriter, r *http.Request) {
d.mu.Lock()
defer d.mu.Unlock()
path := dockerAPIVersionPrefix.ReplaceAllString(r.URL.Path, "")
w.Header().Set("Api-Version", "1.47")
w.Header().Set("Content-Type", "application/json")
switch {
case path == "/_ping":
_, _ = w.Write([]byte("OK"))
case r.Method == http.MethodGet && strings.HasPrefix(path, "/networks/"):
name := strings.TrimPrefix(path, "/networks/")
net, ok := d.networks[name]
if !ok {
w.WriteHeader(http.StatusNotFound)
_, _ = w.Write([]byte(`{"message":"network not found"}`))
return
}
containers := make(map[string]any, len(net.attached))
for id := range net.attached {
containers[id] = map[string]any{"Name": id}
}
labels := map[string]string{}
if net.tenantID != "" {
labels = map[string]string{"pulse.tenant.id": net.tenantID, "pulse.provider-msp.network": "tenant"}
}
_ = json.NewEncoder(w).Encode(map[string]any{
"Name": name, "Id": name, "Containers": containers, "Labels": labels,
})
case r.Method == http.MethodPost && strings.HasPrefix(path, "/networks/") && strings.HasSuffix(path, "/connect"):
name := strings.TrimSuffix(strings.TrimPrefix(path, "/networks/"), "/connect")
var body struct{ Container string }
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
w.WriteHeader(http.StatusBadRequest)
return
}
if _, ok := d.networks[name]; !ok {
w.WriteHeader(http.StatusNotFound)
return
}
d.networks[name].attached[body.Container] = true
d.calls = append(d.calls, "connect "+body.Container+" "+name)
w.WriteHeader(http.StatusOK)
case r.Method == http.MethodGet && path == "/containers/json":
var filters map[string]map[string]bool
if err := json.Unmarshal([]byte(r.URL.Query().Get("filters")), &filters); err != nil {
w.WriteHeader(http.StatusBadRequest)
return
}
var items []map[string]any
for id, support := range d.support {
if matchesDockerFilter(filters["label"], support.role) && matchesDockerFilter(filters["network"], support.ingress) {
items = append(items, map[string]any{"Id": id, "State": "running"})
}
}
_ = json.NewEncoder(w).Encode(items)
case r.Method == http.MethodGet && strings.HasPrefix(path, "/containers/") && strings.HasSuffix(path, "/json"):
id := strings.TrimSuffix(strings.TrimPrefix(path, "/containers/"), "/json")
_ = json.NewEncoder(w).Encode(map[string]any{
"Id": id,
"State": map[string]any{"Status": "running", "Running": true, "Health": map[string]any{"Status": d.health[id]}},
"Config": map[string]any{"Labels": map[string]string{}},
})
case r.Method == http.MethodPost && strings.HasPrefix(path, "/containers/"):
d.calls = append(d.calls, strings.TrimPrefix(path, "/containers/"))
w.WriteHeader(http.StatusNoContent)
default:
w.WriteHeader(http.StatusNotFound)
_, _ = w.Write([]byte(`{"message":"no such object"}`))
}
}
func matchesDockerFilter(values map[string]bool, candidate string) bool {
if len(values) == 0 {
return true
}
return values[candidate]
}
func (d *multiProviderDaemon) recorded() []string {
d.mu.Lock()
defer d.mu.Unlock()
return append([]string(nil), d.calls...)
}
func TestHealthMonitorKeepsProviderInstallationsIsolatedOnRecovery(t *testing.T) {
daemon := &multiProviderDaemon{
networks: map[string]*multiProviderNetwork{
"provider-a-tenant-t-alpha": {tenantID: "t-alpha", attached: map[string]bool{"client-alpha": true}},
"provider-b-tenant-t-beta": {tenantID: "t-beta", attached: map[string]bool{"client-beta": true}},
},
support: map[string]multiProviderSupport{
"traefik-a": {role: "pulse.provider-msp.role=traefik", ingress: "provider-a"},
"control-a": {role: "pulse.provider-msp.role=control-plane", ingress: "provider-a"},
"traefik-b": {role: "pulse.provider-msp.role=traefik", ingress: "provider-b"},
"control-b": {role: "pulse.provider-msp.role=control-plane", ingress: "provider-b"},
},
health: map[string]string{"client-alpha": "healthy", "client-beta": "unhealthy"},
}
srv := httptest.NewServer(daemon)
t.Cleanup(srv.Close)
t.Setenv("DOCKER_HOST", "tcp://"+strings.TrimPrefix(srv.URL, "http://"))
t.Setenv("DOCKER_TLS_VERIFY", "")
t.Setenv("DOCKER_CERT_PATH", "")
for _, client := range []struct {
provider, tenantID, containerID string
passes int
}{
{provider: "provider-a", tenantID: "t-alpha", containerID: "client-alpha", passes: 1},
{provider: "provider-b", tenantID: "t-beta", containerID: "client-beta", passes: 3},
} {
mgr, err := docker.NewManager(docker.ManagerConfig{
Image: "pulse:test", Network: client.provider, IsolateTenantNetworks: true,
TenantNetworkPrefix: client.provider + "-tenant", BaseDomain: "msp.example.com",
})
if err != nil {
t.Fatalf("NewManager(%s): %v", client.provider, err)
}
t.Cleanup(func() { _ = mgr.Close() })
reg, err := registry.NewTenantRegistry(t.TempDir())
if err != nil {
t.Fatalf("NewTenantRegistry(%s): %v", client.provider, err)
}
t.Cleanup(func() { _ = reg.Close() })
if err := reg.Create(&registry.Tenant{
ID: client.tenantID, AccountID: client.provider, State: registry.TenantStateActive,
ContainerID: client.containerID,
}); err != nil {
t.Fatalf("Create(%s): %v", client.tenantID, err)
}
monitor := NewMonitor(reg, mgr, MonitorConfig{RestartOnFail: true, FailThreshold: 3})
for range client.passes {
monitor.checkAll(context.Background())
}
}
got := daemon.recorded()
want := map[string]int{
"connect traefik-a provider-a-tenant-t-alpha": 1,
"connect control-a provider-a-tenant-t-alpha": 1,
"connect traefik-b provider-b-tenant-t-beta": 1,
"connect control-b provider-b-tenant-t-beta": 1,
"client-beta/restart": 1,
}
for _, call := range got {
if want[call] == 0 {
t.Fatalf("unexpected Docker mutation %q (including cross-client attachment or stop); all=%v", call, got)
}
want[call]--
}
for call, count := range want {
if count != 0 {
t.Fatalf("missing Docker mutation %q; all=%v", call, got)
}
}
}