diff --git a/internal/machine/caddyconfig/controller.go b/internal/machine/caddyconfig/controller.go index d98992e1..5055354e 100644 --- a/internal/machine/caddyconfig/controller.go +++ b/internal/machine/caddyconfig/controller.go @@ -95,10 +95,7 @@ func (c *Controller) filterAvailableContainers( ) ([]api.ServiceContainer, error) { containers := make([]api.ServiceContainer, len(containerRecords)) for i, cr := range containerRecords { - containers[i] = api.ServiceContainer{ - Container: cr.Container, - // TODO: restore ServiceSpec from the container record once it's saved in the store. - } + containers[i] = cr.Container } return containers, nil } diff --git a/internal/machine/cluster.go b/internal/machine/cluster.go index 70b03b50..3ae0ecb6 100644 --- a/internal/machine/cluster.go +++ b/internal/machine/cluster.go @@ -12,7 +12,6 @@ import ( "time" "github.com/cenkalti/backoff/v4" - "github.com/docker/docker/client" "github.com/psviderski/uncloud/internal/machine/api/pb" "github.com/psviderski/uncloud/internal/machine/caddyconfig" "github.com/psviderski/uncloud/internal/machine/constants" @@ -36,10 +35,9 @@ type clusterController struct { wgnet *network.WireGuardNetwork endpointChanges <-chan network.EndpointChangeEvent - server *grpc.Server - corroService corroservice.Service - dockerCli *client.Client - dockerManager *docker.Manager + server *grpc.Server + corroService corroservice.Service + dockerCtrl *docker.Controller // dockerReady is signalled when Docker is configured and ready for containers. dockerReady chan<- struct{} caddyconfigCtrl *caddyconfig.Controller @@ -57,7 +55,7 @@ func newClusterController( store *store.Store, server *grpc.Server, corroService corroservice.Service, - dockerCli *client.Client, + dockerService *docker.Service, dockerReady chan<- struct{}, caddyfileCtrl *caddyconfig.Controller, dnsServer *dns.Server, @@ -77,8 +75,7 @@ func newClusterController( endpointChanges: endpointChanges, server: server, corroService: corroService, - dockerCli: dockerCli, - dockerManager: docker.NewManager(dockerCli, state.ID, store), + dockerCtrl: docker.NewController(state.ID, dockerService, store), dockerReady: dockerReady, caddyconfigCtrl: caddyfileCtrl, dnsServer: dnsServer, @@ -238,11 +235,11 @@ func (cc *clusterController) Run(ctx context.Context) error { // ensureDockerNetwork ensures that the Docker network is configured and ready for containers. func (cc *clusterController) ensureDockerNetwork(ctx context.Context) error { - if err := cc.dockerManager.WaitDaemonReady(ctx); err != nil { + if err := cc.dockerCtrl.WaitDaemonReady(ctx); err != nil { return fmt.Errorf("wait for Docker daemon: %w", err) } - if err := cc.dockerManager.EnsureUncloudNetwork( + if err := cc.dockerCtrl.EnsureUncloudNetwork( ctx, cc.state.Network.Subnet, cc.dnsServer.ListenAddr(), @@ -257,6 +254,7 @@ func (cc *clusterController) ensureDockerNetwork(ctx context.Context) error { } // syncDockerContainers watches local Docker containers and syncs them to the cluster store. +// TODO: move this to the Docker controller. func (cc *clusterController) syncDockerContainers(ctx context.Context) error { // Retry to watch and sync containers until the context is done. boff := backoff.WithContext(backoff.NewExponentialBackOff( @@ -265,7 +263,7 @@ func (cc *clusterController) syncDockerContainers(ctx context.Context) error { backoff.WithMaxElapsedTime(0), ), ctx) watchAndSync := func() error { - if wErr := cc.dockerManager.WatchAndSyncContainers(ctx); wErr != nil { + if wErr := cc.dockerCtrl.WatchAndSyncContainers(ctx); wErr != nil { slog.Error("Failed to watch and sync containers to cluster store, retrying.", "err", wErr) return wErr } @@ -413,7 +411,7 @@ func (cc *clusterController) Cleanup() error { <-cc.stopped var errs []error - if err := cc.dockerManager.Cleanup(); err != nil { + if err := cc.dockerCtrl.Cleanup(); err != nil { errs = append(errs, fmt.Errorf("cleanup Docker resources: %w", err)) } if err := cc.wgnet.Cleanup(); err != nil { diff --git a/internal/machine/dns/resolver.go b/internal/machine/dns/resolver.go index e172297c..dd680362 100644 --- a/internal/machine/dns/resolver.go +++ b/internal/machine/dns/resolver.go @@ -5,12 +5,10 @@ import ( "fmt" "log/slog" "net/netip" - "strings" "sync" "time" "github.com/psviderski/uncloud/internal/machine/store" - "github.com/psviderski/uncloud/pkg/api" ) // ClusterResolver implements Resolver by tracking containers in the cluster and resolving service names @@ -84,17 +82,13 @@ func (r *ClusterResolver) updateServiceIPs(containers []store.ContainerRecord) { continue } - ctr := api.ServiceContainer{Container: record.Container} + ctr := record.Container if ctr.ServiceID() == "" || ctr.ServiceName() == "" { // Container is not part of a service, skip it. continue } - // TODO: remove normalisation after implementing service name validation: - //.https://github.com/psviderski/uncloud/issues/53 - serviceName := strings.ToLower(ctr.ServiceName()) - - newServiceIPs[serviceName] = append(newServiceIPs[serviceName], ip) + newServiceIPs[ctr.ServiceName()] = append(newServiceIPs[ctr.ServiceName()], ip) // Also add the service ID as a valid lookup. newServiceIPs[ctr.ServiceID()] = append(newServiceIPs[ctr.ServiceID()], ip) containersCount++ diff --git a/internal/machine/docker/manager.go b/internal/machine/docker/controller.go similarity index 71% rename from internal/machine/docker/manager.go rename to internal/machine/docker/controller.go index d572b6e2..1b5317d1 100644 --- a/internal/machine/docker/manager.go +++ b/internal/machine/docker/controller.go @@ -7,12 +7,11 @@ import ( "log/slog" "time" - dockercontainer "github.com/docker/docker/api/types/container" + "github.com/docker/docker/api/types/container" "github.com/docker/docker/api/types/events" "github.com/docker/docker/api/types/filters" "github.com/docker/docker/client" "github.com/psviderski/uncloud/internal/machine/store" - "github.com/psviderski/uncloud/pkg/api" ) const ( @@ -24,23 +23,26 @@ const ( SyncInterval = 30 * time.Second ) -type Manager struct { - client *client.Client +// Controller monitors Docker events and synchronises service containers with the cluster store. +type Controller struct { // machineID is the ID of the machine where the managed Docker daemon is running. machineID string + client *client.Client + service *Service store *store.Store } -func NewManager(client *client.Client, machineID string, store *store.Store) *Manager { - return &Manager{ - client: client, +func NewController(machineID string, service *Service, store *store.Store) *Controller { + return &Controller{ machineID: machineID, + client: service.Client, + service: service, store: store, } } // WaitDaemonReady waits for the Docker daemon to start and be ready to serve requests. -func (m *Manager) WaitDaemonReady(ctx context.Context) error { +func (c *Controller) WaitDaemonReady(ctx context.Context) error { ticker := time.NewTicker(1 * time.Second) defer ticker.Stop() @@ -50,7 +52,7 @@ func (m *Manager) WaitDaemonReady(ctx context.Context) error { case <-ctx.Done(): return ctx.Err() case <-ticker.C: - _, err := m.client.Ping(ctx) + _, err := c.client.Ping(ctx) if err == nil { ready = true break @@ -67,7 +69,7 @@ func (m *Manager) WaitDaemonReady(ctx context.Context) error { return nil } -func (m *Manager) WatchAndSyncContainers(ctx context.Context) error { +func (c *Controller) WatchAndSyncContainers(ctx context.Context) error { ctx, cancel := context.WithCancel(ctx) defer cancel() // Filter only local container events. @@ -79,9 +81,9 @@ func (m *Manager) WatchAndSyncContainers(ctx context.Context) error { } // Subscribe to Docker events before running the initial sync to avoid missing any events. - eventCh, errCh := m.client.Events(ctx, opts) + eventCh, errCh := c.service.Client.Events(ctx, opts) slog.Debug("Syncing containers to cluster store before processing Docker events.") - if err := m.syncContainersToStore(ctx); err != nil { + if err := c.syncContainersToStore(ctx); err != nil { // The deferred cancel will stop the event subscription. return fmt.Errorf("sync containers to cluster store: %w", err) } @@ -126,13 +128,13 @@ func (m *Manager) WatchAndSyncContainers(ctx context.Context) error { "container_name", e.Actor.Attributes["name"], "action", e.Action) - if err := m.syncContainersToStore(ctx); err != nil { + if err := c.syncContainersToStore(ctx); err != nil { return fmt.Errorf("sync containers to cluster store: %w", err) } case <-ticker.C: slog.Debug("Syncing containers to cluster store triggered by a regular interval.", "interval", SyncInterval) - if err := m.syncContainersToStore(ctx); err != nil { + if err := c.syncContainersToStore(ctx); err != nil { return fmt.Errorf("sync containers to cluster store: %w", err) } case err := <-errCh: @@ -144,33 +146,16 @@ func (m *Manager) WatchAndSyncContainers(ctx context.Context) error { } } -func (m *Manager) syncContainersToStore(ctx context.Context) error { - storeContainers, err := m.store.ListContainers(ctx, store.ListOptions{MachineIDs: []string{m.machineID}}) +func (c *Controller) syncContainersToStore(ctx context.Context) error { + storeContainers, err := c.store.ListContainers(ctx, store.ListOptions{MachineIDs: []string{c.machineID}}) if err != nil { return fmt.Errorf("list containers from store: %w", err) } - // List only Uncloud service containers identified by their labels. - containerSummaries, err := m.client.ContainerList(ctx, dockercontainer.ListOptions{ - Filters: filters.NewArgs( - filters.Arg("label", api.LabelServiceID), - filters.Arg("label", api.LabelServiceName), - filters.Arg("label", api.LabelManaged), - ), - }) + containers, err := c.service.ListServiceContainers(ctx, "", container.ListOptions{}) if err != nil { // TODO: mark all containers as outdated in the store. - return fmt.Errorf("list Docker containers: %w", err) - } - - // Inspect each container to get the full container details. - containers := make([]api.Container, len(containerSummaries)) - for i, cs := range containerSummaries { - ctr, err := m.client.ContainerInspect(ctx, cs.ID) - if err != nil { - return fmt.Errorf("inspect container '%s': %w", cs.ID, err) - } - containers[i] = api.Container{ContainerJSON: ctr} + return fmt.Errorf("list service containers: %w", err) } // Delete containers from the store that are no longer present in the Docker daemon. @@ -190,15 +175,15 @@ func (m *Manager) syncContainersToStore(ctx context.Context) error { var storeErr error if len(deleteIDs) > 0 { - if err = m.store.DeleteContainers(ctx, store.DeleteOptions{IDs: deleteIDs}); err != nil { + if err = c.store.DeleteContainers(ctx, store.DeleteOptions{IDs: deleteIDs}); err != nil { storeErr = fmt.Errorf("delete containers from store: %w", err) } } // Create or update the current Docker containers in the store. - for _, c := range containers { - if err = m.store.CreateOrUpdateContainer(ctx, c, m.machineID); err != nil { - storeErr = errors.Join(storeErr, fmt.Errorf("create or update container %q: %w", c.ID, err)) + for _, ctr := range containers { + if err = c.store.CreateOrUpdateContainer(ctx, ctr, c.machineID); err != nil { + storeErr = errors.Join(storeErr, fmt.Errorf("create or update container '%s': %w", ctr.ID, err)) } } return storeErr diff --git a/internal/machine/docker/manager_darwin.go b/internal/machine/docker/controller_darwin.go similarity index 62% rename from internal/machine/docker/manager_darwin.go rename to internal/machine/docker/controller_darwin.go index e6ce1381..2dc91319 100644 --- a/internal/machine/docker/manager_darwin.go +++ b/internal/machine/docker/controller_darwin.go @@ -9,11 +9,11 @@ import ( ) // EnsureUncloudNetwork is a stub for Darwin. -func (m *Manager) EnsureUncloudNetwork(ctx context.Context, subnet netip.Prefix, dnsServer netip.Addr) error { +func (c *Controller) EnsureUncloudNetwork(ctx context.Context, subnet netip.Prefix, dnsServer netip.Addr) error { return fmt.Errorf("not supported on Darwin") } // Cleanup is a stub for Darwin. -func (m *Manager) Cleanup() error { +func (c *Controller) Cleanup() error { return fmt.Errorf("not supported on Darwin") } diff --git a/internal/machine/docker/manager_linux.go b/internal/machine/docker/controller_linux.go similarity index 92% rename from internal/machine/docker/manager_linux.go rename to internal/machine/docker/controller_linux.go index 98c5272f..8698ef1e 100644 --- a/internal/machine/docker/manager_linux.go +++ b/internal/machine/docker/controller_linux.go @@ -22,10 +22,10 @@ import ( // EnsureUncloudNetwork creates the Docker bridge network NetworkName with the provided machine subnet // if it doesn't exist. If the network exists but has a different subnet, it removes and recreates the network. // It also configures iptables to allow container access from the WireGuard network. -func (m *Manager) EnsureUncloudNetwork(ctx context.Context, subnet netip.Prefix, dnsServer netip.Addr) error { +func (c *Controller) EnsureUncloudNetwork(ctx context.Context, subnet netip.Prefix, dnsServer netip.Addr) error { // Ensure the Docker network 'uncloud' is created with the correct subnet. needsCreation := false - nw, err := m.client.NetworkInspect(ctx, NetworkName, dnetwork.InspectOptions{}) + nw, err := c.client.NetworkInspect(ctx, NetworkName, dnetwork.InspectOptions{}) if err != nil { if !client.IsErrNotFound(err) { return fmt.Errorf("inspect Docker network '%s': %w", NetworkName, err) @@ -37,7 +37,7 @@ func (m *Manager) EnsureUncloudNetwork(ctx context.Context, subnet netip.Prefix, slog.Info( "Removing Docker network with old subnet.", "name", NetworkName, "subnet", nw.IPAM.Config[0].Subnet, ) - if err = m.client.NetworkRemove(ctx, NetworkName); err != nil { + if err = c.client.NetworkRemove(ctx, NetworkName); err != nil { // It can still fail if the network is in use by a container. Leave it to the user to resolve the issue. return fmt.Errorf("remove Docker network '%s': %w", NetworkName, err) } @@ -45,7 +45,7 @@ func (m *Manager) EnsureUncloudNetwork(ctx context.Context, subnet netip.Prefix, } if needsCreation { - if _, err = m.client.NetworkCreate( + if _, err = c.client.NetworkCreate( ctx, NetworkName, dnetwork.CreateOptions{ Driver: "bridge", Scope: "local", @@ -70,7 +70,7 @@ func (m *Manager) EnsureUncloudNetwork(ctx context.Context, subnet netip.Prefix, } slog.Info("Docker network created.", "name", NetworkName, "subnet", subnet.String()) - if nw, err = m.client.NetworkInspect(ctx, NetworkName, dnetwork.InspectOptions{}); err != nil { + if nw, err = c.client.NetworkInspect(ctx, NetworkName, dnetwork.InspectOptions{}); err != nil { return fmt.Errorf("inspect Docker network '%s': %w", NetworkName, err) } } @@ -168,12 +168,12 @@ func cleanupIptables(bridgeName string, subnet netip.Prefix) error { } // Cleanup removes all uncloud-managed containers and the uncloud Docker network. -func (m *Manager) Cleanup() error { +func (c *Controller) Cleanup() error { ctx := context.Background() var errs []error // Remove uncloud-managed Docker containers. - containers, err := m.client.ContainerList(ctx, dockercontainer.ListOptions{ + containers, err := c.client.ContainerList(ctx, dockercontainer.ListOptions{ All: true, // Include stopped containers. Filters: filters.NewArgs( filters.Arg("label", api.LabelManaged), @@ -186,12 +186,12 @@ func (m *Manager) Cleanup() error { removed := 0 for _, ctr := range containers { - err = m.client.ContainerStop(ctx, ctr.ID, dockercontainer.StopOptions{}) + err = c.client.ContainerStop(ctx, ctr.ID, dockercontainer.StopOptions{}) if err != nil && !client.IsErrNotFound(err) { errs = append(errs, fmt.Errorf("stop container '%s': %w", ctr.ID, err)) } - err = m.client.ContainerRemove(ctx, ctr.ID, dockercontainer.RemoveOptions{ + err = c.client.ContainerRemove(ctx, ctr.ID, dockercontainer.RemoveOptions{ // Remove anonymous volumes created by the container. RemoveVolumes: true, }) @@ -205,7 +205,7 @@ func (m *Manager) Cleanup() error { } // Remove the uncloud Docker network and related iptables rules. - nw, err := m.client.NetworkInspect(ctx, NetworkName, dnetwork.InspectOptions{}) + nw, err := c.client.NetworkInspect(ctx, NetworkName, dnetwork.InspectOptions{}) if err == nil { bridgeName := "br-" + nw.ID[:12] var subnet netip.Prefix @@ -221,7 +221,7 @@ func (m *Manager) Cleanup() error { } } - if err = m.client.NetworkRemove(ctx, NetworkName); err == nil { + if err = c.client.NetworkRemove(ctx, NetworkName); err == nil { slog.Info("Docker network removed.", "name", NetworkName) } else if !client.IsErrNotFound(err) { errs = append(errs, fmt.Errorf("remove Docker network '%s': %w", NetworkName, err)) diff --git a/internal/machine/docker/server.go b/internal/machine/docker/server.go index 4d154ead..b01cc347 100644 --- a/internal/machine/docker/server.go +++ b/internal/machine/docker/server.go @@ -73,10 +73,11 @@ func WithWaitForNetworkReady(waitForNetworkReady func(ctx context.Context) error } } -// NewServer creates a new Docker gRPC server with the provided Docker client. -func NewServer(cli *client.Client, db *sqlx.DB, internalDNSIP func() netip.Addr, opts ...ServerOption) *Server { +// NewServer creates a new Docker gRPC server with the provided Docker service. +// TODO: refactor to use methods from the Docker service instead of querying the database directly. +func NewServer(service *Service, db *sqlx.DB, internalDNSIP func() netip.Addr, opts ...ServerOption) *Server { s := &Server{ - client: cli, + client: service.Client, db: db, internalDNSIP: internalDNSIP, } diff --git a/internal/machine/docker/service.go b/internal/machine/docker/service.go new file mode 100644 index 00000000..dc2cb10c --- /dev/null +++ b/internal/machine/docker/service.go @@ -0,0 +1,103 @@ +package docker + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "log/slog" + + "github.com/docker/docker/api/types/container" + "github.com/docker/docker/api/types/filters" + "github.com/docker/docker/client" + "github.com/jmoiron/sqlx" + "github.com/psviderski/uncloud/pkg/api" +) + +// Service provides higher-level Docker operations that extends Docker API with Uncloud-specific data +// from the machine database. +type Service struct { + Client *client.Client + db *sqlx.DB +} + +// NewService creates a new Docker service instance. +func NewService(client *client.Client, db *sqlx.DB) *Service { + return &Service{ + Client: client, + db: db, + } +} + +// InspectServiceContainer inspects a Docker container and retrieves its associated ServiceSpec +// from the machine database, returning a complete ServiceContainer. +func (s *Service) InspectServiceContainer(ctx context.Context, nameOrID string) (api.ServiceContainer, error) { + var serviceCtr api.ServiceContainer + + ctr, err := s.Client.ContainerInspect(ctx, nameOrID) + if err != nil { + return serviceCtr, err + } + if _, ok := ctr.Config.Labels[api.LabelManaged]; !ok { + return serviceCtr, fmt.Errorf("container '%s' is not managed by Uncloud", nameOrID) + } + + serviceCtr.Container = api.Container{ContainerJSON: ctr} + + // Retrieve ServiceSpec from the machine database. + var specBytes []byte + err = s.db.QueryRowContext(ctx, `SELECT service_spec FROM containers WHERE id = $1`, ctr.ID).Scan(&specBytes) + if err != nil { + if errors.Is(err, sql.ErrNoRows) { + // If this happens, there is a bug in the code, or someone manually removed the container from the DB, + // or created a managed container out of band or by previous uncloud installation. + return serviceCtr, fmt.Errorf("service spec not found for container '%s' in machine DB", ctr.ID) + } + return serviceCtr, fmt.Errorf("get service spec for container '%s' from machine DB: %w", ctr.ID, err) + } + + if err = json.Unmarshal(specBytes, &serviceCtr.ServiceSpec); err != nil { + return serviceCtr, fmt.Errorf("unmarshal service spec for container '%s': %w", ctr.ID, err) + } + + return serviceCtr, nil +} + +// ListServiceContainers lists Docker containers that belong to the service with the given name or ID. +// If serviceIDOrName is empty, all service containers are returned. The opts parameter allows additional filtering. +func (s *Service) ListServiceContainers( + ctx context.Context, serviceNameOrID string, opts container.ListOptions, +) ([]api.ServiceContainer, error) { + if opts.Filters.Len() == 0 { + opts.Filters = filters.NewArgs() + } + // Add labels to existing filters to list only Uncloud-managed service containers. + opts.Filters.Add("label", api.LabelServiceID) + opts.Filters.Add("label", api.LabelManaged) + + containerSummaries, err := s.Client.ContainerList(ctx, opts) + if err != nil { + return nil, err + } + + var containers []api.ServiceContainer + for _, cs := range containerSummaries { + // Filter by service name or ID if provided. + if serviceNameOrID != "" && + cs.Labels[api.LabelServiceID] != serviceNameOrID && + cs.Labels[api.LabelServiceName] != serviceNameOrID { + continue + } + + ctr, err := s.InspectServiceContainer(ctx, cs.ID) + if err != nil { + // Log error but continue with other containers. + slog.Error("Failed to inspect service container.", "service", serviceNameOrID, "id", cs.ID, "err", err) + continue + } + containers = append(containers, ctr) + } + + return containers, nil +} diff --git a/internal/machine/machine.go b/internal/machine/machine.go index bc89ddd7..f47388ee 100644 --- a/internal/machine/machine.go +++ b/internal/machine/machine.go @@ -30,7 +30,6 @@ import ( machinedocker "github.com/psviderski/uncloud/internal/machine/docker" "github.com/psviderski/uncloud/internal/machine/network" "github.com/psviderski/uncloud/internal/machine/store" - "github.com/psviderski/uncloud/pkg/api" "github.com/siderolabs/grpc-proxy/proxy" "golang.org/x/sync/errgroup" "google.golang.org/grpc" @@ -162,7 +161,9 @@ type Machine struct { // store is the cluster store backed by a distributed Corrosion database. store *store.Store cluster *cluster.Cluster - docker *machinedocker.Server + // dockerService provides high-level operations for managing Docker containers. + dockerService *machinedocker.Service + dockerServer *machinedocker.Server // localMachineServer is the gRPC server for the machine API listening on the local Unix socket. localMachineServer *grpc.Server @@ -222,17 +223,14 @@ func NewMachine(config *Config) (*Machine, error) { c := cluster.NewCluster(corroStore, corroAdmin) // Init dependencies for a gRPC Docker server that proxies requests to the local Docker daemon. - dockerCli, err := client.NewClientWithOpts(client.FromEnv, client.WithAPIVersionNegotiation()) - if err != nil { - return nil, fmt.Errorf("create Docker client: %w", err) - } - dbFilePath := filepath.Join(config.DataDir, DBFileName) db, err := NewDB(dbFilePath) if err != nil { return nil, fmt.Errorf("init machine database: %w", err) } + dockerService := machinedocker.NewService(config.DockerClient, db) + // Init a local gRPC proxy server that proxies requests to the local or remote machine API servers. proxyDirector := apiproxy.NewDirector(config.MachineSockPath, constants.MachineAPIPort) localProxyServer := grpc.NewServer( @@ -250,6 +248,7 @@ func NewMachine(config *Config) (*Machine, error) { networkReady: make(chan struct{}), store: corroStore, cluster: c, + dockerService: dockerService, localProxyServer: localProxyServer, proxyDirector: proxyDirector, } @@ -258,10 +257,10 @@ func NewMachine(config *Config) (*Machine, error) { internalDNSIP := func() netip.Addr { return m.IP() } - m.docker = machinedocker.NewServer(dockerCli, db, internalDNSIP, + m.dockerServer = machinedocker.NewServer(dockerService, db, internalDNSIP, machinedocker.WithNetworkReady(m.IsNetworkReady), machinedocker.WithWaitForNetworkReady(m.WaitForNetworkReady)) - m.localMachineServer = newGRPCServer(m, c, m.docker) + m.localMachineServer = newGRPCServer(m, c, m.dockerServer) if m.Initialised() { m.initialised <- struct{}{} @@ -405,7 +404,7 @@ func (m *Machine) Run(ctx context.Context) error { m.store, proxyServer, m.config.CorrosionService, - m.config.DockerClient, + m.dockerService, m.networkReady, caddyconfigCtrl, dnsServer, @@ -569,7 +568,7 @@ func (m *Machine) cleanup() error { } // CheckPrerequisites verifies if the machine meets all necessary system requirements to participate in the cluster. -func (m *Machine) CheckPrerequisites(ctx context.Context, _ *emptypb.Empty) (*pb.CheckPrerequisitesResponse, error) { +func (m *Machine) CheckPrerequisites(_ context.Context, _ *emptypb.Empty) (*pb.CheckPrerequisitesResponse, error) { // Check DNS port (UDP) availability. if err := checkDNSPortAvailable(); err != nil { return &pb.CheckPrerequisitesResponse{ @@ -835,7 +834,7 @@ func (m *Machine) WaitForNetworkReady(ctx context.Context) error { // Reset restores the machine to a clean state, scheduling a graceful shutdown and removing all cluster-related // configuration and resource. The uncloud daemon will restart the machine if managed by systemd. -func (m *Machine) Reset(ctx context.Context, _ *pb.ResetRequest) (*emptypb.Empty, error) { +func (m *Machine) Reset(_ context.Context, _ *pb.ResetRequest) (*emptypb.Empty, error) { if !m.Initialised() { return nil, nil } @@ -890,7 +889,7 @@ func (m *Machine) InspectService( } } - ctr := api.ServiceContainer{Container: records[0].Container} + ctr := records[0].Container svc := &pb.Service{ Id: ctr.ServiceID(), Name: ctr.ServiceName(), diff --git a/internal/machine/store/container.go b/internal/machine/store/container.go index 1f3a8524..c5a5da85 100644 --- a/internal/machine/store/container.go +++ b/internal/machine/store/container.go @@ -23,7 +23,7 @@ const ( ) type ContainerRecord struct { - Container api.Container + Container api.ServiceContainer MachineID string SyncStatus string UpdatedAt time.Time @@ -48,10 +48,12 @@ type DeleteOptions struct { // CreateOrUpdateContainer creates a new container record or updates an existing one in the store database. // The container is associated with the given machine ID that indicates which machine the container is running on. -func (s *Store) CreateOrUpdateContainer(ctx context.Context, ctr api.Container, machineID string) error { +func (s *Store) CreateOrUpdateContainer(ctx context.Context, ctr api.ServiceContainer, machineID string) error { // Remove the environment variables from the container record before storing it in the database // to avoid leaking secrets. ctr.Config.Env = nil + ctr.ServiceSpec.Container.Env = nil + cJSON, err := json.Marshal(ctr) if err != nil { return fmt.Errorf("marshal container: %w", err) @@ -118,7 +120,7 @@ func (s *Store) ListContainers(ctx context.Context, opts ListOptions) ([]Contain return nil, fmt.Errorf("scan container record: %w", err) } - var c api.Container + var c api.ServiceContainer if err = json.Unmarshal([]byte(cJSON), &c); err != nil { return nil, fmt.Errorf("unmarshal container: %w", err) }