ateom: share atunnel wiring in ateomtunnel (#1959)

Fixes #687

Both ateoms build atunnel through the new `internal/ateomtunnel`.

- [x] Tests pass
- [x] Appropriate changes to documentation are included in the PR
This commit is contained in:
Eric Curtin
2026-10-01 05:33:56 +00:00
committed by GitHub
parent 7526b9cc75
commit 7317e083cf
13 changed files with 612 additions and 430 deletions
+2 -2
View File
@@ -81,9 +81,9 @@ func (s *AteomService) hostActor(ctx context.Context, attribution resources.Acto
session, err := ateomnet.ServeSandbox(ctx, ateomnet.SandboxNetworkConfig{
ActorUID: uid,
Veth: true,
EgressPort: s.atunnelEgressPort,
EgressPort: s.tunnel.EgressPort,
DNSPort: atunnel.DNSPort,
}, s.atunnelEgress, s.dnsRelay)
}, s.tunnel.Egress, s.tunnel.DNSRelay)
if err != nil {
s.actorsMu.Lock()
delete(s.actors, uid)
+29 -204
View File
@@ -22,7 +22,6 @@ import (
"fmt"
"log/slog"
"net"
"net/url"
"os"
"os/signal"
"slices"
@@ -41,11 +40,10 @@ import (
"github.com/agent-substrate/substrate/internal/ateomcgroup"
"github.com/agent-substrate/substrate/internal/ateomnet"
"github.com/agent-substrate/substrate/internal/ateomstats"
"github.com/agent-substrate/substrate/internal/atunnel"
"github.com/agent-substrate/substrate/internal/ateomtunnel"
"github.com/agent-substrate/substrate/internal/childreap"
"github.com/agent-substrate/substrate/internal/contextlogging"
"github.com/agent-substrate/substrate/internal/imagecache"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/nodepath"
"github.com/agent-substrate/substrate/internal/ocispec"
"github.com/agent-substrate/substrate/internal/otlprelay"
@@ -67,20 +65,10 @@ import (
var (
podUID = pflag.String("pod-uid", "", "The UID of the current pod")
// TODO(liorlieberman) have a sub package for all atunnel releated things like that
tunnelConfig = ateomtunnel.RegisterFlags(pflag.CommandLine)
// Every listen address here is an unspecified wildcard, which Go binds as a
// dual-stack socket.
atunnelListenAddress = pflag.String("atunnel-listen-address", ":443", "Address for actor ingress HTTPS")
atunnelConnectListenAddress = pflag.String("atunnel-connect-listen-address", ":8443", "Address for actor ingress mTLS CONNECT")
workerCredentialBundle = pflag.String("atunnel-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "Worker Pod credential bundle used by atunnel for inbound serving and outbound mTLS")
podIdentityTrustBundle = pflag.String("atunnel-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "Pod identity trust bundle used for router clients and the node-local atelet")
atunnelClientIdentity = pflag.String("atunnel-client-identity", installdefaults.RouterSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity allowed to call actor ingress HTTPS")
ateletIdentity = pflag.String("atunnel-broker-identity", installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity the node-local atelet must present on the credential broker connection. Override when atelet runs outside the default namespace.")
atunnelEgressListenAddress = pflag.String("atunnel-egress-listen-address", "0.0.0.0:15001", "Address for transparently intercepted actor egress TCP")
egressGatewayTrustBundle = pflag.String("atunnel-egress-trust-bundle", "/run/servicedns.podcert.ate.dev/trust-bundle.pem", "Service DNS trust bundle for the remote egress gateway")
readinessListenAddress = pflag.String("readiness-listen-address", "0.0.0.0:8080", "Address for HTTP readiness checks")
maxActors = pflag.Int("max-actors", 1000, "How many actors this worker will host at once")
readinessListenAddress = pflag.String("readiness-listen-address", "0.0.0.0:8080", "Address for HTTP readiness checks")
maxActors = pflag.Int("max-actors", 1000, "How many actors this worker will host at once")
showVersion = pflag.Bool("version", false, "Print version and exit.")
logLevelFlag = pflag.String("log-level", "info", "Minimum log level: debug, info, warn, or error.")
@@ -92,10 +80,6 @@ var (
reaper = childreap.New()
)
// actorHTTPUpstream is the in-sandbox HTTP endpoint atunnel proxies actor
// ingress to.
const actorHTTPUpstream = "http://" + ateomnet.ActorVethIP + ":80"
// workloadGracePeriod is the whole budget for draining the worker on shutdown.
// It needs to stay significantly less than the K8s termination grace period
// for the ateom, so the escalation to SIGKILL happens here rather than as a
@@ -225,29 +209,11 @@ func do(ctx context.Context) error {
}
actorLogger := actorlog.NewActorLogger(syncedWriter, metadata.OnGCE())
upstream, err := url.Parse(actorHTTPUpstream)
if err != nil {
return fmt.Errorf("while parsing atunnel upstream: %w", err)
}
// Use the pod's resolvers for cluster DNS access.
nameservers, err := atunnel.ResolvConfNameservers("/etc/resolv.conf")
if err != nil {
return fmt.Errorf("while reading the worker pod resolvers: %w", err)
}
dnsRelay, err := atunnel.NewDNSRelay(nameservers)
if err != nil {
return fmt.Errorf("while building the actor DNS relay: %w", err)
}
slog.InfoContext(ctx, "Actor DNS relay ready", slog.Any("upstreams", nameservers))
// Construct the service first so atunnel can use its namespace dialer.
ateomService := NewService(dnsRelay, actorLogger, *maxActors, *workerCredentialBundle, *podIdentityTrustBundle, *egressGatewayTrustBundle, *ateletIdentity)
atunnelIngress, atunnelEgress, atunnelEgressPort, err := runAtunnel(ctx, upstream)
tunnel, err := ateomtunnel.Start(ctx, *tunnelConfig, ateomnet.ActorHTTPUpstream)
if err != nil {
return err
}
ateomService.attachAtunnel(atunnelIngress, atunnelEgress, atunnelEgressPort)
ateomService := NewService(tunnel, actorLogger, *maxActors)
svr := grpc.NewServer(
grpc.StatsHandler(otelgrpc.NewServerHandler()),
@@ -281,9 +247,9 @@ func do(ctx context.Context) error {
go func() {
err := ateomcapacity.Report(ctx, ateomcapacity.ReportConfig{
SocketPath: nodepath.AteomSupportSocket,
CredentialBundlePath: *workerCredentialBundle,
TrustBundlePath: *podIdentityTrustBundle,
AteletSPIFFEID: *ateletIdentity,
CredentialBundlePath: tunnelConfig.CredentialBundle,
TrustBundlePath: tunnelConfig.TrustBundle,
AteletSPIFFEID: tunnelConfig.BrokerIdentity,
Actors: *maxActors,
})
if err != nil && ctx.Err() == nil {
@@ -300,49 +266,6 @@ func do(ctx context.Context) error {
return nil
}
func runAtunnel(ctx context.Context, upstream *url.URL) (*atunnel.Server, *atunnel.Egress, uint16, error) {
atunnelIngress, err := atunnel.NewServer(atunnel.Config{
CredentialBundlePath: *workerCredentialBundle,
TrustBundlePath: *podIdentityTrustBundle,
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 := atunnelIngress.Serve(ctx, atunnelListener); err != nil {
serverboot.Fatal(ctx, "Failed to serve actor ingress", err)
}
}()
slog.InfoContext(ctx, "atunnel serving", slog.String("address", *atunnelListenAddress))
atunnelConnectListener, err := net.Listen("tcp", *atunnelConnectListenAddress)
if err != nil {
return nil, nil, 0, fmt.Errorf("while opening atunnel CONNECT listener: %w", err)
}
go func() {
if err := atunnelIngress.ServeConnect(ctx, atunnelConnectListener); err != nil {
serverboot.Fatal(ctx, "Failed to serve actor CONNECT ingress", err)
}
}()
slog.InfoContext(ctx, "atunnel CONNECT serving", slog.String("address", *atunnelConnectListenAddress))
atunnelEgress, err := atunnel.NewEgress(atunnel.TCPOriginalDestination)
if err != nil {
return nil, nil, 0, fmt.Errorf("while configuring atunnel egress: %w", err)
}
// Bind egress only in sandbox namespaces.
atunnelEgressPort, err := atunnel.EgressPort(*atunnelEgressListenAddress)
if err != nil {
return nil, nil, 0, err
}
return atunnelIngress, atunnelEgress, atunnelEgressPort, nil
}
const (
rpcRunWorkload = "RunWorkload"
rpcRestoreWorkload = "RestoreWorkload"
@@ -374,27 +297,8 @@ type AteomService struct {
draining int
maxActors int
actorLogger *actorlog.ActorLogger
atunnelIngress *atunnel.Server
atunnelEgress *atunnel.Egress
// dnsRelay answers the sandbox's DNS from inside its own namespace.
dnsRelay *atunnel.DNSRelay
// atunnelEgressPort is the local atunnel listener used as the target of the
// actor network's transparent TCP redirect.
atunnelEgressPort uint16
// workerCredentialBundlePath contains the worker Pod certificate and key.
// Atunnel uses it for ingress serving and authentication to the atelet broker.
workerCredentialBundlePath string
// podIdentityTrustBundlePath verifies the node-local atelet's Pod identity.
podIdentityTrustBundlePath string
// egressGatewayTrustBundlePath verifies the remote gateway's serving cert.
egressGatewayTrustBundlePath string
// ateletSPIFFEID is the identity the node-local atelet must present on the
// credential broker connection. It names atelet's namespace, not this
// worker's, so it is configured rather than derived from the downward API.
ateletSPIFFEID string
actorLogger *actorlog.ActorLogger
tunnel *ateomtunnel.Tunnel
// shuttingDown is set once SIGTERM has been received. While true, new
// workload RPCs are rejected with codes.Unavailable.
@@ -415,19 +319,15 @@ type AteomService struct {
var _ ateompb.AteomServer = (*AteomService)(nil)
// NewService creates a new AteomService.
func NewService(dnsRelay *atunnel.DNSRelay, actorLogger *actorlog.ActorLogger, maxActors int, workerCredentialBundlePath, podIdentityTrustBundlePath, egressGatewayTrustBundlePath, ateletSPIFFEID string) *AteomService {
func NewService(tunnel *ateomtunnel.Tunnel, actorLogger *actorlog.ActorLogger, maxActors int) *AteomService {
return &AteomService{
locks: actorlock.New(),
inFlight: actorlock.NewInFlight(),
actors: map[string]*hostedActor{},
maxActors: maxActors,
dnsRelay: dnsRelay,
actorLogger: actorLogger,
workerCredentialBundlePath: workerCredentialBundlePath,
podIdentityTrustBundlePath: podIdentityTrustBundlePath,
egressGatewayTrustBundlePath: egressGatewayTrustBundlePath,
ateletSPIFFEID: ateletSPIFFEID,
cgroupRoot: defaultCgroupRoot,
locks: actorlock.New(),
inFlight: actorlock.NewInFlight(),
actors: map[string]*hostedActor{},
maxActors: maxActors,
tunnel: tunnel,
actorLogger: actorLogger,
cgroupRoot: defaultCgroupRoot,
}
}
@@ -629,7 +529,7 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
return nil, err
}
if err := s.deactivateActorNetworking(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
if err := s.tunnel.Deactivate(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
return nil, err
}
@@ -641,7 +541,7 @@ 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 the pause container.
egress, err := s.prepareActorEgress(ctx, req.GetAtespace(), req.GetActorName(), req.GetActorUid(), req.GetEgressGateway())
egress, err := s.tunnel.PrepareEgress(ctx, ateomstats.ActorAttributionFromRequest(req), req.GetEgressGateway())
if err != nil {
return nil, err
}
@@ -661,7 +561,7 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
if retErr != nil {
cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
if err := s.deactivateActorNetworking(cleanupCtx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
if err := s.tunnel.Deactivate(cleanupCtx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
slog.WarnContext(cleanupCtx, "Failed to deactivate actor networking after Run failure", slog.Any("err", err))
}
deleteContainers(cleanupCtx, rcmd, containersToDelete, "Run")
@@ -718,7 +618,7 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
if err := wakeupprobe.WaitAll(ctx, req.GetSpec().GetContainers(), ateomnet.ActorVethIP, wakeupprobe.DialFunc(s.sandboxDialer(req.GetActorUid()))); err != nil {
return nil, fmt.Errorf("while waiting for container wakeup probe: %w", err)
}
if err := s.activateActorNetworking(ateomstats.ActorAttributionFromRequest(req), egress); err != nil {
if err := s.tunnel.Activate(ateomstats.ActorAttributionFromRequest(req), s.sandboxDialer(req.GetActorUid()), egress); err != nil {
return nil, err
}
@@ -744,7 +644,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
// Not cancelable: a checkpoint is saving the actor's state.
defer s.inFlight.Add(req.GetActorUid(), rpcCheckpointWorkload, nil)()
if err := s.deactivateActorNetworking(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
if err := s.tunnel.Deactivate(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
return nil, err
}
@@ -920,7 +820,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
return nil, err
}
if err := s.deactivateActorNetworking(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
if err := s.tunnel.Deactivate(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
return nil, err
}
@@ -933,7 +833,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
// * All OCI bundles are set up, including for the pause container.
// * Checkpoint downloaded and placed on disk
egress, err := s.prepareActorEgress(ctx, req.GetAtespace(), req.GetActorName(), req.GetActorUid(), req.GetEgressGateway())
egress, err := s.tunnel.PrepareEgress(ctx, ateomstats.ActorAttributionFromRequest(req), req.GetEgressGateway())
if err != nil {
return nil, err
}
@@ -952,7 +852,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
if retErr != nil {
cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
if err := s.deactivateActorNetworking(cleanupCtx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
if err := s.tunnel.Deactivate(cleanupCtx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
slog.WarnContext(cleanupCtx, "Failed to deactivate actor networking after Restore failure", slog.Any("err", err))
}
deleteContainers(cleanupCtx, rcmd, containersToDelete, "Restore")
@@ -1040,7 +940,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
if err := wakeupprobe.WaitAll(ctx, req.GetSpec().GetContainers(), ateomnet.ActorVethIP, wakeupprobe.DialFunc(s.sandboxDialer(req.GetActorUid()))); err != nil {
return nil, fmt.Errorf("while waiting for container wakeup probe: %w", err)
}
if err := s.activateActorNetworking(ateomstats.ActorAttributionFromRequest(req), egress); err != nil {
if err := s.tunnel.Activate(ateomstats.ActorAttributionFromRequest(req), s.sandboxDialer(req.GetActorUid()), egress); err != nil {
return nil, err
}
@@ -1050,55 +950,6 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
return &ateompb.RestoreWorkloadResponse{}, nil
}
type actorEgress struct {
// client presents the actor certificate to the remote egress gateway.
client *atunnel.Client
// certificateSource owns the actor key and renews its certificate via atelet.
certificateSource *atunnel.BrokerCertificateSource
expiresAt time.Time
}
func (s *AteomService) prepareActorEgress(ctx context.Context, actorAtespace, actorName, actorUID string, gateway *ateompb.EgressGateway) (*actorEgress, error) {
if gateway == nil {
return nil, nil
}
if gateway.GetAddress() == "" {
return nil, fmt.Errorf("egress gateway address is required")
}
serverName, _, err := net.SplitHostPort(gateway.GetAddress())
if err != nil {
return nil, fmt.Errorf("invalid egress gateway address %q: %w", gateway.GetAddress(), err)
}
certificateSource, err := atunnel.NewBrokerCertificateSource(atunnel.BrokerConfig{
SocketPath: nodepath.AteomSupportSocket,
CredentialBundlePath: s.workerCredentialBundlePath,
TrustBundlePath: s.podIdentityTrustBundlePath,
ActorAtespace: actorAtespace,
ActorName: actorName,
ActorUID: actorUID,
AteletSPIFFEID: s.ateletSPIFFEID,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor certificate broker: %w", err)
}
// Mint before starting the workload so configured tunneled egress fails
// closed. The source retains the private key for mTLS and renewal.
expiresAt, err := certificateSource.MintAteomCertificate(ctx)
if err != nil {
return nil, fmt.Errorf("while obtaining actor certificate: %w", err)
}
gatewayClient, err := atunnel.NewClient(atunnel.ClientConfig{
GatewayAddress: gateway.GetAddress(),
ServerName: serverName,
GetClientCertificate: certificateSource.GetClientCertificate,
TrustBundlePath: s.egressGatewayTrustBundlePath,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor egress client: %w", err)
}
return &actorEgress{client: gatewayClient, certificateSource: certificateSource, expiresAt: expiresAt}, nil
}
func (s *AteomService) TerminateWorkload(ctx context.Context, req *ateompb.TerminateWorkloadRequest) (*ateompb.TerminateWorkloadResponse, error) {
if err := validateActorDirs(req.GetActorDirs()); err != nil {
return nil, err
@@ -1121,7 +972,7 @@ func (s *AteomService) TerminateWorkload(ctx context.Context, req *ateompb.Termi
func (s *AteomService) terminateWorkload(ctx context.Context, actorRef resources.ActorRef, actorUID, runscPath string, actorDirs *ateompb.ActorDirs, containers []*ateompb.Container) error {
var errs []error
if err := s.deactivateActorNetworking(ctx, resources.ActorAttribution{Ref: actorRef, UID: actorUID}); err != nil {
if err := s.tunnel.Deactivate(ctx, resources.ActorAttribution{Ref: actorRef, UID: actorUID}); err != nil {
errs = append(errs, fmt.Errorf("while deactivating actor networking: %w", err))
}
@@ -1166,19 +1017,6 @@ func (s *AteomService) terminateWorkload(ctx context.Context, actorRef resources
return errors.Join(errs...)
}
func (s *AteomService) activateActorNetworking(actor resources.ActorAttribution, egress *actorEgress) error {
if err := s.atunnelIngress.Activate(actor.Ref.Atespace, actor.Ref.Name, actor.UID, s.sandboxDialer(actor.UID)); err != nil {
return fmt.Errorf("while activating actor ingress: %w", err)
}
if egress == nil {
return nil
}
if err := s.atunnelEgress.Activate(actor.UID, egress.client, egress.certificateSource, egress.expiresAt); err != nil {
return fmt.Errorf("while activating actor egress: %w", err)
}
return nil
}
func deleteContainers(ctx context.Context, rcmd *runsc, containers []string, operation string) {
for _, container := range slices.Backward(containers) {
if err := rcmd.cmdDelete(ctx, container); err != nil {
@@ -1187,16 +1025,3 @@ func deleteContainers(ctx context.Context, rcmd *runsc, containers []string, ope
}
}
}
func (s *AteomService) deactivateActorNetworking(ctx context.Context, actor resources.ActorAttribution) error {
// Stop admitting traffic and drain active streams before the Actor network
// is torn down. Attempt both directions even if one fails to deactivate.
err := errors.Join(
s.atunnelIngress.Deactivate(ctx, actor.Ref.Atespace, actor.Ref.Name, actor.UID),
s.atunnelEgress.Deactivate(ctx, actor.UID),
)
if err != nil {
return fmt.Errorf("while deactivating actor networking: %w", err)
}
return nil
}
-8
View File
@@ -27,7 +27,6 @@ import (
"github.com/agent-substrate/substrate/internal/ateomnet"
"github.com/agent-substrate/substrate/internal/ateomnet/dns"
"github.com/agent-substrate/substrate/internal/atunnel"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
)
@@ -57,10 +56,3 @@ func removeActorResolvConf(ctx context.Context, path string) {
slog.WarnContext(ctx, "Failed to remove the actor resolv.conf", slog.Any("err", err))
}
}
// attachAtunnel completes setup after atunnel receives the service's dialer.
func (s *AteomService) attachAtunnel(ingress *atunnel.Server, egress *atunnel.Egress, egressPort uint16) {
s.atunnelIngress = ingress
s.atunnelEgress = egress
s.atunnelEgressPort = egressPort
}
+2 -2
View File
@@ -95,7 +95,7 @@ func (s *AteomService) CheckpointWorkload(ctx context.Context, req *ateompb.Chec
})
}()
if err := s.deactivateActorNetworking(ctx, attribution); err != nil {
if err := s.tunnel.Deactivate(ctx, attribution); err != nil {
return nil, err
}
@@ -412,7 +412,7 @@ func (s *AteomService) stopActorVM(ctx context.Context, actorUID string, actorDi
func (s *AteomService) terminateWorkload(ctx context.Context, actor resources.ActorAttribution, actorDirs *ateompb.ActorDirs) error {
var errs []error
if err := s.deactivateActorNetworking(ctx, actor); err != nil {
if err := s.tunnel.Deactivate(ctx, actor); err != nil {
errs = append(errs, fmt.Errorf("while deactivating actor networking: %w", err))
}
+2 -2
View File
@@ -87,9 +87,9 @@ func (s *AteomService) hostActor(ctx context.Context, attribution resources.Acto
// The tap and atunnel share a namespace; the guest owns the other end.
session, err := ateomnet.ServeSandbox(ctx, ateomnet.SandboxNetworkConfig{
ActorUID: uid,
EgressPort: s.atunnelEgressPort,
EgressPort: s.tunnel.EgressPort,
DNSPort: atunnel.DNSPort,
}, s.atunnelEgress, s.dnsRelay)
}, s.tunnel.Egress, s.tunnel.DNSRelay)
if err != nil {
s.actorsMu.Lock()
delete(s.actors, uid)
+24 -191
View File
@@ -24,19 +24,16 @@ package main
import (
"context"
"errors"
"flag"
"fmt"
"log/slog"
"net"
"net/url"
"os"
"os/signal"
"strings"
"sync"
"sync/atomic"
"syscall"
"time"
"cloud.google.com/go/compute/metadata"
"github.com/agent-substrate/substrate/cmd/ateom-microvm/internal/reaper"
@@ -45,8 +42,8 @@ import (
"github.com/agent-substrate/substrate/internal/ateinterceptors"
"github.com/agent-substrate/substrate/internal/ateomcapacity"
"github.com/agent-substrate/substrate/internal/ateomcgroup"
"github.com/agent-substrate/substrate/internal/atunnel"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/ateomnet"
"github.com/agent-substrate/substrate/internal/ateomtunnel"
"github.com/agent-substrate/substrate/internal/nodepath"
"github.com/agent-substrate/substrate/internal/otlprelay"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
@@ -73,22 +70,10 @@ var (
otlpRelaySocket = flag.String("otlp-relay-socket", nodepath.AteletOTLPSocketPath(),
"Unix socket of atelet's OTLP relay to export telemetry through, keeping it off the pod network. Empty, or absent at startup, exports directly to OTEL_EXPORTER_OTLP_ENDPOINT instead.")
// Every listen address here is an unspecified wildcard, which Go binds as a
// dual-stack socket.
atunnelListenAddress = flag.String("atunnel-listen-address", ":443", "Address for actor ingress HTTPS")
atunnelConnectListenAddress = flag.String("atunnel-connect-listen-address", ":8443", "Address for actor ingress mTLS CONNECT")
workerCredentialBundle = flag.String("atunnel-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "Worker Pod credential bundle used by atunnel for inbound serving and outbound mTLS")
podIdentityTrustBundle = flag.String("atunnel-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "Pod identity trust bundle used for router clients and the node-local atelet")
atunnelClientIdentity = flag.String("atunnel-client-identity", installdefaults.RouterSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity allowed to call actor ingress HTTPS")
ateletIdentity = flag.String("atunnel-broker-identity", installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity the node-local atelet must present on the credential broker connection. Override when atelet runs outside the default namespace.")
atunnelEgressListenAddress = flag.String("atunnel-egress-listen-address", "0.0.0.0:15001", "Address for transparently intercepted actor egress TCP")
egressGatewayTrustBundle = flag.String("atunnel-egress-trust-bundle", "/run/servicedns.podcert.ate.dev/trust-bundle.pem", "Service DNS trust bundle for the remote egress gateway")
readinessListenAddress = flag.String("readiness-listen-address", "0.0.0.0:8080", "Address for HTTP readiness checks")
maxActors = flag.Int("max-actors", 1000, "How many actors this worker will host at once")
)
tunnelConfig = ateomtunnel.RegisterFlags(flag.CommandLine)
const (
actorHTTPUpstream = "http://169.254.17.2:80"
readinessListenAddress = flag.String("readiness-listen-address", "0.0.0.0:8080", "Address for HTTP readiness checks")
maxActors = flag.Int("max-actors", 1000, "How many actors this worker will host at once")
)
func main() {
@@ -225,22 +210,6 @@ 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(actorHTTPUpstream)
if err != nil {
return fmt.Errorf("while parsing atunnel upstream: %w", err)
}
// The pod's own resolvers, so an actor resolves exactly what the worker
// resolves -- cluster DNS included.
nameservers, err := atunnel.ResolvConfNameservers("/etc/resolv.conf")
if err != nil {
return fmt.Errorf("while reading the worker pod resolvers: %w", err)
}
dnsRelay, err := atunnel.NewDNSRelay(nameservers)
if err != nil {
return fmt.Errorf("while building the actor DNS relay: %w", err)
}
slog.InfoContext(ctx, "Actor DNS relay ready", slog.Any("upstreams", nameservers))
// Give each actor's VMM and virtiofsd a cgroup of their own, so one busy
// guest cannot starve the rest.
actorCgroups, err := ateomcgroup.Delegate(ctx)
@@ -248,49 +217,12 @@ func do(ctx context.Context) error {
return fmt.Errorf("while delegating the worker cgroup: %w", err)
}
ateomService := NewService(*podUID, *chBinary, *kataDebug, *vmmMemReserve, *maxActors, dnsRelay, actorLogger, *workerCredentialBundle, *podIdentityTrustBundle, *egressGatewayTrustBundle, *ateletIdentity)
ateomService.actorCgroups = actorCgroups
atunnelIngress, err := atunnel.NewServer(atunnel.Config{
CredentialBundlePath: *workerCredentialBundle,
TrustBundlePath: *podIdentityTrustBundle,
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 := atunnelIngress.Serve(ctx, atunnelListener); err != nil {
serverboot.Fatal(ctx, "Failed to serve actor ingress", err)
}
}()
slog.InfoContext(ctx, "atunnel serving", slog.String("address", *atunnelListenAddress))
atunnelConnectListener, err := net.Listen("tcp", *atunnelConnectListenAddress)
if err != nil {
return fmt.Errorf("while opening atunnel CONNECT listener: %w", err)
}
go func() {
if err := atunnelIngress.ServeConnect(ctx, atunnelConnectListener); err != nil {
serverboot.Fatal(ctx, "Failed to serve actor CONNECT ingress", err)
}
}()
slog.InfoContext(ctx, "atunnel CONNECT serving", slog.String("address", *atunnelConnectListenAddress))
atunnelEgress, err := atunnel.NewEgress(atunnel.TCPOriginalDestination)
if err != nil {
return fmt.Errorf("while configuring atunnel egress: %w", err)
}
// Bind egress only in sandbox namespaces.
atunnelEgressPort, err := atunnel.EgressPort(*atunnelEgressListenAddress)
tunnel, err := ateomtunnel.Start(ctx, *tunnelConfig, ateomnet.ActorHTTPUpstream)
if err != nil {
return err
}
ateomService.attachAtunnel(atunnelIngress, atunnelEgress, atunnelEgressPort)
ateomService := NewService(*podUID, *chBinary, *kataDebug, *vmmMemReserve, *maxActors, tunnel, actorLogger)
ateomService.actorCgroups = actorCgroups
svr := grpc.NewServer(
grpc.StatsHandler(otelgrpc.NewServerHandler()),
@@ -328,9 +260,9 @@ func do(ctx context.Context) error {
go func() {
err := ateomcapacity.Report(ctx, ateomcapacity.ReportConfig{
SocketPath: nodepath.AteomSupportSocket,
CredentialBundlePath: *workerCredentialBundle,
TrustBundlePath: *podIdentityTrustBundle,
AteletSPIFFEID: *ateletIdentity,
CredentialBundlePath: tunnelConfig.CredentialBundle,
TrustBundlePath: tunnelConfig.TrustBundle,
AteletSPIFFEID: tunnelConfig.BrokerIdentity,
Actors: *maxActors,
})
if err != nil && ctx.Err() == nil {
@@ -407,30 +339,11 @@ type AteomService struct {
// with the guest RAM). Set from --vmm-mem-reserve-mib.
memReserveMiB int
// dnsRelay answers the actor's DNS from inside its own namespace.
dnsRelay *atunnel.DNSRelay
// 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
atunnelIngress *atunnel.Server
atunnelEgress *atunnel.Egress
// atunnelEgressPort is the local atunnel listener used as the target of the
// actor network's transparent TCP redirect.
atunnelEgressPort uint16
// workerCredentialBundlePath contains the worker Pod certificate and key.
// Atunnel uses it for ingress serving and authentication to the atelet broker.
workerCredentialBundlePath string
// podIdentityTrustBundlePath verifies the node-local atelet's Pod identity.
podIdentityTrustBundlePath string
// egressGatewayTrustBundlePath verifies the remote gateway's serving cert.
egressGatewayTrustBundlePath string
// ateletSPIFFEID is the identity the node-local atelet must present on the
// credential broker connection. It names atelet's namespace, not this
// worker's, so it is configured rather than derived from the downward API.
ateletSPIFFEID string
actorLogger *actorlog.ActorLogger
tunnel *ateomtunnel.Tunnel
// Guards actors, draining, and mutable hostedActor fields.
actorsMu sync.RWMutex
@@ -447,101 +360,21 @@ type AteomService struct {
var _ ateompb.AteomServer = (*AteomService)(nil)
// NewService creates a new AteomService.
func NewService(podUID, chBinary string, kataDebug bool, memReserveMiB, maxActors int, dnsRelay *atunnel.DNSRelay, actorLogger *actorlog.ActorLogger, workerCredentialBundlePath, podIdentityTrustBundlePath, egressGatewayTrustBundlePath, ateletSPIFFEID string) *AteomService {
func NewService(podUID, chBinary string, kataDebug bool, memReserveMiB, maxActors int, tunnel *ateomtunnel.Tunnel, actorLogger *actorlog.ActorLogger) *AteomService {
return &AteomService{
locks: actorlock.New(),
inFlight: actorlock.NewInFlight(),
actors: map[string]*hostedActor{},
maxActors: maxActors,
podUID: podUID,
chBinary: chBinary,
kataDebug: kataDebug,
memReserveMiB: memReserveMiB,
dnsRelay: dnsRelay,
actorLogger: actorLogger,
workerCredentialBundlePath: workerCredentialBundlePath,
podIdentityTrustBundlePath: podIdentityTrustBundlePath,
egressGatewayTrustBundlePath: egressGatewayTrustBundlePath,
ateletSPIFFEID: ateletSPIFFEID,
locks: actorlock.New(),
inFlight: actorlock.NewInFlight(),
actors: map[string]*hostedActor{},
maxActors: maxActors,
podUID: podUID,
chBinary: chBinary,
kataDebug: kataDebug,
memReserveMiB: memReserveMiB,
tunnel: tunnel,
actorLogger: actorLogger,
}
}
type actorEgress struct {
// client presents the actor certificate to the remote egress gateway.
client *atunnel.Client
// certificateSource owns the actor key and renews its certificate via atelet.
certificateSource *atunnel.BrokerCertificateSource
expiresAt time.Time
}
func (s *AteomService) prepareActorEgress(ctx context.Context, actorAtespace, actorName, actorUID string, gateway *ateompb.EgressGateway) (*actorEgress, error) {
if gateway == nil {
return nil, nil
}
if gateway.GetAddress() == "" {
return nil, fmt.Errorf("egress gateway address is required")
}
serverName, _, err := net.SplitHostPort(gateway.GetAddress())
if err != nil {
return nil, fmt.Errorf("invalid egress gateway address %q: %w", gateway.GetAddress(), err)
}
certificateSource, err := atunnel.NewBrokerCertificateSource(atunnel.BrokerConfig{
SocketPath: nodepath.AteomSupportSocket,
CredentialBundlePath: s.workerCredentialBundlePath,
TrustBundlePath: s.podIdentityTrustBundlePath,
ActorAtespace: actorAtespace,
ActorName: actorName,
ActorUID: actorUID,
AteletSPIFFEID: s.ateletSPIFFEID,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor certificate broker: %w", err)
}
// Mint before starting the workload so configured tunneled egress fails
// closed. The source retains the private key for mTLS and renewal.
expiresAt, err := certificateSource.MintAteomCertificate(ctx)
if err != nil {
return nil, fmt.Errorf("while obtaining actor certificate: %w", err)
}
gatewayClient, err := atunnel.NewClient(atunnel.ClientConfig{
GatewayAddress: gateway.GetAddress(),
ServerName: serverName,
GetClientCertificate: certificateSource.GetClientCertificate,
TrustBundlePath: s.egressGatewayTrustBundlePath,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor egress client: %w", err)
}
return &actorEgress{client: gatewayClient, certificateSource: certificateSource, expiresAt: expiresAt}, nil
}
func (s *AteomService) activateActorNetworking(actor resources.ActorAttribution, egress *actorEgress) error {
if err := s.atunnelIngress.Activate(actor.Ref.Atespace, actor.Ref.Name, actor.UID, s.sandboxDialer(actor.UID)); err != nil {
return fmt.Errorf("while activating actor ingress: %w", err)
}
if egress == nil {
return nil
}
if err := s.atunnelEgress.Activate(actor.UID, egress.client, egress.certificateSource, egress.expiresAt); err != nil {
return fmt.Errorf("while activating actor egress: %w", err)
}
return nil
}
func (s *AteomService) deactivateActorNetworking(ctx context.Context, actor resources.ActorAttribution) error {
// Stop admitting traffic and drain active streams before the Actor network
// is torn down. Attempt both directions even if one fails to deactivate.
err := errors.Join(
s.atunnelIngress.Deactivate(ctx, actor.Ref.Atespace, actor.Ref.Name, actor.UID),
s.atunnelEgress.Deactivate(ctx, actor.UID),
)
if err != nil {
return fmt.Errorf("while deactivating actor networking: %w", err)
}
return nil
}
// beginRPC registers before checking the drain flag so shutdown cannot miss it.
func (s *AteomService) beginRPC(actorUID, name string, cancel context.CancelFunc) (func(), error) {
release := s.inFlight.Add(actorUID, name, cancel)
+4 -5
View File
@@ -24,6 +24,7 @@ import (
"github.com/agent-substrate/substrate/internal/ateomnet"
"github.com/agent-substrate/substrate/internal/ateomnet/netns"
"github.com/agent-substrate/substrate/internal/ateomtunnel"
"github.com/agent-substrate/substrate/internal/atunnel"
"github.com/agent-substrate/substrate/internal/nodepath"
"github.com/agent-substrate/substrate/internal/resources"
@@ -44,11 +45,9 @@ func TestHostActorReplacesSameActor(t *testing.T) {
t.Fatal(err)
}
service := &AteomService{
atunnelEgress: egress,
atunnelEgressPort: 15001,
dnsRelay: dns,
actors: map[string]*hostedActor{},
maxActors: 1,
tunnel: &ateomtunnel.Tunnel{Egress: egress, EgressPort: 15001, DNSRelay: dns},
actors: map[string]*hostedActor{},
maxActors: 1,
}
const actorUID = "microvm-network-replace"
t.Cleanup(func() {
+4 -4
View File
@@ -122,7 +122,7 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
}
defer release()
if err := s.deactivateActorNetworking(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
if err := s.tunnel.Deactivate(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
return nil, err
}
@@ -222,7 +222,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
actorUID := p.actorUID
rr := s.resolveRuntime(p.assetPaths)
egress, err := s.prepareActorEgress(ctx, p.actorRef.Atespace, p.actorRef.Name, p.actorUID, p.egressGateway)
egress, err := s.tunnel.PrepareEgress(ctx, p.attribution(), p.egressGateway)
if err != nil {
return err
}
@@ -321,7 +321,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
if retErr != nil {
cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
if cleanupErr := s.deactivateActorNetworking(cleanupCtx, p.attribution()); cleanupErr != nil {
if cleanupErr := s.tunnel.Deactivate(cleanupCtx, p.attribution()); cleanupErr != nil {
slog.WarnContext(cleanupCtx, "Failed to deactivate actor networking after Restore failure", slog.Any("err", cleanupErr))
}
// Detach any bundle rootfs overlays mounted by buildActorContainers
@@ -502,7 +502,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
s.startActorLogForwarding(guestAC, attribution, c.GetName(), c.GetName())
}
if err := s.activateActorNetworking(p.attribution(), egress); err != nil {
if err := s.tunnel.Activate(p.attribution(), s.sandboxDialer(p.actorUID), egress); err != nil {
return err
}
s.setRunningVM(actorUID, ra)
+4 -4
View File
@@ -244,7 +244,7 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
}
defer release()
if err := s.deactivateActorNetworking(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
if err := s.tunnel.Deactivate(ctx, ateomstats.ActorAttributionFromRequest(req)); err != nil {
return nil, err
}
@@ -390,7 +390,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
return fmt.Errorf("ateom-microvm requires %q and %q asset paths", assetKernel, assetImage)
}
rr := s.resolveRuntime(paths)
egress, err := s.prepareActorEgress(ctx, p.actorRef.Atespace, p.actorRef.Name, p.actorUID, p.egressGateway)
egress, err := s.tunnel.PrepareEgress(ctx, p.attribution(), p.egressGateway)
if err != nil {
return err
}
@@ -402,7 +402,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
if retErr != nil {
cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), 30*time.Second)
defer cancel()
if cleanupErr := s.deactivateActorNetworking(cleanupCtx, p.attribution()); cleanupErr != nil {
if cleanupErr := s.tunnel.Deactivate(cleanupCtx, p.attribution()); cleanupErr != nil {
slog.WarnContext(cleanupCtx, "Failed to deactivate actor networking after Run failure", slog.Any("err", cleanupErr))
}
// Detach any bundle rootfs overlays mounted by buildActorContainers
@@ -587,7 +587,7 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
slog.Duration("since_boot", time.Since(tBooted)))
ra := &runningActor{chCmd: chCmd, vfsdCmd: vfsdCmd, apiSocket: apiSocket, baseID: actorUID, guestAgent: ac, workloadIDs: workloadIDs(ctrs)}
if err := s.activateActorNetworking(p.attribution(), egress); err != nil {
if err := s.tunnel.Activate(p.attribution(), s.sandboxDialer(p.actorUID), egress); err != nil {
return err
}
s.setRunningVM(actorUID, ra)
-8
View File
@@ -22,7 +22,6 @@ import (
"github.com/agent-substrate/substrate/internal/ateomnet"
"github.com/agent-substrate/substrate/internal/ateomnet/dns"
"github.com/agent-substrate/substrate/internal/atunnel"
)
// writeActorResolvConf points the guest resolver at its fixed gateway address.
@@ -33,10 +32,3 @@ func writeActorResolvConf(rootfs string) error {
}
return dns.WriteRootfsResolvConf(rootfs, dns.SandboxResolvConf(ateomnet.ActorVethGateway, pod))
}
// attachAtunnel completes setup after atunnel receives the service's dialer.
func (s *AteomService) attachAtunnel(ingress *atunnel.Server, egress *atunnel.Egress, egressPort uint16) {
s.atunnelIngress = ingress
s.atunnelEgress = egress
s.atunnelEgressPort = egressPort
}
+4
View File
@@ -30,6 +30,10 @@ const (
ActorVethGateway = "169.254.17.1"
ActorVethIP = "169.254.17.2"
// ActorHTTPUpstream is the in-sandbox HTTP endpoint atunnel proxies actor
// ingress to.
ActorHTTPUpstream = "http://" + ActorVethIP + ":80"
// hostVethLocalAddress is the gateway interface's IP address and prefix length.
hostVethLocalAddress = "169.254.17.1/30"
// actorVethLocalAddress is the actor interface's IP address and prefix length.
+244
View File
@@ -0,0 +1,244 @@
//go:build linux
// Copyright 2026 Google LLC
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
// Package ateomtunnel wires atunnel into an ateom: its flags, listeners, DNS
// relay, and per-actor ingress and egress. Both ateoms build their tunnel here
// so they cannot drift.
package ateomtunnel
import (
"context"
"errors"
"fmt"
"log/slog"
"net"
"net/url"
"time"
"github.com/agent-substrate/substrate/internal/atunnel"
"github.com/agent-substrate/substrate/internal/installdefaults"
"github.com/agent-substrate/substrate/internal/nodepath"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/agent-substrate/substrate/internal/serverboot"
)
// resolvConfPath holds the worker pod's resolvers. Tests override it.
var resolvConfPath = "/etc/resolv.conf"
// FlagSet is the registration surface that flag.FlagSet and pflag.FlagSet share.
type FlagSet interface {
StringVar(p *string, name, value, usage string)
}
// Config is the atunnel setup an ateom takes from its flags.
type Config struct {
// ListenAddress serves actor ingress HTTPS.
ListenAddress string
// ConnectListenAddress serves actor ingress mTLS CONNECT.
ConnectListenAddress string
// CredentialBundle holds the worker Pod certificate and key, used for
// inbound serving, outbound mTLS, and authentication to the atelet broker.
CredentialBundle string
// TrustBundle verifies router clients and the node-local atelet.
TrustBundle string
// ClientIdentity is the SPIFFE ID allowed to call actor ingress HTTPS.
ClientIdentity string
// BrokerIdentity is the SPIFFE ID the node-local atelet must present on the
// credential broker connection. It names atelet's namespace, not this
// worker's, so it is configured rather than derived from the downward API.
BrokerIdentity string
// EgressListenAddress receives transparently intercepted actor egress TCP.
EgressListenAddress string
// EgressTrustBundle verifies the remote egress gateway's serving cert.
EgressTrustBundle string
}
// RegisterFlags registers the atunnel flags on fs and returns the Config they
// fill once fs is parsed.
func RegisterFlags(fs FlagSet) *Config {
c := &Config{}
// Every listen address here is an unspecified wildcard, which Go binds as a
// dual-stack socket.
fs.StringVar(&c.ListenAddress, "atunnel-listen-address", ":443", "Address for actor ingress HTTPS")
fs.StringVar(&c.ConnectListenAddress, "atunnel-connect-listen-address", ":8443", "Address for actor ingress mTLS CONNECT")
fs.StringVar(&c.CredentialBundle, "atunnel-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "Worker Pod credential bundle used by atunnel for inbound serving and outbound mTLS")
fs.StringVar(&c.TrustBundle, "atunnel-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "Pod identity trust bundle used for router clients and the node-local atelet")
fs.StringVar(&c.ClientIdentity, "atunnel-client-identity", installdefaults.RouterSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity allowed to call actor ingress HTTPS")
fs.StringVar(&c.BrokerIdentity, "atunnel-broker-identity", installdefaults.AteletSPIFFEID(installdefaults.SystemNamespace), "SPIFFE identity the node-local atelet must present on the credential broker connection. Override when atelet runs outside the default namespace.")
fs.StringVar(&c.EgressListenAddress, "atunnel-egress-listen-address", "0.0.0.0:15001", "Address for transparently intercepted actor egress TCP")
fs.StringVar(&c.EgressTrustBundle, "atunnel-egress-trust-bundle", "/run/servicedns.podcert.ate.dev/trust-bundle.pem", "Service DNS trust bundle for the remote egress gateway")
return c
}
// Tunnel is the atunnel state one ateom shares across its actors.
type Tunnel struct {
// Ingress proxies router traffic to the active actors.
Ingress *atunnel.Server
// Egress tunnels the active actors' outbound TCP to the egress gateway.
Egress *atunnel.Egress
// EgressPort is the local atunnel listener used as the target of the
// actor network's transparent TCP redirect.
EgressPort uint16
// DNSRelay answers the sandbox's DNS from inside its own namespace.
DNSRelay *atunnel.DNSRelay
cfg Config
}
// Start builds the DNS relay and starts serving actor ingress, which proxies
// to upstream inside each sandbox. Egress binds only in sandbox namespaces, so
// nothing listens for it here.
func Start(ctx context.Context, cfg Config, upstream string) (*Tunnel, error) {
upstreamURL, err := url.Parse(upstream)
if err != nil {
return nil, fmt.Errorf("while parsing atunnel upstream: %w", err)
}
// The pod's own resolvers, so an actor resolves exactly what the worker
// resolves, cluster DNS included.
nameservers, err := atunnel.ResolvConfNameservers(resolvConfPath)
if err != nil {
return nil, fmt.Errorf("while reading the worker pod resolvers: %w", err)
}
dnsRelay, err := atunnel.NewDNSRelay(nameservers)
if err != nil {
return nil, fmt.Errorf("while building the actor DNS relay: %w", err)
}
slog.InfoContext(ctx, "Actor DNS relay ready", slog.Any("upstreams", nameservers))
ingress, err := atunnel.NewServer(atunnel.Config{
CredentialBundlePath: cfg.CredentialBundle,
TrustBundlePath: cfg.TrustBundle,
AllowedClientID: cfg.ClientIdentity,
Upstream: upstreamURL,
})
if err != nil {
return nil, fmt.Errorf("while configuring atunnel: %w", err)
}
if err := serve(ctx, "atunnel", cfg.ListenAddress, ingress.Serve); err != nil {
return nil, err
}
if err := serve(ctx, "atunnel CONNECT", cfg.ConnectListenAddress, ingress.ServeConnect); err != nil {
return nil, err
}
egress, err := atunnel.NewEgress(atunnel.TCPOriginalDestination)
if err != nil {
return nil, fmt.Errorf("while configuring atunnel egress: %w", err)
}
egressPort, err := atunnel.EgressPort(cfg.EgressListenAddress)
if err != nil {
return nil, err
}
return &Tunnel{Ingress: ingress, Egress: egress, EgressPort: egressPort, DNSRelay: dnsRelay, cfg: cfg}, nil
}
// serve listens on address and runs serveFn on it, exiting the process if it
// fails.
func serve(ctx context.Context, name, address string, serveFn func(context.Context, net.Listener) error) error {
lis, err := net.Listen("tcp", address)
if err != nil {
return fmt.Errorf("while opening %s listener: %w", name, err)
}
go func() {
if err := serveFn(ctx, lis); err != nil {
serverboot.Fatal(ctx, "Failed to serve "+name, err)
}
}()
slog.InfoContext(ctx, name+" serving", slog.String("address", address))
return nil
}
// ActorEgress is an actor's tunneled egress, ready to activate.
type ActorEgress struct {
// client presents the actor certificate to the remote egress gateway.
client *atunnel.Client
// certificateSource owns the actor key and renews its certificate via atelet.
certificateSource *atunnel.BrokerCertificateSource
expiresAt time.Time
}
// PrepareEgress mints the actor's certificate and builds its gateway client. A
// nil gateway means the actor has no tunneled egress and yields a nil result.
func (t *Tunnel) PrepareEgress(ctx context.Context, actor resources.ActorAttribution, gateway *ateompb.EgressGateway) (*ActorEgress, error) {
if gateway == nil {
return nil, nil
}
if gateway.GetAddress() == "" {
return nil, fmt.Errorf("egress gateway address is required")
}
serverName, _, err := net.SplitHostPort(gateway.GetAddress())
if err != nil {
return nil, fmt.Errorf("invalid egress gateway address %q: %w", gateway.GetAddress(), err)
}
certificateSource, err := atunnel.NewBrokerCertificateSource(atunnel.BrokerConfig{
SocketPath: nodepath.AteomSupportSocket,
CredentialBundlePath: t.cfg.CredentialBundle,
TrustBundlePath: t.cfg.TrustBundle,
ActorAtespace: actor.Ref.Atespace,
ActorName: actor.Ref.Name,
ActorUID: actor.UID,
AteletSPIFFEID: t.cfg.BrokerIdentity,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor certificate broker: %w", err)
}
// Mint before starting the workload so configured tunneled egress fails
// closed. The source retains the private key for mTLS and renewal.
expiresAt, err := certificateSource.MintAteomCertificate(ctx)
if err != nil {
return nil, fmt.Errorf("while obtaining actor certificate: %w", err)
}
gatewayClient, err := atunnel.NewClient(atunnel.ClientConfig{
GatewayAddress: gateway.GetAddress(),
ServerName: serverName,
GetClientCertificate: certificateSource.GetClientCertificate,
TrustBundlePath: t.cfg.EgressTrustBundle,
})
if err != nil {
return nil, fmt.Errorf("while configuring actor egress client: %w", err)
}
return &ActorEgress{client: gatewayClient, certificateSource: certificateSource, expiresAt: expiresAt}, nil
}
// Activate starts admitting the actor's traffic. Ingress reaches the actor
// through dial; egress is activated only when it was prepared.
func (t *Tunnel) Activate(actor resources.ActorAttribution, dial atunnel.DialFunc, egress *ActorEgress) error {
if err := t.Ingress.Activate(actor.Ref.Atespace, actor.Ref.Name, actor.UID, dial); err != nil {
return fmt.Errorf("while activating actor ingress: %w", err)
}
if egress == nil {
return nil
}
if err := t.Egress.Activate(actor.UID, egress.client, egress.certificateSource, egress.expiresAt); err != nil {
return fmt.Errorf("while activating actor egress: %w", err)
}
return nil
}
// Deactivate stops admitting the actor's traffic and drains its active streams,
// before the actor network is torn down. It attempts both directions even if
// one fails.
func (t *Tunnel) Deactivate(ctx context.Context, actor resources.ActorAttribution) error {
err := errors.Join(
t.Ingress.Deactivate(ctx, actor.Ref.Atespace, actor.Ref.Name, actor.UID),
t.Egress.Deactivate(ctx, actor.UID),
)
if err != nil {
return fmt.Errorf("while deactivating actor networking: %w", err)
}
return nil
}
+293
View File
@@ -0,0 +1,293 @@
//go:build linux
// Copyright 2026 Google LLC
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package ateomtunnel
import (
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/x509"
"crypto/x509/pkix"
"encoding/pem"
"flag"
"math/big"
"net"
"os"
"path/filepath"
"strings"
"testing"
"time"
"github.com/agent-substrate/substrate/internal/proto/ateompb"
"github.com/agent-substrate/substrate/internal/resources"
"github.com/spf13/pflag"
)
func TestRegisterFlags(t *testing.T) {
args := []string{
"--atunnel-listen-address=:1",
"--atunnel-connect-listen-address=:2",
"--atunnel-credential-bundle=cred",
"--atunnel-trust-bundle=trust",
"--atunnel-client-identity=client",
"--atunnel-broker-identity=broker",
"--atunnel-egress-listen-address=0.0.0.0:3",
"--atunnel-egress-trust-bundle=egress-trust",
}
want := Config{
ListenAddress: ":1",
ConnectListenAddress: ":2",
CredentialBundle: "cred",
TrustBundle: "trust",
ClientIdentity: "client",
BrokerIdentity: "broker",
EgressListenAddress: "0.0.0.0:3",
EgressTrustBundle: "egress-trust",
}
std := flag.NewFlagSet("std", flag.ContinueOnError)
stdCfg := RegisterFlags(std)
pf := pflag.NewFlagSet("pflag", pflag.ContinueOnError)
pfCfg := RegisterFlags(pf)
for name, tc := range map[string]struct {
cfg *Config
parse func([]string) error
}{
"flag": {stdCfg, std.Parse},
"pflag": {pfCfg, pf.Parse},
} {
t.Run(name, func(t *testing.T) {
if tc.cfg.ListenAddress != ":443" || tc.cfg.ConnectListenAddress != ":8443" || tc.cfg.EgressListenAddress != "0.0.0.0:15001" {
t.Errorf("unexpected listen defaults: %+v", *tc.cfg)
}
if tc.cfg.CredentialBundle == "" || tc.cfg.TrustBundle == "" || tc.cfg.ClientIdentity == "" || tc.cfg.BrokerIdentity == "" || tc.cfg.EgressTrustBundle == "" {
t.Errorf("unset default: %+v", *tc.cfg)
}
if err := tc.parse(args); err != nil {
t.Fatal(err)
}
if *tc.cfg != want {
t.Errorf("parsed config = %+v, want %+v", *tc.cfg, want)
}
})
}
}
// testConfig returns a Config over fresh credentials and free loopback ports.
func testConfig(t *testing.T) Config {
t.Helper()
dir := t.TempDir()
key, err := ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
t.Fatal(err)
}
now := time.Now()
template := &x509.Certificate{
SerialNumber: big.NewInt(1),
Subject: pkix.Name{CommonName: "test"},
NotBefore: now.Add(-time.Minute),
NotAfter: now.Add(time.Hour),
IsCA: true,
KeyUsage: x509.KeyUsageCertSign | x509.KeyUsageDigitalSignature,
ExtKeyUsage: []x509.ExtKeyUsage{x509.ExtKeyUsageClientAuth, x509.ExtKeyUsageServerAuth},
BasicConstraintsValid: true,
}
der, err := x509.CreateCertificate(rand.Reader, template, template, &key.PublicKey, key)
if err != nil {
t.Fatal(err)
}
keyDER, err := x509.MarshalPKCS8PrivateKey(key)
if err != nil {
t.Fatal(err)
}
certPEM := pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der})
keyPEM := pem.EncodeToMemory(&pem.Block{Type: "PRIVATE KEY", Bytes: keyDER})
bundle := filepath.Join(dir, "bundle.pem")
trust := filepath.Join(dir, "trust.pem")
if err := os.WriteFile(bundle, append(certPEM, keyPEM...), 0o600); err != nil {
t.Fatal(err)
}
if err := os.WriteFile(trust, certPEM, 0o600); err != nil {
t.Fatal(err)
}
resolvConf := filepath.Join(dir, "resolv.conf")
if err := os.WriteFile(resolvConf, []byte("nameserver 127.0.0.1\n"), 0o600); err != nil {
t.Fatal(err)
}
old := resolvConfPath
resolvConfPath = resolvConf
t.Cleanup(func() { resolvConfPath = old })
return Config{
ListenAddress: freeAddress(t),
ConnectListenAddress: freeAddress(t),
CredentialBundle: bundle,
TrustBundle: trust,
ClientIdentity: "spiffe://cluster.local/ns/ate-system/sa/atenet-router",
BrokerIdentity: "spiffe://cluster.local/ns/ate-system/sa/atelet",
EgressListenAddress: "0.0.0.0:15001",
EgressTrustBundle: trust,
}
}
// freeAddress returns a loopback address nothing is listening on.
func freeAddress(t *testing.T) string {
t.Helper()
lis, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
defer lis.Close()
return lis.Addr().String()
}
func startTunnel(t *testing.T, cfg Config) *Tunnel {
t.Helper()
ctx, cancel := context.WithCancel(context.Background())
t.Cleanup(cancel)
tunnel, err := Start(ctx, cfg, "http://169.254.17.2:80")
if err != nil {
t.Fatal(err)
}
return tunnel
}
func TestStart(t *testing.T) {
cfg := testConfig(t)
tunnel := startTunnel(t, cfg)
if tunnel.EgressPort != 15001 {
t.Errorf("EgressPort = %d, want 15001", tunnel.EgressPort)
}
if tunnel.Ingress == nil || tunnel.Egress == nil || tunnel.DNSRelay == nil {
t.Errorf("tunnel is missing a component: %+v", tunnel)
}
for _, address := range []string{cfg.ListenAddress, cfg.ConnectListenAddress} {
conn, err := net.DialTimeout("tcp", address, 5*time.Second)
if err != nil {
t.Fatalf("nothing is serving %s: %v", address, err)
}
conn.Close()
}
}
func TestStartErrors(t *testing.T) {
inUse, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatal(err)
}
defer inUse.Close()
for name, tc := range map[string]struct {
edit func(t *testing.T, cfg *Config)
wantErr string
}{
"no resolvers": {
edit: func(t *testing.T, _ *Config) {
empty := filepath.Join(t.TempDir(), "resolv.conf")
if err := os.WriteFile(empty, nil, 0o600); err != nil {
t.Fatal(err)
}
resolvConfPath = empty
},
wantErr: "worker pod resolvers",
},
"missing credentials": {
edit: func(t *testing.T, cfg *Config) { cfg.CredentialBundle = filepath.Join(t.TempDir(), "absent.pem") },
wantErr: "while configuring atunnel",
},
"ingress address in use": {
edit: func(_ *testing.T, cfg *Config) { cfg.ListenAddress = inUse.Addr().String() },
wantErr: "while opening atunnel listener",
},
"CONNECT address in use": {
edit: func(_ *testing.T, cfg *Config) { cfg.ConnectListenAddress = inUse.Addr().String() },
wantErr: "while opening atunnel CONNECT listener",
},
"egress address without a port": {
edit: func(_ *testing.T, cfg *Config) { cfg.EgressListenAddress = "0.0.0.0" },
wantErr: "egress listen address",
},
} {
t.Run(name, func(t *testing.T) {
cfg := testConfig(t)
tc.edit(t, &cfg)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
_, err := Start(ctx, cfg, "http://169.254.17.2:80")
if err == nil || !strings.Contains(err.Error(), tc.wantErr) {
t.Fatalf("Start error = %v, want one containing %q", err, tc.wantErr)
}
})
}
t.Run("bad upstream", func(t *testing.T) {
_, err := Start(context.Background(), testConfig(t), "http://[::1")
if err == nil || !strings.Contains(err.Error(), "atunnel upstream") {
t.Fatalf("Start error = %v, want an upstream error", err)
}
})
}
func TestPrepareEgress(t *testing.T) {
tunnel := &Tunnel{}
actor := resources.ActorAttribution{Ref: resources.ActorRef{Atespace: "space", Name: "actor"}, UID: "uid"}
got, err := tunnel.PrepareEgress(context.Background(), actor, nil)
if err != nil || got != nil {
t.Errorf("PrepareEgress(nil gateway) = %v, %v; want nil, nil", got, err)
}
for name, address := range map[string]string{
"empty address": "",
"missing a port": "gateway.example",
"unbalanced port": "[::1",
} {
t.Run(name, func(t *testing.T) {
got, err := tunnel.PrepareEgress(context.Background(), actor, &ateompb.EgressGateway{Address: address})
if err == nil || got != nil {
t.Errorf("PrepareEgress(%q) = %v, %v; want an error", address, got, err)
}
})
}
}
func TestActivateAndDeactivate(t *testing.T) {
tunnel := startTunnel(t, testConfig(t))
ctx := context.Background()
actor := resources.ActorAttribution{Ref: resources.ActorRef{Atespace: "space", Name: "actor"}, UID: "uid"}
dial := func(context.Context, string, string) (net.Conn, error) { return nil, net.ErrClosed }
if err := tunnel.Activate(actor, dial, nil); err != nil {
t.Fatalf("Activate without egress: %v", err)
}
if err := tunnel.Deactivate(ctx, actor); err != nil {
t.Fatalf("Deactivate: %v", err)
}
// Deactivating an actor that is not active is not an error.
if err := tunnel.Deactivate(ctx, actor); err != nil {
t.Fatalf("second Deactivate: %v", err)
}
err := tunnel.Activate(resources.ActorAttribution{UID: "uid"}, dial, nil)
if err == nil || !strings.Contains(err.Error(), "while activating actor ingress") {
t.Fatalf("Activate with no actor reference = %v, want an ingress error", err)
}
}