feat: auto-provision volumes on all machines for global services (#243)

* feat: auto-provision volumes on all machines for global services

* fix: schedule global service volumes only on machines that need them

* fix: use union of eligible machines for volumes shared by global services

* test: restore global with missing volume fails test

* refactor: simplify volume scheduler by removing redundant check

Move isVolumeSharedBetweenGlobalAndReplicated check earlier to fail
fast, then use isVolumeForGlobalService instead of the now-redundant
isVolumeOnlyForGlobalServices function.
This commit is contained in:
Zasda Yusuf Mikail
2026-02-02 12:30:57 +01:00
committed by GitHub
parent a4a8e70f9c
commit 15d9ceb4d1
4 changed files with 562 additions and 28 deletions
+121 -28
View File
@@ -11,6 +11,7 @@ import (
// VolumeScheduler determines what missing volumes should be created and where for a multi-service deployment.
// It satisfies the following constraints:
// - Volumes used by global services will be created on all eligible machines.
// - Services that share a volume must be placed on the same machine where the volume is located.
// If the volume is located on multiple machines, services can be placed on any of them.
// - Services must respect their individual placement constraints.
@@ -128,12 +129,20 @@ func (s *VolumeScheduler) Schedule() (map[string][]api.VolumeSpec, error) {
serviceEligibleMachines[spec.Name] = machineIDs
}
// For each volume that exists on any machine(s) (which shouldn't be created), intersect each service's
// For each volume that exists on any machine(s), intersect each non-global service's
// eligible machines that use the volume with the machines the volume is located on.
// Global services skip this constraint as they need the volume on ALL eligible machines,
// and the volume will be created on machines that don't have it.
//
// Service name -> list of processed volume names (quoted) to format the error message.
quotedServiceVolumes := make(map[string][]string)
for volumeName, volumeMachines := range s.existingVolumeMachines {
// Skip constraint narrowing for global services - they don't need to be constrained
// to machines that already have the volume.
if s.isVolumeForGlobalService(volumeName) {
continue
}
for _, serviceName := range s.volumeServices[volumeName] {
quotedServiceVolumes[serviceName] = append(quotedServiceVolumes[serviceName],
fmt.Sprintf("'%s'", volumeName))
@@ -147,52 +156,97 @@ func (s *VolumeScheduler) Schedule() (map[string][]api.VolumeSpec, error) {
}
}
// Skip constraints propagation for volumes that already exist on machines as the propagation only works
// for missing volumes.
// Skip constraints propagation for:
// 1. Volumes that already exist on machines (for replicated services) as the propagation only works
// for missing volumes. Global service volumes are NOT marked as placed here since they still need
// to be scheduled on machines that don't have them.
// 2. Volumes only used by global services - these need UNION of eligible machines, not intersection.
placedVolumes := make(map[string]struct{})
for volumeName := range s.existingVolumeMachines {
placedVolumes[volumeName] = struct{}{}
if !s.isVolumeForGlobalService(volumeName) {
placedVolumes[volumeName] = struct{}{}
}
}
// Skip constraint propagation for volumes used by global services.
// Also check for invalid configuration: volume shared between global and replicated services.
for volumeName := range s.volumeSpecs {
if s.isVolumeSharedBetweenGlobalAndReplicated(volumeName) {
return nil, fmt.Errorf("volume '%s' cannot be shared between global and replicated services: "+
"global services require the volume on all machines while replicated services require "+
"co-location with the volume", volumeName)
}
if s.isVolumeForGlobalService(volumeName) {
placedVolumes[volumeName] = struct{}{}
}
}
if err := s.propagateConstraintsUntilConvergence(serviceEligibleMachines, placedVolumes); err != nil {
return nil, err
}
// Schedule each missing volume on one of its eligible machines.
// Schedule each missing volume on eligible machines.
// For global services: schedule on ALL eligible machines that don't already have the volume.
// For replicated services: schedule on ONE eligible machine (skip if volume exists anywhere).
scheduledVolumes := make(map[string][]api.VolumeSpec)
for missingVolumeName, missingVolumeSpec := range s.volumeSpecs {
// Skip volumes that already exist on machines.
if _, ok := s.existingVolumeMachines[missingVolumeName]; ok {
continue
}
serviceNames := s.volumeServices[missingVolumeName]
for volumeName, volumeSpec := range s.volumeSpecs {
existingMachines := s.existingVolumeMachines[volumeName]
serviceNames := s.volumeServices[volumeName]
if len(serviceNames) == 0 {
return nil, fmt.Errorf("bug detected: no services using volume '%s'", missingVolumeName)
return nil, fmt.Errorf("bug detected: no services using volume '%s'", volumeName)
}
// Get the current eligible machines (any service using the volume will have the same set after convergence).
eligibleMachines := serviceEligibleMachines[serviceNames[0]]
// Get the eligible machines for this volume.
// For volumes used by global services: compute UNION of all services' eligible machines.
// For other volumes: any service will have the same set after constraint convergence.
var eligibleMachines mapset.Set[string]
if s.isVolumeForGlobalService(volumeName) {
// Compute union of eligible machines for all global services using this volume.
eligibleMachines = mapset.NewSet[string]()
for _, serviceName := range serviceNames {
eligibleMachines = eligibleMachines.Union(serviceEligibleMachines[serviceName])
}
} else {
eligibleMachines = serviceEligibleMachines[serviceNames[0]]
}
if eligibleMachines.Cardinality() == 0 {
return nil, fmt.Errorf("bug detected: no eligible machines for volume '%s'", missingVolumeName)
return nil, fmt.Errorf("bug detected: no eligible machines for volume '%s'", volumeName)
}
// Choose the first machine in the sorted eligible machines to schedule the volume on.
// Sort the eligible machines to ensure deterministic behavior.
sortedEligibleMachines := eligibleMachines.ToSlice()
slices.Sort(sortedEligibleMachines)
machineID := sortedEligibleMachines[0]
// Update constraints for all services that use this volume to be placed on the selected machine.
for _, serviceName := range serviceNames {
serviceEligibleMachines[serviceName] = mapset.NewSet(machineID)
}
placedVolumes[missingVolumeName] = struct{}{}
scheduledVolumes[machineID] = append(scheduledVolumes[machineID], missingVolumeSpec)
// Propagate the updated constraints.
if err := s.propagateConstraintsUntilConvergence(serviceEligibleMachines, placedVolumes); err != nil {
return nil, fmt.Errorf("unexpected error while propagating constraints after "+
"scheduling volume '%s' on machine '%s': %w", missingVolumeName, machineID, err)
if s.isVolumeForGlobalService(volumeName) {
// Global service: schedule volume on eligible machines that don't already have it.
for _, machineID := range sortedEligibleMachines {
if existingMachines != nil && existingMachines.Contains(machineID) {
// Volume already exists on this machine, skip it.
continue
}
scheduledVolumes[machineID] = append(scheduledVolumes[machineID], volumeSpec)
}
// Mark volume as placed - no constraint propagation needed since volume will be on all machines.
placedVolumes[volumeName] = struct{}{}
} else {
// Replicated service: skip if volume already exists on any machine (services will use that location).
if existingMachines != nil && existingMachines.Cardinality() > 0 {
continue
}
// Schedule volume on ONE machine (first in sorted order).
machineID := sortedEligibleMachines[0]
// Update constraints for all services that use this volume to be placed on the selected machine.
for _, serviceName := range serviceNames {
serviceEligibleMachines[serviceName] = mapset.NewSet(machineID)
}
placedVolumes[volumeName] = struct{}{}
scheduledVolumes[machineID] = append(scheduledVolumes[machineID], volumeSpec)
// Propagate the updated constraints.
if err := s.propagateConstraintsUntilConvergence(serviceEligibleMachines, placedVolumes); err != nil {
return nil, fmt.Errorf("unexpected error while propagating constraints after "+
"scheduling volume '%s' on machine '%s': %w", volumeName, machineID, err)
}
}
}
@@ -298,3 +352,42 @@ func (s *VolumeScheduler) propagateConstraintsUntilConvergence(
return nil
}
// isVolumeForGlobalService returns true if any service using this volume is a global service.
func (s *VolumeScheduler) isVolumeForGlobalService(volumeName string) bool {
serviceNames := s.volumeServices[volumeName]
for _, serviceName := range serviceNames {
for _, spec := range s.serviceSpecs {
if spec.Name == serviceName && spec.Mode == api.ServiceModeGlobal {
return true
}
}
}
return false
}
// isVolumeSharedBetweenGlobalAndReplicated returns true if a volume is used by both
// global and replicated services, which is an invalid configuration.
func (s *VolumeScheduler) isVolumeSharedBetweenGlobalAndReplicated(volumeName string) bool {
serviceNames := s.volumeServices[volumeName]
hasGlobal := false
hasReplicated := false
for _, serviceName := range serviceNames {
for _, spec := range s.serviceSpecs {
if spec.Name == serviceName {
mode := spec.Mode
if mode == "" {
mode = api.ServiceModeReplicated
}
if mode == api.ServiceModeGlobal {
hasGlobal = true
} else {
hasReplicated = true
}
}
}
}
return hasGlobal && hasReplicated
}