mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
worker: serve actor ingress through atunnel
Every worker pod now hosts an atunnel ingress server on :443 and an
atunnel egress listener, both long-lived, with per-activation
Activate/Deactivate bracketing every Run, Restore and Checkpoint.
The nftables rules ateom installs change accordingly:
* The pod-IP:80 -> actor-veth:80 DNAT is gone. Worker port 80 is no
longer an Actor ingress path; the only way in is the mTLS listener
on :443, which authorizes the caller and checks the Actor is the one
currently assigned to this worker.
* A prerouting REDIRECT sends Actor TCP egress to the local atunnel
egress listener, preserving SO_ORIGINAL_DST. It is only installed
when the activation carries an egress gateway address, which nothing
populates yet, so the masquerade path is unchanged and Actor egress
behaves exactly as before.
The ateom Run/Restore protos gain the two fields that arm that path:
egress_gateway_address, which decides whether the redirect is installed
at all, and actor_version, the Actor resource version ate-api observed
when assigning the worker, which atunnel asserts to the egress gateway
as a lower bound on trustworthy Actor metadata. Both are consumed here
and left unset. atelet and ate-api start populating them in the egress
gateway change, which is the point at which they mean anything, so this
change adds no new requirement to the atelet wire contract.
atecontroller gives worker pods the podidentity credential + trust
bundles (the atunnel server identity), the servicedns trust bundle, and
container port 443. podidentitysigner now issues certs with
ExtKeyUsageServerAuth as well as ClientAuth, without which the worker
cannot present its podidentity cert as a TLS server cert and the
gateway handshake fails.
This commit is contained in:
@@ -48,6 +48,13 @@ type ateomOTelSettings struct {
|
||||
MetricExportTimeout string
|
||||
}
|
||||
|
||||
const (
|
||||
atunnelIdentityVolume = "atunnel-identity"
|
||||
atunnelIdentityMountPath = "/run/podidentity.podcert.ate.dev"
|
||||
atunnelEgressTrustVolume = "atunnel-egress-trust"
|
||||
atunnelEgressTrustMountPath = "/run/servicedns.podcert.ate.dev"
|
||||
)
|
||||
|
||||
// buildDeploymentApplyConfig constructs the SSA apply configuration for the
|
||||
// Deployment managed by a WorkerPool. Only fields owned by this controller
|
||||
// are declared here. otel, when it carries an endpoint, is propagated to the
|
||||
@@ -58,23 +65,73 @@ func buildDeploymentApplyConfig(wp *atev1alpha1.WorkerPool, otel ateomOTelSettin
|
||||
WithImage(wp.Spec.AteomImage).
|
||||
WithArgs(
|
||||
"--pod-uid=$(POD_UID)",
|
||||
"--atunnel-listen-address=0.0.0.0:443",
|
||||
"--atunnel-credential-bundle="+atunnelIdentityMountPath+"/credential-bundle.pem",
|
||||
"--atunnel-trust-bundle="+atunnelIdentityMountPath+"/trust-bundle.pem",
|
||||
"--atunnel-egress-listen-address=0.0.0.0:15001",
|
||||
"--atunnel-egress-trust-bundle="+atunnelEgressTrustMountPath+"/trust-bundle.pem",
|
||||
).
|
||||
WithPorts(corev1ac.ContainerPort().
|
||||
WithName("https").
|
||||
WithContainerPort(443).
|
||||
WithProtocol(corev1.ProtocolTCP)).
|
||||
WithSecurityContext(ateomSecurityContext(wp.Spec.SandboxClass)).
|
||||
WithEnv(ateomContainerEnv(otel)...).
|
||||
WithVolumeMounts(corev1ac.VolumeMount().
|
||||
WithName("run-ateom").
|
||||
WithMountPath(ateompath.BasePath).
|
||||
WithMountPropagation(corev1.MountPropagationHostToContainer))
|
||||
WithVolumeMounts(
|
||||
corev1ac.VolumeMount().
|
||||
WithName("run-ateom").
|
||||
WithMountPath(ateompath.BasePath).
|
||||
WithMountPropagation(corev1.MountPropagationHostToContainer),
|
||||
corev1ac.VolumeMount().
|
||||
WithName(atunnelIdentityVolume).
|
||||
WithMountPath(atunnelIdentityMountPath).
|
||||
WithReadOnly(true),
|
||||
corev1ac.VolumeMount().
|
||||
WithName(atunnelEgressTrustVolume).
|
||||
WithMountPath(atunnelEgressTrustMountPath).
|
||||
WithReadOnly(true),
|
||||
)
|
||||
|
||||
podSpecAC := corev1ac.PodSpec().
|
||||
WithSecurityContext(corev1ac.PodSecurityContext().
|
||||
WithRunAsUser(0).
|
||||
WithRunAsGroup(0)).
|
||||
WithVolumes(corev1ac.Volume().
|
||||
WithName("run-ateom").
|
||||
WithHostPath(corev1ac.HostPathVolumeSource().
|
||||
WithPath(ateompath.BasePath).
|
||||
WithType(corev1.HostPathDirectoryOrCreate)))
|
||||
WithVolumes(
|
||||
corev1ac.Volume().
|
||||
WithName("run-ateom").
|
||||
WithHostPath(corev1ac.HostPathVolumeSource().
|
||||
WithPath(ateompath.BasePath).
|
||||
WithType(corev1.HostPathDirectoryOrCreate)),
|
||||
corev1ac.Volume().
|
||||
WithName(atunnelIdentityVolume).
|
||||
WithProjected(corev1ac.ProjectedVolumeSource().
|
||||
WithSources(
|
||||
corev1ac.VolumeProjection().
|
||||
WithPodCertificate(corev1ac.PodCertificateProjection().
|
||||
WithSignerName("podidentity.podcert.ate.dev/identity").
|
||||
WithKeyType("ECDSAP256").
|
||||
WithCredentialBundlePath("credential-bundle.pem")),
|
||||
corev1ac.VolumeProjection().
|
||||
WithClusterTrustBundle(corev1ac.ClusterTrustBundleProjection().
|
||||
WithSignerName("podidentity.podcert.ate.dev/identity").
|
||||
WithLabelSelector(metav1ac.LabelSelector().
|
||||
WithMatchLabels(map[string]string{"podcert.ate.dev/canarying": "live"})).
|
||||
WithPath("trust-bundle.pem")),
|
||||
),
|
||||
),
|
||||
corev1ac.Volume().
|
||||
WithName(atunnelEgressTrustVolume).
|
||||
WithProjected(corev1ac.ProjectedVolumeSource().
|
||||
WithSources(
|
||||
corev1ac.VolumeProjection().
|
||||
WithClusterTrustBundle(corev1ac.ClusterTrustBundleProjection().
|
||||
WithSignerName("servicedns.podcert.ate.dev/identity").
|
||||
WithLabelSelector(metav1ac.LabelSelector().
|
||||
WithMatchLabels(map[string]string{"podcert.ate.dev/canarying": "live"})).
|
||||
WithPath("trust-bundle.pem")),
|
||||
),
|
||||
),
|
||||
)
|
||||
|
||||
applyWorkerPoolPodTemplate(podSpecAC, containerAC, wp.Spec.Template)
|
||||
maybeApplyMicroVMPodShape(podSpecAC, containerAC, wp.Spec.SandboxClass)
|
||||
|
||||
@@ -485,15 +485,57 @@ func expectedDeploymentApplyConfig(mutatePodSpec func(*corev1ac.PodSpecApplyConf
|
||||
WithSecurityContext(corev1ac.PodSecurityContext().
|
||||
WithRunAsUser(0).
|
||||
WithRunAsGroup(0)).
|
||||
WithVolumes(corev1ac.Volume().
|
||||
WithName("run-ateom").
|
||||
WithHostPath(corev1ac.HostPathVolumeSource().
|
||||
WithPath(ateompath.BasePath).
|
||||
WithType(corev1.HostPathDirectoryOrCreate))).
|
||||
WithVolumes(
|
||||
corev1ac.Volume().
|
||||
WithName("run-ateom").
|
||||
WithHostPath(corev1ac.HostPathVolumeSource().
|
||||
WithPath(ateompath.BasePath).
|
||||
WithType(corev1.HostPathDirectoryOrCreate)),
|
||||
corev1ac.Volume().
|
||||
WithName(atunnelIdentityVolume).
|
||||
WithProjected(corev1ac.ProjectedVolumeSource().
|
||||
WithSources(
|
||||
corev1ac.VolumeProjection().
|
||||
WithPodCertificate(corev1ac.PodCertificateProjection().
|
||||
WithSignerName("podidentity.podcert.ate.dev/identity").
|
||||
WithKeyType("ECDSAP256").
|
||||
WithCredentialBundlePath("credential-bundle.pem")),
|
||||
corev1ac.VolumeProjection().
|
||||
WithClusterTrustBundle(corev1ac.ClusterTrustBundleProjection().
|
||||
WithSignerName("podidentity.podcert.ate.dev/identity").
|
||||
WithLabelSelector(metav1ac.LabelSelector().
|
||||
WithMatchLabels(map[string]string{"podcert.ate.dev/canarying": "live"})).
|
||||
WithPath("trust-bundle.pem")),
|
||||
),
|
||||
),
|
||||
corev1ac.Volume().
|
||||
WithName(atunnelEgressTrustVolume).
|
||||
WithProjected(corev1ac.ProjectedVolumeSource().
|
||||
WithSources(
|
||||
corev1ac.VolumeProjection().
|
||||
WithClusterTrustBundle(corev1ac.ClusterTrustBundleProjection().
|
||||
WithSignerName("servicedns.podcert.ate.dev/identity").
|
||||
WithLabelSelector(metav1ac.LabelSelector().
|
||||
WithMatchLabels(map[string]string{"podcert.ate.dev/canarying": "live"})).
|
||||
WithPath("trust-bundle.pem")),
|
||||
),
|
||||
),
|
||||
).
|
||||
WithContainers(corev1ac.Container().
|
||||
WithName("ateom").
|
||||
WithImage(wp.Spec.AteomImage).
|
||||
WithArgs("--pod-uid=$(POD_UID)").
|
||||
WithArgs(
|
||||
"--pod-uid=$(POD_UID)",
|
||||
"--atunnel-listen-address=0.0.0.0:443",
|
||||
"--atunnel-credential-bundle="+atunnelIdentityMountPath+"/credential-bundle.pem",
|
||||
"--atunnel-trust-bundle="+atunnelIdentityMountPath+"/trust-bundle.pem",
|
||||
"--atunnel-egress-listen-address=0.0.0.0:15001",
|
||||
"--atunnel-egress-trust-bundle="+atunnelEgressTrustMountPath+"/trust-bundle.pem",
|
||||
).
|
||||
WithPorts(corev1ac.ContainerPort().
|
||||
WithName("https").
|
||||
WithContainerPort(443).
|
||||
WithProtocol(corev1.ProtocolTCP)).
|
||||
WithSecurityContext(corev1ac.SecurityContext().
|
||||
WithRunAsUser(0).
|
||||
WithRunAsGroup(0).
|
||||
@@ -508,10 +550,20 @@ func expectedDeploymentApplyConfig(mutatePodSpec func(*corev1ac.PodSpecApplyConf
|
||||
WithValueFrom(corev1ac.EnvVarSource().
|
||||
WithFieldRef(corev1ac.ObjectFieldSelector().
|
||||
WithFieldPath("metadata.uid")))).
|
||||
WithVolumeMounts(corev1ac.VolumeMount().
|
||||
WithName("run-ateom").
|
||||
WithMountPath(ateompath.BasePath).
|
||||
WithMountPropagation(corev1.MountPropagationHostToContainer)).
|
||||
WithVolumeMounts(
|
||||
corev1ac.VolumeMount().
|
||||
WithName("run-ateom").
|
||||
WithMountPath(ateompath.BasePath).
|
||||
WithMountPropagation(corev1.MountPropagationHostToContainer),
|
||||
corev1ac.VolumeMount().
|
||||
WithName(atunnelIdentityVolume).
|
||||
WithMountPath(atunnelIdentityMountPath).
|
||||
WithReadOnly(true),
|
||||
corev1ac.VolumeMount().
|
||||
WithName(atunnelEgressTrustVolume).
|
||||
WithMountPath(atunnelEgressTrustMountPath).
|
||||
WithReadOnly(true),
|
||||
).
|
||||
WithResources(corev1ac.ResourceRequirements()))
|
||||
|
||||
podSpecAC.NodeSelector = map[string]string{}
|
||||
|
||||
@@ -146,8 +146,10 @@ func TestWorkerPoolCreatesDeployment(t *testing.T) {
|
||||
if len(dep.OwnerReferences) == 0 || dep.OwnerReferences[0].Name != wp.Name {
|
||||
return false, nil
|
||||
}
|
||||
return len(dep.Spec.Template.Spec.Volumes) == 1 &&
|
||||
dep.Spec.Template.Spec.Volumes[0].Name == "run-ateom", nil
|
||||
return len(dep.Spec.Template.Spec.Volumes) == 3 &&
|
||||
dep.Spec.Template.Spec.Volumes[0].Name == "run-ateom" &&
|
||||
dep.Spec.Template.Spec.Volumes[1].Name == atunnelIdentityVolume &&
|
||||
dep.Spec.Template.Spec.Volumes[2].Name == atunnelEgressTrustVolume, nil
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
+168
-8
@@ -18,9 +18,11 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net"
|
||||
"net/url"
|
||||
"os"
|
||||
"sort"
|
||||
"strings"
|
||||
@@ -31,6 +33,7 @@ import (
|
||||
"github.com/agent-substrate/substrate/internal/ateinterceptors"
|
||||
"github.com/agent-substrate/substrate/internal/ateomnet"
|
||||
"github.com/agent-substrate/substrate/internal/ateompath"
|
||||
"github.com/agent-substrate/substrate/internal/atunnel"
|
||||
"github.com/agent-substrate/substrate/internal/contextlogging"
|
||||
"github.com/agent-substrate/substrate/internal/imagecache"
|
||||
"github.com/agent-substrate/substrate/internal/proto/ateompb"
|
||||
@@ -51,11 +54,22 @@ import (
|
||||
var (
|
||||
podUID = pflag.String("pod-uid", "", "The UID of the current pod")
|
||||
|
||||
atunnelListenAddress = pflag.String("atunnel-listen-address", "0.0.0.0:443", "Address for actor ingress HTTPS")
|
||||
atunnelCredentialBundle = pflag.String("atunnel-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "PEM credential bundle for actor ingress HTTPS")
|
||||
atunnelTrustBundle = pflag.String("atunnel-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "PEM trust bundle for actor ingress clients")
|
||||
atunnelClientIdentity = pflag.String("atunnel-client-identity", "spiffe://cluster.local/ns/ate-system/sa/atenet-router", "SPIFFE identity allowed to call actor ingress HTTPS")
|
||||
atunnelEgressListenAddress = pflag.String("atunnel-egress-listen-address", "0.0.0.0:15001", "Address for transparently intercepted actor egress TCP")
|
||||
atunnelEgressTrustBundle = pflag.String("atunnel-egress-trust-bundle", "/run/servicedns.podcert.ate.dev/trust-bundle.pem", "PEM trust bundle for the egress gateway")
|
||||
|
||||
showVersion = pflag.Bool("version", false, "Print version and exit.")
|
||||
|
||||
reapLock sync.RWMutex
|
||||
)
|
||||
|
||||
// actorHTTPUpstream is the in-sandbox HTTP endpoint atunnel proxies actor
|
||||
// ingress to.
|
||||
const actorHTTPUpstream = "http://" + ateomnet.ActorVethIP + ":80"
|
||||
|
||||
func main() {
|
||||
pflag.Parse()
|
||||
if *showVersion {
|
||||
@@ -137,7 +151,16 @@ func do(ctx context.Context) error {
|
||||
}
|
||||
|
||||
actorLogger := actorlog.NewActorLogger(syncedWriter, metadata.OnGCE())
|
||||
ateomService := NewService(interiorNetNS, actorLogger)
|
||||
upstream, err := url.Parse(actorHTTPUpstream)
|
||||
if err != nil {
|
||||
return fmt.Errorf("while parsing atunnel upstream: %w", err)
|
||||
}
|
||||
atunnelServer, atunnelEgress, atunnelEgressPort, err := runAtunnel(ctx, upstream)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
ateomService := NewService(interiorNetNS, actorLogger, atunnelServer, atunnelEgress, atunnelEgressPort, *atunnelCredentialBundle, *atunnelEgressTrustBundle)
|
||||
|
||||
svr := grpc.NewServer(
|
||||
grpc.StatsHandler(otelgrpc.NewServerHandler()),
|
||||
@@ -154,6 +177,50 @@ func do(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func runAtunnel(ctx context.Context, upstream *url.URL) (*atunnel.Server, *atunnel.Egress, uint16, error) {
|
||||
atunnelServer, err := atunnel.NewServer(atunnel.Config{
|
||||
CredentialBundlePath: *atunnelCredentialBundle,
|
||||
TrustBundlePath: *atunnelTrustBundle,
|
||||
AllowedClientID: *atunnelClientIdentity,
|
||||
Upstream: upstream,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, nil, 0, fmt.Errorf("while configuring atunnel: %w", err)
|
||||
}
|
||||
atunnelListener, err := net.Listen("tcp", *atunnelListenAddress)
|
||||
if err != nil {
|
||||
return nil, nil, 0, fmt.Errorf("while opening atunnel listener: %w", err)
|
||||
}
|
||||
go func() {
|
||||
if err := atunnelServer.Serve(ctx, atunnelListener); err != nil {
|
||||
serverboot.Fatal(ctx, "Failed to serve actor ingress", err)
|
||||
}
|
||||
}()
|
||||
slog.InfoContext(ctx, "atunnel serving", slog.String("address", *atunnelListenAddress))
|
||||
|
||||
atunnelEgress, err := atunnel.NewEgress(atunnel.TCPOriginalDestination)
|
||||
if err != nil {
|
||||
return nil, nil, 0, fmt.Errorf("while configuring atunnel egress: %w", err)
|
||||
}
|
||||
egressListener, err := net.Listen("tcp", *atunnelEgressListenAddress)
|
||||
if err != nil {
|
||||
return nil, nil, 0, fmt.Errorf("while opening atunnel egress listener: %w", err)
|
||||
}
|
||||
egressTCPAddr, ok := egressListener.Addr().(*net.TCPAddr)
|
||||
if !ok || egressTCPAddr.Port < 1 || egressTCPAddr.Port > 65535 {
|
||||
_ = egressListener.Close()
|
||||
return nil, nil, 0, fmt.Errorf("atunnel egress listener has invalid address %q", egressListener.Addr())
|
||||
}
|
||||
atunnelEgressPort := uint16(egressTCPAddr.Port)
|
||||
go func() {
|
||||
if err := atunnelEgress.Serve(ctx, egressListener); err != nil {
|
||||
serverboot.Fatal(ctx, "Failed to serve actor egress", err)
|
||||
}
|
||||
}()
|
||||
slog.InfoContext(ctx, "atunnel egress serving", slog.String("address", *atunnelEgressListenAddress))
|
||||
return atunnelServer, atunnelEgress, atunnelEgressPort, nil
|
||||
}
|
||||
|
||||
// AteomService is a service for shepherding single microvm.
|
||||
type AteomService struct {
|
||||
ateompb.UnimplementedAteomServer
|
||||
@@ -164,22 +231,36 @@ type AteomService struct {
|
||||
|
||||
interiorNetNS netns.NsHandle
|
||||
actorLogger *actorlog.ActorLogger
|
||||
atunnel *atunnel.Server
|
||||
atunnelEgress *atunnel.Egress
|
||||
// atunnelEgressPort is zero when tunneled egress is disabled. Otherwise,
|
||||
// actor TCP connections are transparently redirected to this local port.
|
||||
atunnelEgressPort uint16
|
||||
atunnelCredentialBundle string
|
||||
atunnelEgressTrustBundle string
|
||||
}
|
||||
|
||||
var _ ateompb.AteomServer = (*AteomService)(nil)
|
||||
|
||||
// NewService creates a new AteomService.
|
||||
func NewService(interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger) *AteomService {
|
||||
svc := &AteomService{
|
||||
interiorNetNS: interiorNetNS,
|
||||
actorLogger: actorLogger,
|
||||
func NewService(interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger, atunnelServer *atunnel.Server, atunnelEgress *atunnel.Egress, atunnelEgressPort uint16, credentialBundle, egressTrustBundle string) *AteomService {
|
||||
return &AteomService{
|
||||
interiorNetNS: interiorNetNS,
|
||||
actorLogger: actorLogger,
|
||||
atunnel: atunnelServer,
|
||||
atunnelEgress: atunnelEgress,
|
||||
atunnelEgressPort: atunnelEgressPort,
|
||||
atunnelCredentialBundle: credentialBundle,
|
||||
atunnelEgressTrustBundle: egressTrustBundle,
|
||||
}
|
||||
return svc
|
||||
}
|
||||
|
||||
func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkloadRequest) (resp *ateompb.RunWorkloadResponse, retErr error) {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
if err := s.deactivateActorNetworking(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()}
|
||||
s.actorLogger.EmitLifecycleLog("Actor starting", actorRef, req.GetActorUid(), req.GetActorTemplateNamespace(), req.GetActorTemplateName())
|
||||
@@ -189,7 +270,11 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
|
||||
// * Correct runsc version is downloaded and placed on disk.
|
||||
// * All OCI bundles are set up, including for "pause" container.
|
||||
|
||||
if err := ateomnet.SetupActorNetwork(ctx, ateomnet.NetworkConfig{InteriorNetNS: s.interiorNetNS, DumpNetInfo: true}); err != nil {
|
||||
if err := ateomnet.SetupActorNetwork(ctx, ateomnet.NetworkConfig{
|
||||
InteriorNetNS: s.interiorNetNS,
|
||||
DumpNetInfo: true,
|
||||
EgressRedirectPort: s.egressRedirectPort(req.GetEgressGatewayAddress() != ""),
|
||||
}); err != nil {
|
||||
return nil, fmt.Errorf("while setting up actor network: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
@@ -251,6 +336,9 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
|
||||
if err := readyz.WaitAll(ctx, req.GetSpec().GetContainers(), ateomnet.ActorVethIP); err != nil {
|
||||
return nil, fmt.Errorf("while waiting for container readyz: %w", err)
|
||||
}
|
||||
if err := s.activateActorNetworking(req.GetAtespace(), req.GetActorName(), req.GetActorVersion(), req.GetEgressGatewayAddress()); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
s.actorLogger.EmitLifecycleLog("Actor started", actorRef, req.GetActorUid(), req.GetActorTemplateNamespace(), req.GetActorTemplateName())
|
||||
|
||||
@@ -260,6 +348,9 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
|
||||
func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.CheckpointWorkloadRequest) (*ateompb.CheckpointWorkloadResponse, error) {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
if err := s.deactivateActorNetworking(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()}
|
||||
s.actorLogger.EmitLifecycleLog("Actor checkpointing", actorRef, req.GetActorUid(), req.GetActorTemplateNamespace(), req.GetActorTemplateName())
|
||||
@@ -387,6 +478,9 @@ func (r *runsc) cleanupContainersAfterCheckpoint(ctx context.Context, containers
|
||||
func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.RestoreWorkloadRequest) (resp *ateompb.RestoreWorkloadResponse, retErr error) {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
if err := s.deactivateActorNetworking(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()}
|
||||
s.actorLogger.EmitLifecycleLog("Actor restoring", actorRef, req.GetActorUid(), req.GetActorTemplateNamespace(), req.GetActorTemplateName())
|
||||
@@ -397,7 +491,11 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
|
||||
// * All OCI bundles are set up, including for "pause" container.
|
||||
// * Checkpoint downloaded and placed on disk
|
||||
|
||||
if err := ateomnet.SetupActorNetwork(ctx, ateomnet.NetworkConfig{InteriorNetNS: s.interiorNetNS, DumpNetInfo: true}); err != nil {
|
||||
if err := ateomnet.SetupActorNetwork(ctx, ateomnet.NetworkConfig{
|
||||
InteriorNetNS: s.interiorNetNS,
|
||||
DumpNetInfo: true,
|
||||
EgressRedirectPort: s.egressRedirectPort(req.GetEgressGatewayAddress() != ""),
|
||||
}); err != nil {
|
||||
return nil, fmt.Errorf("while setting up actor network: %w", err)
|
||||
}
|
||||
defer func() {
|
||||
@@ -483,12 +581,74 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
|
||||
if err := readyz.WaitAll(ctx, req.GetSpec().GetContainers(), ateomnet.ActorVethIP); err != nil {
|
||||
return nil, fmt.Errorf("while waiting for container readyz: %w", err)
|
||||
}
|
||||
if err := s.activateActorNetworking(req.GetAtespace(), req.GetActorName(), req.GetActorVersion(), req.GetEgressGatewayAddress()); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
s.actorLogger.EmitLifecycleLog("Actor restored", actorRef, req.GetActorUid(), req.GetActorTemplateNamespace(), req.GetActorTemplateName())
|
||||
|
||||
return &ateompb.RestoreWorkloadResponse{}, nil
|
||||
}
|
||||
|
||||
func (s *AteomService) activateActorNetworking(atespace, actorName string, actorVersion int64, egressGatewayAddress string) error {
|
||||
var egressClient atunnel.EgressDialer
|
||||
if s.atunnelEgress != nil && egressGatewayAddress != "" {
|
||||
serverName, _, err := net.SplitHostPort(egressGatewayAddress)
|
||||
if err != nil {
|
||||
return fmt.Errorf("invalid egress gateway address %q: %w", egressGatewayAddress, err)
|
||||
}
|
||||
egressClient, err = atunnel.NewClient(atunnel.ClientConfig{
|
||||
GatewayAddress: egressGatewayAddress,
|
||||
ServerName: serverName,
|
||||
CredentialBundlePath: s.atunnelCredentialBundle,
|
||||
TrustBundlePath: s.atunnelEgressTrustBundle,
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("while configuring actor egress client: %w", err)
|
||||
}
|
||||
}
|
||||
if s.atunnel != nil {
|
||||
if err := s.atunnel.Activate(atespace, actorName); err != nil {
|
||||
return fmt.Errorf("while activating actor ingress: %w", err)
|
||||
}
|
||||
}
|
||||
if egressClient != nil {
|
||||
if err := s.atunnelEgress.Activate(egressClient, atespace, actorName, actorVersion, ""); err != nil {
|
||||
if s.atunnel != nil {
|
||||
_ = s.atunnel.Deactivate(context.Background())
|
||||
}
|
||||
return fmt.Errorf("while activating actor egress: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *AteomService) deactivateActorNetworking(ctx context.Context) error {
|
||||
// Stop admitting traffic and drain active streams before the Actor network
|
||||
// is torn down. Attempt both directions even if one fails to deactivate.
|
||||
var err error
|
||||
if s.atunnel != nil {
|
||||
err = errors.Join(err, s.atunnel.Deactivate(ctx))
|
||||
}
|
||||
if s.atunnelEgress != nil {
|
||||
err = errors.Join(err, s.atunnelEgress.Deactivate(ctx))
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("while deactivating actor networking: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// egressRedirectPort returns the local atunnel egress listener port when the
|
||||
// activation arms tunneled egress, and zero otherwise, which leaves the
|
||||
// prerouting redirect uninstalled and actor egress on the masquerade path.
|
||||
func (s *AteomService) egressRedirectPort(redirectEgress bool) uint16 {
|
||||
if !redirectEgress {
|
||||
return 0
|
||||
}
|
||||
return s.atunnelEgressPort
|
||||
}
|
||||
|
||||
// setupCgroupDelegation prepares the worker pod's cgroup so runsc can create a
|
||||
// per-actor-container leaf under it with real cpu/memory/pids accounting.
|
||||
//
|
||||
|
||||
@@ -61,6 +61,9 @@ import (
|
||||
func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.CheckpointWorkloadRequest) (*ateompb.CheckpointWorkloadResponse, error) {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
if err := s.deactivateActorNetworking(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()}
|
||||
actorUID := req.GetActorUid()
|
||||
|
||||
+134
-10
@@ -24,10 +24,12 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"flag"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net"
|
||||
"net/url"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
@@ -38,6 +40,7 @@ import (
|
||||
"github.com/agent-substrate/substrate/internal/ateinterceptors"
|
||||
"github.com/agent-substrate/substrate/internal/ateomnet"
|
||||
"github.com/agent-substrate/substrate/internal/ateompath"
|
||||
"github.com/agent-substrate/substrate/internal/atunnel"
|
||||
"github.com/agent-substrate/substrate/internal/proto/ateompb"
|
||||
"github.com/agent-substrate/substrate/internal/serverboot"
|
||||
"github.com/agent-substrate/substrate/internal/version"
|
||||
@@ -55,6 +58,13 @@ var (
|
||||
kataConfig = flag.String("kata-config", "", "Path to a kata configuration.toml (passed to the shim as KATA_CONF_FILE). Empty uses kata's default. atelet generates one pointing at runtime-fetched assets.")
|
||||
kataDebug = flag.Bool("kata-debug", false, "Verbose kata-agent debugging: raise the guest agent log level and forward the guest console (incl. agent logs) into the pod logs.")
|
||||
showVersion = flag.Bool("version", false, "Print version and exit.")
|
||||
|
||||
atunnelListenAddress = flag.String("atunnel-listen-address", "0.0.0.0:443", "Address for actor ingress HTTPS")
|
||||
atunnelCredentialBundle = flag.String("atunnel-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "PEM credential bundle for actor ingress HTTPS")
|
||||
atunnelTrustBundle = flag.String("atunnel-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "PEM trust bundle for actor ingress clients")
|
||||
atunnelClientIdentity = flag.String("atunnel-client-identity", "spiffe://cluster.local/ns/ate-system/sa/atenet-router", "SPIFFE identity allowed to call actor ingress HTTPS")
|
||||
atunnelEgressListenAddress = flag.String("atunnel-egress-listen-address", "0.0.0.0:15001", "Address for transparently intercepted actor egress TCP")
|
||||
atunnelEgressTrustBundle = flag.String("atunnel-egress-trust-bundle", "/run/servicedns.podcert.ate.dev/trust-bundle.pem", "PEM trust bundle for the egress gateway")
|
||||
)
|
||||
|
||||
func main() {
|
||||
@@ -146,12 +156,55 @@ func do(ctx context.Context) error {
|
||||
// logWriter with the runtime logger so the two streams to os.Stdout are
|
||||
// serialized through one SyncedWriter and never interleave-corrupt lines.
|
||||
actorLogger := actorlog.NewActorLogger(logWriter, metadata.OnGCE())
|
||||
upstream, err := url.Parse("http://169.254.17.2:80")
|
||||
if err != nil {
|
||||
return fmt.Errorf("while parsing atunnel upstream: %w", err)
|
||||
}
|
||||
atunnelServer, err := atunnel.NewServer(atunnel.Config{
|
||||
CredentialBundlePath: *atunnelCredentialBundle,
|
||||
TrustBundlePath: *atunnelTrustBundle,
|
||||
AllowedClientID: *atunnelClientIdentity,
|
||||
Upstream: upstream,
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("while configuring atunnel: %w", err)
|
||||
}
|
||||
atunnelListener, err := net.Listen("tcp", *atunnelListenAddress)
|
||||
if err != nil {
|
||||
return fmt.Errorf("while opening atunnel listener: %w", err)
|
||||
}
|
||||
go func() {
|
||||
if err := atunnelServer.Serve(ctx, atunnelListener); err != nil {
|
||||
serverboot.Fatal(ctx, "Failed to serve actor ingress", err)
|
||||
}
|
||||
}()
|
||||
slog.InfoContext(ctx, "atunnel serving", slog.String("address", *atunnelListenAddress))
|
||||
atunnelEgress, err := atunnel.NewEgress(atunnel.TCPOriginalDestination)
|
||||
if err != nil {
|
||||
return fmt.Errorf("while configuring atunnel egress: %w", err)
|
||||
}
|
||||
egressListener, err := net.Listen("tcp", *atunnelEgressListenAddress)
|
||||
if err != nil {
|
||||
return fmt.Errorf("while opening atunnel egress listener: %w", err)
|
||||
}
|
||||
egressTCPAddr, ok := egressListener.Addr().(*net.TCPAddr)
|
||||
if !ok || egressTCPAddr.Port < 1 || egressTCPAddr.Port > 65535 {
|
||||
_ = egressListener.Close()
|
||||
return fmt.Errorf("atunnel egress listener has invalid address %q", egressListener.Addr())
|
||||
}
|
||||
atunnelEgressPort := uint16(egressTCPAddr.Port)
|
||||
go func() {
|
||||
if err := atunnelEgress.Serve(ctx, egressListener); err != nil {
|
||||
serverboot.Fatal(ctx, "Failed to serve actor egress", err)
|
||||
}
|
||||
}()
|
||||
slog.InfoContext(ctx, "atunnel egress serving", slog.String("address", *atunnelEgressListenAddress))
|
||||
|
||||
svr := grpc.NewServer(
|
||||
grpc.StatsHandler(otelgrpc.NewServerHandler()),
|
||||
grpc.UnaryInterceptor(ateinterceptors.InternalServerUnaryInterceptor),
|
||||
)
|
||||
ateompb.RegisterAteomServer(svr, NewService(*podUID, *chBinary, *kataConfig, *kataDebug, interiorNetNS, actorLogger))
|
||||
ateompb.RegisterAteomServer(svr, NewService(*podUID, *chBinary, *kataConfig, *kataDebug, interiorNetNS, actorLogger, atunnelServer, atunnelEgress, atunnelEgressPort, *atunnelCredentialBundle, *atunnelEgressTrustBundle))
|
||||
reflection.Register(svr)
|
||||
|
||||
slog.InfoContext(ctx, "ateom-microvm serving", slog.String("socket", sockPath))
|
||||
@@ -209,7 +262,14 @@ type AteomService struct {
|
||||
// actorLogger forwards the actor container's stdout/stderr to the worker pod's
|
||||
// stdout as ate.dev/*-labeled JSON and emits actor lifecycle events (parity
|
||||
// with ateom-gvisor).
|
||||
actorLogger *actorlog.ActorLogger
|
||||
actorLogger *actorlog.ActorLogger
|
||||
atunnel *atunnel.Server
|
||||
atunnelEgress *atunnel.Egress
|
||||
// atunnelEgressPort is zero when tunneled egress is disabled. Otherwise,
|
||||
// actor TCP connections are transparently redirected to this local port.
|
||||
atunnelEgressPort uint16
|
||||
atunnelCredentialBundle string
|
||||
atunnelEgressTrustBundle string
|
||||
|
||||
// running maps actor UID -> the live micro-VM, kept so CheckpointWorkload can
|
||||
// pause+snapshot+teardown the same sandbox (and RestoreWorkload can track the
|
||||
@@ -220,14 +280,78 @@ type AteomService struct {
|
||||
var _ ateompb.AteomServer = (*AteomService)(nil)
|
||||
|
||||
// NewService creates a new AteomService.
|
||||
func NewService(podUID, chBinary, kataConfig string, kataDebug bool, interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger) *AteomService {
|
||||
func NewService(podUID, chBinary, kataConfig string, kataDebug bool, interiorNetNS netns.NsHandle, actorLogger *actorlog.ActorLogger, atunnelServer *atunnel.Server, atunnelEgress *atunnel.Egress, atunnelEgressPort uint16, credentialBundle, egressTrustBundle string) *AteomService {
|
||||
return &AteomService{
|
||||
podUID: podUID,
|
||||
chBinary: chBinary,
|
||||
kataConfig: kataConfig,
|
||||
kataDebug: kataDebug,
|
||||
interiorNetNS: interiorNetNS,
|
||||
actorLogger: actorLogger,
|
||||
running: map[string]*runningActor{},
|
||||
podUID: podUID,
|
||||
chBinary: chBinary,
|
||||
kataConfig: kataConfig,
|
||||
kataDebug: kataDebug,
|
||||
interiorNetNS: interiorNetNS,
|
||||
actorLogger: actorLogger,
|
||||
atunnel: atunnelServer,
|
||||
atunnelEgress: atunnelEgress,
|
||||
atunnelEgressPort: atunnelEgressPort,
|
||||
atunnelCredentialBundle: credentialBundle,
|
||||
atunnelEgressTrustBundle: egressTrustBundle,
|
||||
running: map[string]*runningActor{},
|
||||
}
|
||||
}
|
||||
|
||||
func (s *AteomService) activateActorNetworking(atespace, actorName string, actorVersion int64, egressGatewayAddress string) error {
|
||||
var egressClient atunnel.EgressDialer
|
||||
if s.atunnelEgress != nil && egressGatewayAddress != "" {
|
||||
serverName, _, err := net.SplitHostPort(egressGatewayAddress)
|
||||
if err != nil {
|
||||
return fmt.Errorf("invalid egress gateway address %q: %w", egressGatewayAddress, err)
|
||||
}
|
||||
egressClient, err = atunnel.NewClient(atunnel.ClientConfig{
|
||||
GatewayAddress: egressGatewayAddress,
|
||||
ServerName: serverName,
|
||||
CredentialBundlePath: s.atunnelCredentialBundle,
|
||||
TrustBundlePath: s.atunnelEgressTrustBundle,
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("while configuring actor egress client: %w", err)
|
||||
}
|
||||
}
|
||||
if s.atunnel != nil {
|
||||
if err := s.atunnel.Activate(atespace, actorName); err != nil {
|
||||
return fmt.Errorf("while activating actor ingress: %w", err)
|
||||
}
|
||||
}
|
||||
if egressClient != nil {
|
||||
if err := s.atunnelEgress.Activate(egressClient, atespace, actorName, actorVersion, ""); err != nil {
|
||||
if s.atunnel != nil {
|
||||
_ = s.atunnel.Deactivate(context.Background())
|
||||
}
|
||||
return fmt.Errorf("while activating actor egress: %w", err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *AteomService) deactivateActorNetworking(ctx context.Context) error {
|
||||
// Stop admitting traffic and drain active streams before the Actor network
|
||||
// is torn down. Attempt both directions even if one fails to deactivate.
|
||||
var err error
|
||||
if s.atunnel != nil {
|
||||
err = errors.Join(err, s.atunnel.Deactivate(ctx))
|
||||
}
|
||||
if s.atunnelEgress != nil {
|
||||
err = errors.Join(err, s.atunnelEgress.Deactivate(ctx))
|
||||
}
|
||||
if err != nil {
|
||||
return fmt.Errorf("while deactivating actor networking: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// egressRedirectPort returns the local atunnel egress listener port when the
|
||||
// activation arms tunneled egress, and zero otherwise, which leaves the
|
||||
// prerouting redirect uninstalled and actor egress on the masquerade path.
|
||||
func (s *AteomService) egressRedirectPort(redirectEgress bool) uint16 {
|
||||
if !redirectEgress {
|
||||
return 0
|
||||
}
|
||||
return s.atunnelEgressPort
|
||||
}
|
||||
|
||||
@@ -54,6 +54,10 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
|
||||
if err := s.deactivateActorNetworking(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
p := actorBootParams{
|
||||
actorRef: resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()},
|
||||
actorUID: req.GetActorUid(),
|
||||
@@ -61,6 +65,9 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
|
||||
templateName: req.GetActorTemplateName(),
|
||||
containers: req.GetSpec().GetContainers(),
|
||||
assetPaths: req.GetRuntimeAssetPaths(),
|
||||
|
||||
actorVersion: req.GetActorVersion(),
|
||||
egressGatewayAddress: req.GetEgressGatewayAddress(),
|
||||
}
|
||||
restoreDir := ateompath.RestoreStateDir(p.actorUID)
|
||||
durableDir := ateompath.DurableDirVolumeMountsDir(p.actorUID)
|
||||
@@ -187,6 +194,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
|
||||
InteriorNetNS: s.interiorNetNS,
|
||||
HostVethHWAddr: hostVethHWAddr,
|
||||
SweepInteriorLinks: true,
|
||||
EgressRedirectPort: s.egressRedirectPort(p.egressGatewayAddress != ""),
|
||||
}); err != nil {
|
||||
return fmt.Errorf("while setting up actor network: %w", err)
|
||||
}
|
||||
@@ -281,6 +289,9 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
|
||||
}
|
||||
}
|
||||
|
||||
if err := s.activateActorNetworking(p.actorRef.Atespace, p.actorRef.Name, p.actorVersion, p.egressGatewayAddress); err != nil {
|
||||
return err
|
||||
}
|
||||
s.running[actorUID] = ra
|
||||
slog.InfoContext(ctx, "Actor restored (overlay rootfs)",
|
||||
slog.String("id", actorUID), slog.Duration("total", time.Since(tStart)))
|
||||
|
||||
@@ -200,6 +200,9 @@ func writeGuestResolvConf(rootfs string) error {
|
||||
func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkloadRequest) (*ateompb.RunWorkloadResponse, error) {
|
||||
s.lock.Lock()
|
||||
defer s.lock.Unlock()
|
||||
if err := s.deactivateActorNetworking(ctx); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
p := actorBootParams{
|
||||
actorRef: resources.ActorRef{Atespace: req.GetAtespace(), Name: req.GetActorName()},
|
||||
@@ -208,6 +211,9 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
|
||||
templateName: req.GetActorTemplateName(),
|
||||
containers: req.GetSpec().GetContainers(),
|
||||
assetPaths: req.GetRuntimeAssetPaths(),
|
||||
|
||||
actorVersion: req.GetActorVersion(),
|
||||
egressGatewayAddress: req.GetEgressGatewayAddress(),
|
||||
}
|
||||
|
||||
s.actorLogger.EmitLifecycleLog("Actor starting", p.actorRef, p.actorUID, p.templateNS, p.templateName)
|
||||
@@ -229,6 +235,12 @@ type actorBootParams struct {
|
||||
templateName string
|
||||
containers []*ateompb.Container
|
||||
assetPaths map[string]string
|
||||
// actorVersion is the Actor resource version ate-api observed when it
|
||||
// assigned this worker; atunnel asserts it to the egress gateway.
|
||||
actorVersion int64
|
||||
// egressGatewayAddress is empty unless an egress gateway is configured, in
|
||||
// which case actor TCP egress is redirected to atunnel's local listener.
|
||||
egressGatewayAddress string
|
||||
}
|
||||
|
||||
// coldBootAttempts is how many times a cold boot is tried when the micro-VM
|
||||
@@ -293,6 +305,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
|
||||
InteriorNetNS: s.interiorNetNS,
|
||||
HostVethHWAddr: hostVethHWAddr,
|
||||
SweepInteriorLinks: true,
|
||||
EgressRedirectPort: s.egressRedirectPort(p.egressGatewayAddress != ""),
|
||||
}); err != nil {
|
||||
return fmt.Errorf("while setting up actor network: %w", err)
|
||||
}
|
||||
@@ -446,6 +459,9 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
|
||||
}
|
||||
|
||||
ra := &runningActor{chCmd: chCmd, vfsdCmd: vfsdCmd, durableVfsdCmd: durableVfsdCmd, apiSocket: apiSocket, baseID: actorUID, logAgent: ac}
|
||||
if err := s.activateActorNetworking(p.actorRef.Atespace, p.actorRef.Name, p.actorVersion, p.egressGatewayAddress); err != nil {
|
||||
return err
|
||||
}
|
||||
s.running[actorUID] = ra
|
||||
|
||||
// Forward each container's stdout/stderr into the pod logs. The overlay workload's
|
||||
|
||||
@@ -30,6 +30,7 @@ import (
|
||||
"github.com/agent-substrate/substrate/internal/localca"
|
||||
"github.com/agent-substrate/substrate/internal/substratex509"
|
||||
certsv1beta1 "k8s.io/api/certificates/v1beta1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/utils/clock"
|
||||
@@ -47,9 +48,18 @@ const (
|
||||
ateletServiceAccount = "atelet"
|
||||
)
|
||||
|
||||
func extKeyUsages(namespace, serviceAccount string) []x509.ExtKeyUsage {
|
||||
// workerPoolLabel marks pods created by atecontroller for a WorkerPool. Worker
|
||||
// pods host the atunnel ingress server, presenting this cert as a TLS server
|
||||
// cert to atenet-router, so they need the serverAuth EKU too. They run as the
|
||||
// actor namespace's default ServiceAccount, so the label is what distinguishes
|
||||
// them rather than their identity.
|
||||
const workerPoolLabel = "ate.dev/worker-pool"
|
||||
|
||||
func extKeyUsages(pod *corev1.Pod, namespace, serviceAccount string) []x509.ExtKeyUsage {
|
||||
usages := []x509.ExtKeyUsage{x509.ExtKeyUsageClientAuth}
|
||||
if namespace == ateletNamespace && serviceAccount == ateletServiceAccount {
|
||||
_, isWorker := pod.ObjectMeta.Labels[workerPoolLabel]
|
||||
isAtelet := namespace == ateletNamespace && serviceAccount == ateletServiceAccount
|
||||
if isAtelet || isWorker {
|
||||
usages = append(usages, x509.ExtKeyUsageServerAuth)
|
||||
}
|
||||
return usages
|
||||
@@ -146,7 +156,7 @@ func (h *Impl) MakeCert(ctx context.Context, pcr *certsv1beta1.PodCertificateReq
|
||||
NotAfter: notAfter,
|
||||
URIs: []*url.URL{spiffeURI},
|
||||
KeyUsage: x509.KeyUsageDigitalSignature,
|
||||
ExtKeyUsage: extKeyUsages(pcr.ObjectMeta.Namespace, pcr.Spec.ServiceAccountName),
|
||||
ExtKeyUsage: extKeyUsages(pod, pcr.ObjectMeta.Namespace, pcr.Spec.ServiceAccountName),
|
||||
// Link the leaf to its issuing CA by key id so verifiers can disambiguate
|
||||
// a multi-CA trust bundle (e.g. valkey trusts both the servicedns and
|
||||
// podidentity CAs).
|
||||
|
||||
@@ -97,6 +97,7 @@ func TestMakeCert(t *testing.T) {
|
||||
namespace string
|
||||
podName string
|
||||
serviceAccount string
|
||||
podLabels map[string]string
|
||||
maxExpirationSeconds int32
|
||||
wantLifetime time.Duration
|
||||
wantURI string
|
||||
@@ -122,6 +123,28 @@ func TestMakeCert(t *testing.T) {
|
||||
NodeUID: "node-uid-1",
|
||||
},
|
||||
},
|
||||
{
|
||||
// Worker pods host the atunnel ingress server, so they serve TLS
|
||||
// despite running as the actor namespace's default ServiceAccount.
|
||||
name: "worker pod also serves",
|
||||
namespace: "ate-demo-counter",
|
||||
podName: "counter-abcde",
|
||||
serviceAccount: "default",
|
||||
podLabels: map[string]string{"ate.dev/worker-pool": "counter"},
|
||||
maxExpirationSeconds: 86400,
|
||||
wantLifetime: 24 * time.Hour,
|
||||
wantURI: "spiffe://cluster.local/ns/ate-demo-counter/sa/default",
|
||||
wantEKUs: []x509.ExtKeyUsage{x509.ExtKeyUsageClientAuth, x509.ExtKeyUsageServerAuth},
|
||||
wantIdentity: &substratex509.PodIdentity{
|
||||
Namespace: "ate-demo-counter",
|
||||
ServiceAccountName: "default",
|
||||
ServiceAccountUID: "sa-uid-1",
|
||||
PodName: "counter-abcde",
|
||||
PodUID: "pod-uid-1",
|
||||
NodeName: "node-1",
|
||||
NodeUID: "node-uid-1",
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "ordinary workload is client-only",
|
||||
namespace: "default",
|
||||
@@ -197,6 +220,7 @@ func TestMakeCert(t *testing.T) {
|
||||
}
|
||||
|
||||
pod, pcr := makePodAndPCR(tc.namespace, tc.podName, tc.serviceAccount, tc.maxExpirationSeconds)
|
||||
pod.ObjectMeta.Labels = tc.podLabels
|
||||
pcr.Spec.StubPKCS10Request = stubCSR(t, subjectPriv)
|
||||
|
||||
kc := fake.NewSimpleClientset(pod, pcr)
|
||||
|
||||
+42
-62
@@ -232,27 +232,27 @@ func EnableIPv4Forwarding() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// InstallActorNftablesRules configures the NAT and filtering rules for the actor.
|
||||
func InstallActorNftablesRules(podIP net.IP) error {
|
||||
// InstallActorNftablesRules configures the NAT and filtering rules for the
|
||||
// actor. egressPort, when non-zero, is the local atunnel egress listener actor
|
||||
// TCP egress is redirected to; zero leaves the redirect uninstalled.
|
||||
func InstallActorNftablesRules(egressPort uint16) error {
|
||||
// Install a dedicated nftables table for the active actor. Keeping all
|
||||
// rules in an ateom-owned table makes cleanup simple and avoids mutating
|
||||
// Kubernetes or CNI-managed chains directly.
|
||||
//
|
||||
// TODO: Add IPv6 veth addressing, forwarding, and nftables rules once actor
|
||||
// networking supports dual-stack pods. The current compatibility path is
|
||||
// IPv4-only.
|
||||
// networking supports dual-stack pods. The current actor network is IPv4-only.
|
||||
//
|
||||
// The temporary compatibility rules do three things:
|
||||
// The rules do three things:
|
||||
//
|
||||
// * postrouting: masquerade actor egress from 169.254.17.2 behind the worker
|
||||
// pod IP so replies route back to the pod.
|
||||
// * prerouting: DNAT traffic sent to the worker pod IP on TCP/80 to the
|
||||
// actor veth IP on TCP/80, preserving existing inbound behavior.
|
||||
// * prerouting: redirect new actor TCP connections to atunnel's local
|
||||
// listener. REDIRECT preserves SO_ORIGINAL_DST for the CONNECT authority.
|
||||
// * postrouting: masquerade traffic not handled by the TCP tunnel, notably
|
||||
// DNS over UDP, so hostname resolution continues to work.
|
||||
// * forward: accept forwarded packets between the actor veth and pod eth0.
|
||||
//
|
||||
// This is not the final egress policy path. The later AgentGateway phase
|
||||
// should replace the broad masquerade path with transparent TCP capture and
|
||||
// default-deny rules.
|
||||
// TODO: Restrict the compatibility masquerade to DNS traffic sent to the
|
||||
// configured cluster resolver and drop all other non-tunneled actor egress.
|
||||
if err := RemoveActorNftablesRules(); err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -271,32 +271,9 @@ func InstallActorNftablesRules(podIP net.IP) error {
|
||||
Hooknum: nftables.ChainHookPrerouting,
|
||||
Priority: nftables.ChainPriorityNATDest,
|
||||
})
|
||||
// TODO: Support inbound UDP DNAT for actors that expose UDP protocols such
|
||||
// as QUIC.
|
||||
// TODO: Replace the hard-coded HTTP port with the actor's configured
|
||||
// inbound ports, either by adding one rule per port or by matching a set.
|
||||
preroutingExprs := append(IPDestinationEqual(podIP.String()), TCPDestinationPortEqual(80)...)
|
||||
preroutingExprs = append(preroutingExprs,
|
||||
&expr.Immediate{
|
||||
Register: 1,
|
||||
Data: net.ParseIP(ActorVethIP).To4(),
|
||||
},
|
||||
&expr.Immediate{
|
||||
Register: 2,
|
||||
Data: binaryutil.BigEndian.PutUint16(80),
|
||||
},
|
||||
&expr.NAT{
|
||||
Type: expr.NATTypeDestNAT,
|
||||
Family: unix.NFPROTO_IPV4,
|
||||
RegAddrMin: 1,
|
||||
RegProtoMin: 2,
|
||||
},
|
||||
)
|
||||
c.AddRule(&nftables.Rule{
|
||||
Table: table,
|
||||
Chain: prerouting,
|
||||
Exprs: preroutingExprs,
|
||||
})
|
||||
if redirectRule := ActorEgressRedirectRule(table, prerouting, egressPort); redirectRule != nil {
|
||||
c.AddRule(redirectRule)
|
||||
}
|
||||
|
||||
postrouting := c.AddChain(&nftables.Chain{
|
||||
Name: "postrouting",
|
||||
@@ -361,10 +338,6 @@ func IPSourceEqual(ip string) []expr.Any {
|
||||
return IPPayloadEqual(12, ip)
|
||||
}
|
||||
|
||||
func IPDestinationEqual(ip string) []expr.Any {
|
||||
return IPPayloadEqual(16, ip)
|
||||
}
|
||||
|
||||
func IPPayloadEqual(offset uint32, ip string) []expr.Any {
|
||||
return []expr.Any{
|
||||
&expr.Payload{
|
||||
@@ -381,7 +354,7 @@ func IPPayloadEqual(offset uint32, ip string) []expr.Any {
|
||||
}
|
||||
}
|
||||
|
||||
func TCPDestinationPortEqual(port uint16) []expr.Any {
|
||||
func TCPProtocol() []expr.Any {
|
||||
return []expr.Any{
|
||||
&expr.Meta{Key: expr.MetaKeyL4PROTO, Register: 1},
|
||||
&expr.Cmp{
|
||||
@@ -389,18 +362,25 @@ func TCPDestinationPortEqual(port uint16) []expr.Any {
|
||||
Register: 1,
|
||||
Data: []byte{unix.IPPROTO_TCP},
|
||||
},
|
||||
&expr.Payload{
|
||||
DestRegister: 1,
|
||||
Base: expr.PayloadBaseTransportHeader,
|
||||
Offset: 2,
|
||||
Len: 2,
|
||||
},
|
||||
&expr.Cmp{
|
||||
Op: expr.CmpOpEq,
|
||||
}
|
||||
}
|
||||
|
||||
// ActorEgressRedirectRule returns the prerouting rule that redirects actor TCP
|
||||
// egress to the local atunnel egress listener on port, or nil when port is zero
|
||||
// (tunneled egress disabled, so actor egress stays on the masquerade path).
|
||||
func ActorEgressRedirectRule(table *nftables.Table, chain *nftables.Chain, port uint16) *nftables.Rule {
|
||||
if port == 0 {
|
||||
return nil
|
||||
}
|
||||
exprs := append(IPSourceEqual(ActorVethIP), TCPProtocol()...)
|
||||
exprs = append(exprs,
|
||||
&expr.Immediate{
|
||||
Register: 1,
|
||||
Data: binaryutil.BigEndian.PutUint16(port),
|
||||
},
|
||||
}
|
||||
&expr.Redir{RegisterProtoMin: 1},
|
||||
)
|
||||
return &nftables.Rule{Table: table, Chain: chain, Exprs: exprs}
|
||||
}
|
||||
|
||||
// CreateNetNSWithoutSwitching creates a named netns and returns its handle,
|
||||
@@ -499,6 +479,12 @@ type NetworkConfig struct {
|
||||
// DumpNetInfo indicates whether to dump network information to the logs for debugging purposes.
|
||||
// Used by: gVisor.
|
||||
DumpNetInfo bool
|
||||
|
||||
// EgressRedirectPort is the local atunnel egress listener port actor TCP
|
||||
// egress is redirected to. Zero installs no redirect, leaving actor egress
|
||||
// on the masquerade path.
|
||||
// Used by: Both gVisor and MicroVM.
|
||||
EgressRedirectPort uint16
|
||||
}
|
||||
|
||||
// SetupActorNetwork builds a fresh point-to-point network between the worker
|
||||
@@ -511,10 +497,9 @@ func SetupActorNetwork(ctx context.Context, cfg NetworkConfig) (retErr error) {
|
||||
// the worker-side veth address. This replaces the old behavior of moving the
|
||||
// Kubernetes-provided eth0 out of the worker pod.
|
||||
//
|
||||
// The nftables rules installed here are a compatibility bridge for the
|
||||
// current router assumptions: actor egress is masqueraded behind the worker
|
||||
// pod IP, and inbound traffic to the worker pod's HTTP port is DNAT'd to the
|
||||
// actor veth IP.
|
||||
// The nftables rules installed here redirect actor TCP egress to atunnel
|
||||
// when configured and masquerade traffic the TCP tunnel does not handle
|
||||
// (notably DNS over UDP).
|
||||
//
|
||||
// Clean up stale state from a failed prior activation before creating the
|
||||
// next actor-side network. The worker currently runs one actor at a time.
|
||||
@@ -529,11 +514,6 @@ func SetupActorNetwork(ctx context.Context, cfg NetworkConfig) (retErr error) {
|
||||
}
|
||||
}()
|
||||
|
||||
podIP, err := PodIPv4()
|
||||
if err != nil {
|
||||
return fmt.Errorf("while resolving pod IPv4 address: %w", err)
|
||||
}
|
||||
|
||||
if cfg.SweepInteriorLinks {
|
||||
if err := NetNSDo(ctx, cfg.InteriorNetNS, func(ctx context.Context) error {
|
||||
links, err := netlink.LinkList()
|
||||
@@ -594,7 +574,7 @@ func SetupActorNetwork(ctx context.Context, cfg NetworkConfig) (retErr error) {
|
||||
if err := EnableIPv4Forwarding(); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := InstallActorNftablesRules(podIP); err != nil {
|
||||
if err := InstallActorNftablesRules(cfg.EgressRedirectPort); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
@@ -231,7 +231,7 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
s.reject(w)
|
||||
return
|
||||
}
|
||||
atespace, actorName, err := resources.ParseActorDNSName(host)
|
||||
actorRef, err := resources.ParseActorDNSName(host)
|
||||
if err != nil {
|
||||
s.reject(w)
|
||||
return
|
||||
@@ -239,7 +239,7 @@ func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
|
||||
s.mu.Lock()
|
||||
active := s.active
|
||||
if active == nil || active.atespace != atespace || active.actorName != actorName {
|
||||
if active == nil || active.atespace != actorRef.Atespace || active.actorName != actorRef.Name {
|
||||
s.mu.Unlock()
|
||||
s.reject(w)
|
||||
return
|
||||
|
||||
@@ -91,21 +91,26 @@ func (SnapshotScope) EnumDescriptor() ([]byte, []int) {
|
||||
}
|
||||
|
||||
type RunWorkloadRequest struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
Atespace string `protobuf:"bytes,1,opt,name=atespace,proto3" json:"atespace,omitempty"`
|
||||
ActorName string `protobuf:"bytes,2,opt,name=actor_name,json=actorName,proto3" json:"actor_name,omitempty"`
|
||||
ActorUid string `protobuf:"bytes,3,opt,name=actor_uid,json=actorUid,proto3" json:"actor_uid,omitempty"`
|
||||
ActorTemplateNamespace string `protobuf:"bytes,4,opt,name=actor_template_namespace,json=actorTemplateNamespace,proto3" json:"actor_template_namespace,omitempty"`
|
||||
ActorTemplateName string `protobuf:"bytes,5,opt,name=actor_template_name,json=actorTemplateName,proto3" json:"actor_template_name,omitempty"`
|
||||
RunscPath string `protobuf:"bytes,6,opt,name=runsc_path,json=runscPath,proto3" json:"runsc_path,omitempty"`
|
||||
Spec *WorkloadSpec `protobuf:"bytes,7,opt,name=spec,proto3" json:"spec,omitempty"`
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
Atespace string `protobuf:"bytes,1,opt,name=atespace,proto3" json:"atespace,omitempty"`
|
||||
ActorName string `protobuf:"bytes,2,opt,name=actor_name,json=actorName,proto3" json:"actor_name,omitempty"`
|
||||
ActorUid string `protobuf:"bytes,3,opt,name=actor_uid,json=actorUid,proto3" json:"actor_uid,omitempty"`
|
||||
// Actor resource version observed by ate-api when assigning this worker.
|
||||
ActorVersion int64 `protobuf:"varint,9,opt,name=actor_version,json=actorVersion,proto3" json:"actor_version,omitempty"`
|
||||
ActorTemplateNamespace string `protobuf:"bytes,4,opt,name=actor_template_namespace,json=actorTemplateNamespace,proto3" json:"actor_template_namespace,omitempty"`
|
||||
ActorTemplateName string `protobuf:"bytes,5,opt,name=actor_template_name,json=actorTemplateName,proto3" json:"actor_template_name,omitempty"`
|
||||
RunscPath string `protobuf:"bytes,6,opt,name=runsc_path,json=runscPath,proto3" json:"runsc_path,omitempty"`
|
||||
Spec *WorkloadSpec `protobuf:"bytes,7,opt,name=spec,proto3" json:"spec,omitempty"`
|
||||
// runtime_asset_paths maps a runtime asset name (e.g. "cloud-hypervisor",
|
||||
// "virtiofsd", "kata-kernel", "kata-image", "kata-config")
|
||||
// to the local on-disk path atelet fetched it to (content-addressed, like
|
||||
// runsc_path). Empty for the gVisor runtime, which uses runsc_path.
|
||||
RuntimeAssetPaths map[string]string `protobuf:"bytes,8,rep,name=runtime_asset_paths,json=runtimeAssetPaths,proto3" json:"runtime_asset_paths,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
// Remote egress gateway selected for this activation. When absent, actor
|
||||
// traffic uses direct egress instead of being redirected through atunnel.
|
||||
EgressGatewayAddress *string `protobuf:"bytes,10,opt,name=egress_gateway_address,json=egressGatewayAddress,proto3,oneof" json:"egress_gateway_address,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
|
||||
func (x *RunWorkloadRequest) Reset() {
|
||||
@@ -159,6 +164,13 @@ func (x *RunWorkloadRequest) GetActorUid() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func (x *RunWorkloadRequest) GetActorVersion() int64 {
|
||||
if x != nil {
|
||||
return x.ActorVersion
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func (x *RunWorkloadRequest) GetActorTemplateNamespace() string {
|
||||
if x != nil {
|
||||
return x.ActorTemplateNamespace
|
||||
@@ -194,6 +206,13 @@ func (x *RunWorkloadRequest) GetRuntimeAssetPaths() map[string]string {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (x *RunWorkloadRequest) GetEgressGatewayAddress() string {
|
||||
if x != nil && x.EgressGatewayAddress != nil {
|
||||
return *x.EgressGatewayAddress
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// WorkloadSpec parallels Pod, but with far fewer configurable fields.
|
||||
type WorkloadSpec struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
@@ -669,23 +688,28 @@ func (x *CheckpointWorkloadResponse) GetSnapshotFiles() []string {
|
||||
}
|
||||
|
||||
type RestoreWorkloadRequest struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
Atespace string `protobuf:"bytes,1,opt,name=atespace,proto3" json:"atespace,omitempty"`
|
||||
ActorName string `protobuf:"bytes,2,opt,name=actor_name,json=actorName,proto3" json:"actor_name,omitempty"`
|
||||
ActorUid string `protobuf:"bytes,3,opt,name=actor_uid,json=actorUid,proto3" json:"actor_uid,omitempty"`
|
||||
ActorTemplateNamespace string `protobuf:"bytes,4,opt,name=actor_template_namespace,json=actorTemplateNamespace,proto3" json:"actor_template_namespace,omitempty"`
|
||||
ActorTemplateName string `protobuf:"bytes,5,opt,name=actor_template_name,json=actorTemplateName,proto3" json:"actor_template_name,omitempty"`
|
||||
RunscPath string `protobuf:"bytes,6,opt,name=runsc_path,json=runscPath,proto3" json:"runsc_path,omitempty"`
|
||||
Spec *WorkloadSpec `protobuf:"bytes,7,opt,name=spec,proto3" json:"spec,omitempty"`
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
Atespace string `protobuf:"bytes,1,opt,name=atespace,proto3" json:"atespace,omitempty"`
|
||||
ActorName string `protobuf:"bytes,2,opt,name=actor_name,json=actorName,proto3" json:"actor_name,omitempty"`
|
||||
ActorUid string `protobuf:"bytes,3,opt,name=actor_uid,json=actorUid,proto3" json:"actor_uid,omitempty"`
|
||||
// Actor resource version observed by ate-api when assigning this worker.
|
||||
ActorVersion int64 `protobuf:"varint,11,opt,name=actor_version,json=actorVersion,proto3" json:"actor_version,omitempty"`
|
||||
ActorTemplateNamespace string `protobuf:"bytes,4,opt,name=actor_template_namespace,json=actorTemplateNamespace,proto3" json:"actor_template_namespace,omitempty"`
|
||||
ActorTemplateName string `protobuf:"bytes,5,opt,name=actor_template_name,json=actorTemplateName,proto3" json:"actor_template_name,omitempty"`
|
||||
RunscPath string `protobuf:"bytes,6,opt,name=runsc_path,json=runscPath,proto3" json:"runsc_path,omitempty"`
|
||||
Spec *WorkloadSpec `protobuf:"bytes,7,opt,name=spec,proto3" json:"spec,omitempty"`
|
||||
// The object storage URI prefix of the snapshot to restore.
|
||||
SnapshotUriPrefix string `protobuf:"bytes,8,opt,name=snapshot_uri_prefix,json=snapshotUriPrefix,proto3" json:"snapshot_uri_prefix,omitempty"`
|
||||
// runtime_asset_paths maps a runtime asset name to the local on-disk path
|
||||
// atelet fetched it to (see RunWorkloadRequest). Empty for gVisor.
|
||||
RuntimeAssetPaths map[string]string `protobuf:"bytes,9,rep,name=runtime_asset_paths,json=runtimeAssetPaths,proto3" json:"runtime_asset_paths,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"`
|
||||
// What content to restore from the snapshot.
|
||||
Scope SnapshotScope `protobuf:"varint,10,opt,name=scope,proto3,enum=ateom.SnapshotScope" json:"scope,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
Scope SnapshotScope `protobuf:"varint,10,opt,name=scope,proto3,enum=ateom.SnapshotScope" json:"scope,omitempty"`
|
||||
// Remote egress gateway selected for this activation. When absent, actor
|
||||
// traffic uses direct egress instead of being redirected through atunnel.
|
||||
EgressGatewayAddress *string `protobuf:"bytes,12,opt,name=egress_gateway_address,json=egressGatewayAddress,proto3,oneof" json:"egress_gateway_address,omitempty"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
sizeCache protoimpl.SizeCache
|
||||
}
|
||||
|
||||
func (x *RestoreWorkloadRequest) Reset() {
|
||||
@@ -739,6 +763,13 @@ func (x *RestoreWorkloadRequest) GetActorUid() string {
|
||||
return ""
|
||||
}
|
||||
|
||||
func (x *RestoreWorkloadRequest) GetActorVersion() int64 {
|
||||
if x != nil {
|
||||
return x.ActorVersion
|
||||
}
|
||||
return 0
|
||||
}
|
||||
|
||||
func (x *RestoreWorkloadRequest) GetActorTemplateNamespace() string {
|
||||
if x != nil {
|
||||
return x.ActorTemplateNamespace
|
||||
@@ -788,6 +819,13 @@ func (x *RestoreWorkloadRequest) GetScope() SnapshotScope {
|
||||
return SnapshotScope_SNAPSHOT_SCOPE_UNSPECIFIED
|
||||
}
|
||||
|
||||
func (x *RestoreWorkloadRequest) GetEgressGatewayAddress() string {
|
||||
if x != nil && x.EgressGatewayAddress != nil {
|
||||
return *x.EgressGatewayAddress
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
type RestoreWorkloadResponse struct {
|
||||
state protoimpl.MessageState `protogen:"open.v1"`
|
||||
unknownFields protoimpl.UnknownFields
|
||||
@@ -828,21 +866,25 @@ var File_ateom_proto protoreflect.FileDescriptor
|
||||
|
||||
const file_ateom_proto_rawDesc = "" +
|
||||
"\n" +
|
||||
"\vateom.proto\x12\x05ateom\"\xc6\x03\n" +
|
||||
"\vateom.proto\x12\x05ateom\"\xc1\x04\n" +
|
||||
"\x12RunWorkloadRequest\x12\x1a\n" +
|
||||
"\batespace\x18\x01 \x01(\tR\batespace\x12\x1d\n" +
|
||||
"\n" +
|
||||
"actor_name\x18\x02 \x01(\tR\tactorName\x12\x1b\n" +
|
||||
"\tactor_uid\x18\x03 \x01(\tR\bactorUid\x128\n" +
|
||||
"\tactor_uid\x18\x03 \x01(\tR\bactorUid\x12#\n" +
|
||||
"\ractor_version\x18\t \x01(\x03R\factorVersion\x128\n" +
|
||||
"\x18actor_template_namespace\x18\x04 \x01(\tR\x16actorTemplateNamespace\x12.\n" +
|
||||
"\x13actor_template_name\x18\x05 \x01(\tR\x11actorTemplateName\x12\x1d\n" +
|
||||
"\n" +
|
||||
"runsc_path\x18\x06 \x01(\tR\trunscPath\x12'\n" +
|
||||
"\x04spec\x18\a \x01(\v2\x13.ateom.WorkloadSpecR\x04spec\x12`\n" +
|
||||
"\x13runtime_asset_paths\x18\b \x03(\v20.ateom.RunWorkloadRequest.RuntimeAssetPathsEntryR\x11runtimeAssetPaths\x1aD\n" +
|
||||
"\x13runtime_asset_paths\x18\b \x03(\v20.ateom.RunWorkloadRequest.RuntimeAssetPathsEntryR\x11runtimeAssetPaths\x129\n" +
|
||||
"\x16egress_gateway_address\x18\n" +
|
||||
" \x01(\tH\x00R\x14egressGatewayAddress\x88\x01\x01\x1aD\n" +
|
||||
"\x16RuntimeAssetPathsEntry\x12\x10\n" +
|
||||
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
|
||||
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"@\n" +
|
||||
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01B\x19\n" +
|
||||
"\x17_egress_gateway_address\"@\n" +
|
||||
"\fWorkloadSpec\x120\n" +
|
||||
"\n" +
|
||||
"containers\x18\x01 \x03(\v2\x10.ateom.ContainerR\n" +
|
||||
@@ -880,12 +922,13 @@ const file_ateom_proto_rawDesc = "" +
|
||||
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
|
||||
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"C\n" +
|
||||
"\x1aCheckpointWorkloadResponse\x12%\n" +
|
||||
"\x0esnapshot_files\x18\x01 \x03(\tR\rsnapshotFiles\"\xaa\x04\n" +
|
||||
"\x0esnapshot_files\x18\x01 \x03(\tR\rsnapshotFiles\"\xa5\x05\n" +
|
||||
"\x16RestoreWorkloadRequest\x12\x1a\n" +
|
||||
"\batespace\x18\x01 \x01(\tR\batespace\x12\x1d\n" +
|
||||
"\n" +
|
||||
"actor_name\x18\x02 \x01(\tR\tactorName\x12\x1b\n" +
|
||||
"\tactor_uid\x18\x03 \x01(\tR\bactorUid\x128\n" +
|
||||
"\tactor_uid\x18\x03 \x01(\tR\bactorUid\x12#\n" +
|
||||
"\ractor_version\x18\v \x01(\x03R\factorVersion\x128\n" +
|
||||
"\x18actor_template_namespace\x18\x04 \x01(\tR\x16actorTemplateNamespace\x12.\n" +
|
||||
"\x13actor_template_name\x18\x05 \x01(\tR\x11actorTemplateName\x12\x1d\n" +
|
||||
"\n" +
|
||||
@@ -894,10 +937,12 @@ const file_ateom_proto_rawDesc = "" +
|
||||
"\x13snapshot_uri_prefix\x18\b \x01(\tR\x11snapshotUriPrefix\x12d\n" +
|
||||
"\x13runtime_asset_paths\x18\t \x03(\v24.ateom.RestoreWorkloadRequest.RuntimeAssetPathsEntryR\x11runtimeAssetPaths\x12*\n" +
|
||||
"\x05scope\x18\n" +
|
||||
" \x01(\x0e2\x14.ateom.SnapshotScopeR\x05scope\x1aD\n" +
|
||||
" \x01(\x0e2\x14.ateom.SnapshotScopeR\x05scope\x129\n" +
|
||||
"\x16egress_gateway_address\x18\f \x01(\tH\x00R\x14egressGatewayAddress\x88\x01\x01\x1aD\n" +
|
||||
"\x16RuntimeAssetPathsEntry\x12\x10\n" +
|
||||
"\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" +
|
||||
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"\x19\n" +
|
||||
"\x05value\x18\x02 \x01(\tR\x05value:\x028\x01B\x19\n" +
|
||||
"\x17_egress_gateway_address\"\x19\n" +
|
||||
"\x17RestoreWorkloadResponse*a\n" +
|
||||
"\rSnapshotScope\x12\x1e\n" +
|
||||
"\x1aSNAPSHOT_SCOPE_UNSPECIFIED\x10\x00\x12\x17\n" +
|
||||
@@ -970,6 +1015,8 @@ func file_ateom_proto_init() {
|
||||
if File_ateom_proto != nil {
|
||||
return
|
||||
}
|
||||
file_ateom_proto_msgTypes[0].OneofWrappers = []any{}
|
||||
file_ateom_proto_msgTypes[8].OneofWrappers = []any{}
|
||||
type x struct{}
|
||||
out := protoimpl.TypeBuilder{
|
||||
File: protoimpl.DescBuilder{
|
||||
|
||||
@@ -51,6 +51,8 @@ message RunWorkloadRequest {
|
||||
string atespace = 1;
|
||||
string actor_name = 2;
|
||||
string actor_uid = 3;
|
||||
// Actor resource version observed by ate-api when assigning this worker.
|
||||
int64 actor_version = 9;
|
||||
|
||||
string actor_template_namespace = 4;
|
||||
string actor_template_name = 5;
|
||||
@@ -64,6 +66,10 @@ message RunWorkloadRequest {
|
||||
// to the local on-disk path atelet fetched it to (content-addressed, like
|
||||
// runsc_path). Empty for the gVisor runtime, which uses runsc_path.
|
||||
map<string, string> runtime_asset_paths = 8;
|
||||
|
||||
// Remote egress gateway selected for this activation. When absent, actor
|
||||
// traffic uses direct egress instead of being redirected through atunnel.
|
||||
optional string egress_gateway_address = 10;
|
||||
}
|
||||
|
||||
// WorkloadSpec parallels Pod, but with far fewer configurable fields.
|
||||
@@ -164,6 +170,8 @@ message RestoreWorkloadRequest {
|
||||
string atespace = 1;
|
||||
string actor_name = 2;
|
||||
string actor_uid = 3;
|
||||
// Actor resource version observed by ate-api when assigning this worker.
|
||||
int64 actor_version = 11;
|
||||
|
||||
string actor_template_namespace = 4;
|
||||
string actor_template_name = 5;
|
||||
@@ -181,6 +189,10 @@ message RestoreWorkloadRequest {
|
||||
|
||||
// What content to restore from the snapshot.
|
||||
SnapshotScope scope = 10;
|
||||
|
||||
// Remote egress gateway selected for this activation. When absent, actor
|
||||
// traffic uses direct egress instead of being redirected through atunnel.
|
||||
optional string egress_gateway_address = 12;
|
||||
}
|
||||
|
||||
message RestoreWorkloadResponse {
|
||||
|
||||
Reference in New Issue
Block a user