refactor(client): speed up service list call (N -> 1 broadcast RPC)

This commit is contained in:
Pasha Sviderski committed 2026-10-01 19:52:18 +10:00
1 parent 9c21ae6b08
commit 0e53f505c0
1 file changed
+25 -41
+25 -41
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"maps"
"slices"
"sync"
@@ -320,40 +321,18 @@ func (cli *Client) StartService(ctx context.Context, id string) error {
// ListServices returns a list of all services and their containers.
func (cli *Client) ListServices(ctx context.Context) ([]api.Service, error) {
machines, err := cli.ListMachines(ctx, nil)
if err != nil {
return nil, fmt.Errorf("list machines: %w", err)
}
// Broadcast the container list request to all available machines.
md := metadata.New(nil)
for _, m := range machines {
if m.State == pb.MachineMember_UP || m.State == pb.MachineMember_SUSPECT {
md.Append("machines", m.Machine.Id)
} else {
tui.PrintWarning(fmt.Sprintf("failed to list service containers on machine '%s' (state is %s). "+
"The results may be incomplete.", m.Machine.Name, m.State.String()))
}
}
listCtx := metadata.NewOutgoingContext(ctx, md)
// List all containers including stopped ones.
// Broadcast the container list request to all machines in the cluster.
listCtx := ProxyMachinesContext(ctx, nil)
// List all containers including stopped ones and deployment hooks.
opts := container.ListOptions{All: true}
machineContainers, err := cli.Docker.ListServiceContainers(listCtx, "", opts)
if err != nil {
return nil, fmt.Errorf("list containers: %w", err)
}
// TODO: optimise by extracting services from the list of all containers instead of inspecting each service.
// Most of the code can be reused in both InspectService and ListServices.
// Group containers by service ID using the broadcast results.
servicesByID := make(map[string]api.Service)
for _, mc := range machineContainers {
// NOTE: Metadata should never be nil in practice. This is legacy fallback that will be removed.
if mc.Metadata == nil {
tui.PrintWarning("metadata is missing in response from unknown server")
continue
}
if mc.Metadata.Error != "" {
// TODO: return failed machines in the response.
tui.PrintWarning(fmt.Sprintf("failed to list containers on machine '%s': %s",
@@ -362,25 +341,30 @@ func (cli *Client) ListServices(ctx context.Context) ([]api.Service, error) {
}
for _, ctr := range append(mc.Containers, mc.HookContainers...) {
if _, ok := servicesByID[ctr.ServiceID()]; ok {
continue
}
svc, err := cli.InspectService(ctx, ctr.ServiceID())
if err != nil {
if errors.Is(err, api.ErrNotFound) {
continue
serviceID := ctr.ServiceID()
svc, ok := servicesByID[serviceID]
if !ok {
svc = api.Service{
ID: serviceID,
Name: ctr.ServiceName(),
Mode: ctr.ServiceMode(),
}
return nil, fmt.Errorf("inspect service: %w", err)
}
servicesByID[ctr.ServiceID()] = svc
machineContainer := api.MachineServiceContainer{
MachineID: mc.Metadata.MachineId,
MachineName: mc.Metadata.MachineName,
Container: ctr,
}
if ctr.IsHook() {
svc.HookContainers = append(svc.HookContainers, machineContainer)
} else {
svc.Containers = append(svc.Containers, machineContainer)
}
servicesByID[serviceID] = svc
}
}
services := make([]api.Service, 0, len(servicesByID))
for _, svc := range servicesByID {
services = append(services, svc)
}
return services, nil
return slices.Collect(maps.Values(servicesByID)), nil
}