Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 16 additions & 8 deletions cmd/atenet/internal/router/egress/egress.go
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,7 @@ func (h *Handler) handleConnect(ctx context.Context, md *extproc.RequestMetadata
"egress denied: invalid actor certificate")
}

if err := h.validateActor(ctx, actorRef); err != nil {
if err := h.validateActor(ctx, actorRef, md.Host); err != nil {
return extproc.Result{}, err
}

Expand Down Expand Up @@ -255,9 +255,10 @@ func (h *Handler) lookupPolicy(ctx context.Context, leg string, ref resources.Ac
}

// validateActor checks the actor certified by the certificate against the control
// plane's current view of that actor: it still exists and it is running. Every
// error it returns is already a client-facing ext_proc denial.
func (h *Handler) validateActor(ctx context.Context, actorRef resources.ActorRef) error {
// plane's current view of that actor: it still exists and it is placed on a
// worker. A refusal is logged with destination. Every error it returns is
// already a client-facing ext_proc denial.
func (h *Handler) validateActor(ctx context.Context, actorRef resources.ActorRef, destination string) error {
// Confirm the certified actor still exists.
// TODO: this can cause heavy load on ate api server. Change it based on https://github.com/agent-substrate/substrate/issues/592.
actor, err := h.apiClient.GetActor(ctx, &ateapipb.GetActorRequest{
Expand All @@ -267,12 +268,19 @@ func (h *Handler) validateActor(ctx context.Context, actorRef resources.ActorRef
return mapEgressIdentityError(actorRef.Atespace, actorRef.Name, err)
}

// The actor performing egress must actually be running.
if actor.GetStatus().GetState() != ateapipb.ActorState_ACTOR_STATE_RUNNING {
// The actor performing egress must be placed on a worker. RESUMING is
// saved together with the worker assignment, so a booting workload can
// reach the network before it answers its wakeup probe. Every other state
// has left its worker or is leaving it.
switch state := actor.GetStatus().GetState(); state {
case ateapipb.ActorState_ACTOR_STATE_RUNNING, ateapipb.ActorState_ACTOR_STATE_RESUMING:
return nil
default:
slog.WarnContext(ctx, "egress denied: actor is not placed on a worker",
slog.Any("actor", actorRef), slog.String("state", state.String()), slog.String("destination", destination))
return extproc.NewReqError(envoy_type.StatusCode_Forbidden,
"egress denied: actor %q/%q is %s, not running", actorRef.Atespace, actorRef.Name, actor.GetStatus().GetState())
"egress denied: actor %q/%q is %s, not placed on a worker", actorRef.Atespace, actorRef.Name, state)
}
return nil
}

// authenticateActorCertificate turns the mTLS peer certificate Envoy recorded
Expand Down
33 changes: 32 additions & 1 deletion cmd/atenet/internal/router/egress/egress_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -689,7 +689,7 @@ func TestHandleRequestHeadersAuthorization(t *testing.T) {
want envoy_type.StatusCode
}{
{
name: "actor is not running",
name: "actor is suspended",
actor: &ateapipb.Actor{
Metadata: &ateapipb.ResourceMetadata{
Atespace: testEgressAtespace,
Expand Down Expand Up @@ -722,6 +722,37 @@ func TestHandleRequestHeadersAuthorization(t *testing.T) {
}
}

// Egress follows the actor's placement on a worker: a resuming actor is booting
// or restoring there and may fetch what it needs to become ready; every other
// state has left its worker or is leaving it.
func TestHandleRequestHeadersActorState(t *testing.T) {
ca := newTestCA(t, "actor-identity-ca")
allowed := map[ateapipb.ActorState]bool{
ateapipb.ActorState_ACTOR_STATE_RESUMING: true,
ateapipb.ActorState_ACTOR_STATE_RUNNING: true,
}
for value := range ateapipb.ActorState_name {
state := ateapipb.ActorState(value)
if state == ateapipb.ActorState_ACTOR_STATE_UNSPECIFIED {
continue
}
t.Run(state.String(), func(t *testing.T) {
actor := runningActor()
actor.Status.State = state
h := egressHandler(ca.roots(), actor, nil)
leaf := ca.issueActorCert(t, "spiffe://substrate-actor.local/ateom-for-actor/foo/bar", actorCertOptions{})
_, err := h.HandleRequestHeaders(context.Background(), egressMetadata(xfccHeader(leaf)))
if allowed[state] {
if err != nil {
t.Fatalf("HandleRequestHeaders() error = %v, want nil", err)
}
return
}
wantStatus(t, err, envoy_type.StatusCode_Forbidden)
})
}
}

// An ingress-only router has no actor-identity CA. If an egress CONNECT somehow
// reaches it, it must fail closed rather than tunnel unauthenticated traffic.
func TestHandleRequestHeadersWithoutConfiguredCA(t *testing.T) {
Expand Down
18 changes: 15 additions & 3 deletions cmd/ateom-gvisor/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -617,6 +617,11 @@ func (s *AteomService) RunWorkload(ctx context.Context, req *ateompb.RunWorkload
}
}
}()
// Egress before the first container starts, ingress after the wakeup
// probe (see ateomtunnel.Tunnel.ActivateEgress).
if err := s.tunnel.ActivateEgress(attribution, egress); err != nil {
return nil, err
}
// Create and start pause container. The bundle rootfs is composed here —
// an overlay of the node's cached image layers plus the bundle's private
// upper — because mounting is ateom's job (atelet runs with no
Expand Down Expand Up @@ -657,7 +662,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.tunnel.Activate(ateomstats.ActorAttributionFromRequest(req), s.sandboxDialer(req.GetActorUid()), egress); err != nil {
if err := s.tunnel.ActivateIngress(attribution, s.sandboxDialer(req.GetActorUid())); err != nil {
return nil, err
}

Expand Down Expand Up @@ -971,6 +976,13 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
}
}
}()
// As in RunWorkload: egress before the containers restore, ingress after
// the wakeup probe.
err = s.tunnel.ActivateEgress(attribution, egress)
timing.activate = lap(&tLast)
if err != nil {
return nil, err
}
checkpointDir := req.GetActorDirs().GetRestoreDir()

if hasDurableVolumes(containers) {
Expand Down Expand Up @@ -1065,8 +1077,8 @@ func (s *AteomService) RestoreWorkload(ctx context.Context, req *ateompb.Restore
if err != nil {
return nil, fmt.Errorf("while waiting for container wakeup probe: %w", err)
}
err = s.tunnel.Activate(attribution, s.sandboxDialer(req.GetActorUid()), egress)
timing.activate = lap(&tLast)
err = s.tunnel.ActivateIngress(attribution, s.sandboxDialer(req.GetActorUid()))
timing.activate += lap(&tLast)
if err != nil {
return nil, err
}
Expand Down
4 changes: 3 additions & 1 deletion cmd/ateom-gvisor/phaselog.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,9 @@ const containerCountKey = "ate.actor.container.count"

// restoreTiming is the elapsed time of each RestoreWorkload phase. A phase
// left at zero never ran. Under a Data scope pauseRestore and appRestore time
// the cold start that stands in for the restore.
// the cold start that stands in for the restore. activate sums the tunnel's
// egress half, armed before the containers restore, and its ingress half,
// after the wakeup probe.
type restoreTiming struct {
prep, egressPrepare, netSetup, durableDir time.Duration
pauseRootfs, pauseCreate, pauseRestore time.Duration
Expand Down
7 changes: 6 additions & 1 deletion cmd/ateom-microvm/restore.go
Original file line number Diff line number Diff line change
Expand Up @@ -321,6 +321,11 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
}
}
}()
// As in coldBootActor: egress before the VM restores, ingress after the
// wakeup probe.
if err := s.tunnel.ActivateEgress(p.attribution(), egress); err != nil {
return err
}
netDevs, err := ch.SnapshotNetDevices(restoreDir)
if err != nil {
return fmt.Errorf("while reading snapshot net devices: %w", err)
Expand Down Expand Up @@ -478,7 +483,7 @@ func (s *AteomService) restoreFullScope(ctx context.Context, p actorBootParams,
s.startActorLogForwarding(guestAC, attribution, c.GetName(), c.GetName())
}

if err := s.tunnel.Activate(p.attribution(), s.sandboxDialer(p.actorUID), egress); err != nil {
if err := s.tunnel.ActivateIngress(p.attribution(), s.sandboxDialer(p.actorUID)); err != nil {
return err
}
s.setRunningVM(actorUID, ra)
Expand Down
8 changes: 7 additions & 1 deletion cmd/ateom-microvm/run.go
Original file line number Diff line number Diff line change
Expand Up @@ -415,6 +415,12 @@ func (s *AteomService) coldBootActor(ctx context.Context, p actorBootParams) (re
}
}()

// Egress before the guest boots, ingress after the wakeup probe (see
// ateomtunnel.Tunnel.ActivateEgress).
if err := s.tunnel.ActivateEgress(p.attribution(), egress); err != nil {
return err
}

// Guest sizing + agent kernel params.
memMiB, vcpus, kparams := s.guestConfig()

Expand Down Expand Up @@ -589,7 +595,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.tunnel.Activate(p.attribution(), s.sandboxDialer(p.actorUID), egress); err != nil {
if err := s.tunnel.ActivateIngress(p.attribution(), s.sandboxDialer(p.actorUID)); err != nil {
return err
}
s.setRunningVM(actorUID, ra)
Expand Down
1 change: 1 addition & 0 deletions docs/api-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -322,6 +322,7 @@ How it behaves:
- **When it runs.** On every `ResumeActor` that actually wakes the actor. A `ResumeActor` on an actor that is already `RUNNING` is a no-op and does not probe.
- **Block-until-ready semantics.** `ResumeActor` returns successfully only after every container with a `wakeupProbe` has returned HTTP 200. If any container has not done so within `timeoutSeconds` (30s by default), `ResumeActor` fails and the actor moves to `ACTOR_STATE_CRASHED`, as it does for any other error while waking. A crashed actor cannot be resumed directly: `ResumeActor` rejects it with `FAILED_PRECONDITION`. Call `RevertActor` first, which returns it to `ACTOR_STATE_SUSPENDED` with its last snapshot, then resume it again.
- **Aggressive polling.** The poll loop is tuned for single-millisecond detection latency: a keep-alive HTTP client with a 1ms interval and 250ms per-request timeout. While the workload is still booting, kernel `RST`s return in microseconds, so the loop spends almost no time blocked; once the listener is up, the next attempt completes on veth-local latency.
- **Egress while waking.** The actor's egress is active before its containers start, so a workload can download what it needs to become ready (a model, a Git repository) before it answers the probe. Its EgressPolicy applies as it does once the actor is `RUNNING`; inbound traffic still waits for the probe.
- **Golden snapshot warm-up shortcut.** When **every** container in a template declares `wakeupProbe`, the actor template controller skips its default ~20s "give the workload time to settle" delay before taking the golden snapshot — `ResumeActor` already blocked until the workload reported 200, so the workload is known to be initialized. Templates that omit `wakeupProbe` on any container keep the 20s warm-up as a safety net.
- **Snapshot/restore interaction.** The TCP listener is part of the checkpointed RAM, so on resume `wakeupProbe` typically returns 200 on the first attempt, with no observable latency penalty.

Expand Down
4 changes: 2 additions & 2 deletions docs/network-egress.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,9 +22,9 @@ In the default Kubernetes deployment, the PEP obtains its certificate and privat

## Client Authorization

`atunnel` MUST present an actor-specific client certificate when connecting to the PEP. At actor activation, `atunnel` generates a private key and requests a short-lived certificate from substrate through the node-local atelet. The private key remains in `atunnel`. The certificate identifies the actor by its atespace, name, and UID and is scoped to the `atunnel` purpose.
`atunnel` MUST present an actor-specific client certificate when connecting to the PEP. Before the actor's workload starts, `atunnel` generates a private key and requests a short-lived certificate from substrate through the node-local atelet, and egress is active from then on: a workload may need the network to become ready, so egress does not wait for the wakeup probe the way ingress does. The private key remains in `atunnel`. The certificate identifies the actor by its atespace, name, and UID and is scoped to the `atunnel` purpose.

The PEP MUST require and verify the client certificate before accepting a CONNECT request. It MUST verify the certificate chain and lifetime against the Actor Identity CA, require the client authentication usage, and require exactly one valid `ActorIdentity` extension whose atespace, actor name, actor UID, and purpose are present. The current extension OID is `1.3.6.1.4.1.11129.2.12.2`; it is allocated under Google's enterprise number and will change after the CNCF donation is complete. The purpose MUST be `atunnel`, and the certificate's actor URI SAN MUST identify the same actor as the extension. The PEP MUST then confirm with ate-api-server that the actor still exists, that its current UID matches the certificate, and that it is running. The PEP MUST reject the CONNECT request if any of these checks fail.
The PEP MUST require and verify the client certificate before accepting a CONNECT request. It MUST verify the certificate chain and lifetime against the Actor Identity CA, require the client authentication usage, and require exactly one valid `ActorIdentity` extension whose atespace, actor name, actor UID, and purpose are present. The current extension OID is `1.3.6.1.4.1.11129.2.12.2`; it is allocated under Google's enterprise number and will change after the CNCF donation is complete. The purpose MUST be `atunnel`, and the certificate's actor URI SAN MUST identify the same actor as the extension. The PEP MUST then confirm with ate-api-server that the actor still exists, that its current UID matches the certificate, and that it is placed on a worker: `RESUMING` (booting or restoring there) or `RUNNING`. The PEP MUST reject the CONNECT request if any of these checks fail.

```text
CURRENT CONNECT EGRESS PATH (one tunnel per actor TCP connection except destination port 53)
Expand Down
19 changes: 13 additions & 6 deletions internal/ateomtunnel/ateomtunnel.go
Original file line number Diff line number Diff line change
Expand Up @@ -218,12 +218,10 @@ func (t *Tunnel) PrepareEgress(ctx context.Context, actor resources.ActorAttribu
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)
}
// ActivateEgress starts tunneling the actor's outbound traffic. Call it before
// the workload starts, since a workload may need the network to become ready.
// A nil egress is a no-op.
func (t *Tunnel) ActivateEgress(actor resources.ActorAttribution, egress *ActorEgress) error {
if egress == nil {
return nil
}
Expand All @@ -233,6 +231,15 @@ func (t *Tunnel) Activate(actor resources.ActorAttribution, dial atunnel.DialFun
return nil
}

// ActivateIngress starts admitting the actor's inbound traffic through dial.
// Call it once the workload is ready to serve.
func (t *Tunnel) ActivateIngress(actor resources.ActorAttribution, dial atunnel.DialFunc) 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)
}
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.
Expand Down
11 changes: 7 additions & 4 deletions internal/ateomtunnel/ateomtunnel_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -275,8 +275,11 @@ func TestActivateAndDeactivate(t *testing.T) {
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.ActivateEgress(actor, nil); err != nil {
t.Fatalf("ActivateEgress without a gateway: %v", err)
}
if err := tunnel.ActivateIngress(actor, dial); err != nil {
t.Fatalf("ActivateIngress: %v", err)
}
if err := tunnel.Deactivate(ctx, actor); err != nil {
t.Fatalf("Deactivate: %v", err)
Expand All @@ -286,8 +289,8 @@ func TestActivateAndDeactivate(t *testing.T) {
t.Fatalf("second Deactivate: %v", err)
}

err := tunnel.Activate(resources.ActorAttribution{UID: "uid"}, dial, nil)
err := tunnel.ActivateIngress(resources.ActorAttribution{UID: "uid"}, dial)
if err == nil || !strings.Contains(err.Error(), "while activating actor ingress") {
t.Fatalf("Activate with no actor reference = %v, want an ingress error", err)
t.Fatalf("ActivateIngress with no actor reference = %v, want an ingress error", err)
}
}
11 changes: 11 additions & 0 deletions internal/atunnel/egress.go
Original file line number Diff line number Diff line change
Expand Up @@ -274,10 +274,21 @@ func (e *Egress) handle(downstream net.Conn, active *egressActivation) {
_ = downstream.Close()
return
}
if active.dialer == nil {
e.mu.Unlock()
// From inside the sandbox this is a connection that opens and closes at
// once, indistinguishable from a network failure.
slog.WarnContext(active.ctx, "atunnel dropped an egress connection: actor egress is not active",
slog.String("peer", downstream.RemoteAddr().String()))
_ = downstream.Close()
return
}
if time.Now().Compare(active.expiresAt) >= 0 {
// Expiry blocks only new tunnels. Connections admitted with a valid
// certificate have completed mTLS and are allowed to drain normally.
e.mu.Unlock()
slog.WarnContext(active.ctx, "atunnel dropped an egress connection: actor certificate expired",
slog.String("peer", downstream.RemoteAddr().String()))
_ = downstream.Close()
return
}
Expand Down
6 changes: 6 additions & 0 deletions internal/e2e/atenet_dataplane.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ type AtenetDataplane interface {
IsEgressPolicyDenied(status int, body string) bool
PlatformMetricPrefixes([]string) []string
RouteDurationSeen(context.Context, string) (bool, error)
SupportsEgressWhileResuming() bool
}

// CurrentAtenetDataplane returns the implementation selected for this test
Expand Down Expand Up @@ -71,6 +72,8 @@ func (envoyAtenetDataplane) IsEgressPolicyDenied(status int, body string) bool {
(status == http.StatusBadGateway && strings.Contains(body, "request failed"))
}

func (envoyAtenetDataplane) SupportsEgressWhileResuming() bool { return true }

func (envoyAtenetDataplane) PlatformMetricPrefixes(prefixes []string) []string { return prefixes }

func (envoyAtenetDataplane) RouteDurationSeen(_ context.Context, collectorScrape string) (bool, error) {
Expand All @@ -93,6 +96,9 @@ func (agentGatewayAtenetDataplane) IsEgressPolicyDenied(status int, body string)
return status == http.StatusForbidden && strings.Contains(body, "actor egress policy denied")
}

// AgentGateway's substrate egress actor resolution admits RUNNING actors only.
func (agentGatewayAtenetDataplane) SupportsEgressWhileResuming() bool { return false }

func (agentGatewayAtenetDataplane) PlatformMetricPrefixes(prefixes []string) []string {
filtered := make([]string, 0, len(prefixes))
for _, prefix := range prefixes {
Expand Down
Loading
Loading