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
10 changes: 8 additions & 2 deletions components/egress/docs/opentelemetry.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ This page lists the OpenTelemetry metrics currently implemented in egress.
| `egress.dns.query.duration` | Histogram | `s` | Upstream DNS forward latency (recorded for allowed queries). |
| `egress.dns.query.failed_total` | Counter | - | Queries the proxy could not resolve, by `reason`. |
| `egress.policy.denied_total` | Counter | - | Number of DNS queries denied by policy. |
| `egress.nftables.rules.count` | Observable Gauge | `{element}` | Approximate policy size after last successful static apply. |
| `egress.nftables.rules.count` | Observable Gauge | `{element}` | Approximate policy size after last successful static apply (fleet profile: summed across every installed subject's policy, 0 while deny-first). |
| `egress.nftables.updates.count` | Counter | - | Number of successful nftables updates (static apply + dynamic IP add). |
| `egress.nftables.updates.failed_total` | Counter | - | nftables updates that failed, by `operation`. |
| `egress.system.memory.usage_bytes` | Observable Gauge | `By` | System memory used bytes (Linux: gopsutil; non-Linux build: `0`). |
Expand Down Expand Up @@ -65,11 +65,17 @@ queried name nor the error text is ever attached:
| `rcode` | The last resolver answered with a failover-worthy rcode, e.g. `SERVFAIL`. |

`egress.nftables.updates.failed_total` covers the other silent failure. Its `operation`
attribute is one of `static_apply`, `dynamic_add` or `remove`; `dynamic_add` is the one to
attribute is one of `static_apply`, `dynamic_add`, `remove`, or — in the fleet profile
(OSEP-0022) — `deny_first`, `dispatch_update`, `reset`; `dynamic_add` is the one to
alert on, because a failed add means the kernel never learned about IPs the policy allows,
so the chain drops traffic that should pass — which looks exactly like a policy denial from
inside the sandbox while `egress.policy.denied_total` stays flat.

The per-sandbox netns layer (fleet profile) counts its updates under the same operations;
two expected cases are deliberately NOT counted as failures: a sandbox-layer removal whose
netns is already destroyed (the rules died with it), and the startup recovery sweep of
netns that never had a table installed.

A `static_apply` failure happens during startup, where the sidecar logs and exits. Metrics
leave through a periodic reader and `os.Exit` skips the deferred shutdown, so that path
flushes telemetry explicitly before terminating — otherwise the one sample explaining why the
Expand Down
16 changes: 14 additions & 2 deletions components/egress/docs/policy-traffic-vault-flow.md
Original file line number Diff line number Diff line change
Expand Up @@ -62,8 +62,20 @@ server's reconciliation re-pushes policies.

## 3. Data plane: outbound traffic flow

Two paths per sandbox: DNS via the rewritten resolv.conf, and everything else
via the Pod netns `forward` hook. Dispatch is a verdict map keyed by
Two enforcement layers per sandbox: the authoritative Pod netns `forward`
hook (below), plus a per-sandbox netns OUTPUT chain mirroring the same policy
as defense in depth (`nsenter` from the host, table `opensandbox-fleet-ns`).
The sandbox layer allows loopback, DNS to the slot gateway only (dport 53,
gateway-scoped — the Pod layer enforces DNS policy via the proxy), and the
mirrored deny/dyn/allow verdicts; it catches traffic the forward hook never
sees (sandbox → host-local destinations take the INPUT path). DNS-learned
leases are refreshed in lockstep between both layers by the per-subject
connection refresh loop (Pod netns conntrack, bucketed by source IP, one
batched transaction per tick). Only TCP sessions are renewed — UDP/QUIC
(HTTP/3) relies on DNS lease TTLs; a sandbox-layer mirror miss marks the IPs
pending and redelivers them on the next tick.

Dispatch is a verdict map keyed by
`ip saddr . iifname` (the host veth binding is defense in depth against UDP
spoofing); the master chain defaults to **drop** so unregistered sources are
denied before their slot is even observed.
Expand Down
151 changes: 149 additions & 2 deletions components/egress/fleet.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ import (
"net/http"
"net/netip"
"os"
"path/filepath"
"strings"
"time"

"github.com/alibaba/opensandbox/egress/pkg/constants"
Expand All @@ -39,8 +41,10 @@ import (
"github.com/alibaba/opensandbox/egress/pkg/log"
"github.com/alibaba/opensandbox/egress/pkg/nftables"
"github.com/alibaba/opensandbox/egress/pkg/policy"
"github.com/alibaba/opensandbox/egress/pkg/sandboxnft"
"github.com/alibaba/opensandbox/egress/pkg/slotsource"
"github.com/alibaba/opensandbox/egress/pkg/subject"
"github.com/alibaba/opensandbox/egress/pkg/telemetry"
"github.com/alibaba/opensandbox/internal/safego"
)

Expand All @@ -49,6 +53,19 @@ import (
func runFleetProfile(ctx context.Context) {
log.Infof("egress profile: fleet (multi-sandbox control plane)")

otelShutdown, err := telemetry.Init(ctx)
if err != nil {
log.Warnf("OpenTelemetry metrics disabled (continuing without OTLP): %v", err)
otelShutdown = nil
}
if otelShutdown != nil {
defer func() {
shutdownCtx, shutdownCancel := context.WithTimeout(context.Background(), 5*time.Second)
defer shutdownCancel()
_ = otelShutdown(shutdownCtx)
}()
}

slotDir := envOrDefault(constants.EnvSlotStoreDir, constants.DefaultSlotStoreDir)
pollSec := constants.EnvIntOrDefault(constants.EnvSlotPollInterval, constants.DefaultSlotPollIntervalSeconds)
src := slotsource.NewFileSource(slotDir, time.Duration(pollSec)*time.Second)
Expand All @@ -59,13 +76,21 @@ func runFleetProfile(ctx context.Context) {
log.Fatalf("failed to load always allow/deny rule files: %v", err)
}

nftMgr := fleetnft.NewApplier(nil)
podNft := fleetnft.NewApplier(nil, fleetDoHOptions())
sandboxNft := sandboxnft.NewApplier(nil, sandboxDoHOptions())
nftMgr := &fleetEnforcer{pod: podNft, sandbox: sandboxNft}
// Recovery: wipe stale rules from a previous egress generation BEFORE
// rescanning, so no dead subject's policy survives into a new sandbox.
if err := nftMgr.ApplyReset(ctx); err != nil {
if err := podNft.ApplyReset(ctx); err != nil {
log.Fatalf("fleet nftables reset failed: %v", err)
}
log.Infof("fleet nftables table reset (stale rules cleared)")
// The sandbox layer needs the same wipe: sandbox netns can outlive the
// egress process, and their OUTPUT tables are the ONLY enforcement for
// host-local traffic (never seen by the Pod forward hook). Reset every
// netns the previous generation could have installed into — from the
// slot store and from the shared netns mount dir.
wipeSandboxTables(ctx, src, sandboxNft)

reg := subject.NewRegistry(alwaysDeny, alwaysAllow)
pendingTTL := time.Duration(constants.EnvIntOrDefault(constants.EnvPendingPushTTL, constants.DefaultPendingPushTTL)) * time.Second
Expand Down Expand Up @@ -124,6 +149,12 @@ func runFleetProfile(ctx context.Context) {
fleetSrv.StartPendingSweep(ctx)
controllerErr := controller.StartWatch(ctx, src)

// Per-subject connection refresh: active TCP connections keep their
// dynamic leases alive (bucketed by source IP from the Pod netns
// conntrack table); the sandbox-netns mirror is refreshed in lockstep.
nftMgr.StartConnectionRefresh(ctx)
log.Infof("fleet connection refresh started (bucketed per subject, every 30s)")

// Block until shutdown or a fatal control-plane failure (slot store
// unreadable = fail closed: the daemon must exit, not run unenforced).
select {
Expand All @@ -150,3 +181,119 @@ func runFleetProfile(ctx context.Context) {
log.Infof("fleet profile shutdown complete")
_ = os.Stderr.Sync()
}

// fleetDoHOptions parses the shared DoH-443 blocking env for the fleet
// profile: OPENSANDBOX_EGRESS_BLOCK_DOH_443 (strict all-443 drop when the
// blocklist is empty) + OPENSANDBOX_EGRESS_DOH_BLOCKLIST (comma-separated
// IP/CIDR list), same semantics as the sidecar profile.
func fleetDoHOptions() fleetnft.Options {
opts := fleetnft.Options{BlockDoH443: constants.IsTruthy(os.Getenv(constants.EnvBlockDoH443))}
if raw := strings.TrimSpace(os.Getenv(constants.EnvDoHBlocklist)); raw != "" {
opts.DoHBlocklistV4, opts.DoHBlocklistV6 = parseDoHBlocklist(raw)
}
return opts
}

// sandboxDoHOptions mirrors fleetDoHOptions for the per-sandbox netns layer,
// so both layers carry identical encrypted-DNS blocking.
func sandboxDoHOptions() sandboxnft.Options {
opts := sandboxnft.Options{BlockDoH443: constants.IsTruthy(os.Getenv(constants.EnvBlockDoH443))}
if raw := strings.TrimSpace(os.Getenv(constants.EnvDoHBlocklist)); raw != "" {
opts.DoHBlocklistV4, opts.DoHBlocklistV6 = parseDoHBlocklist(raw)
}
return opts
}

// fleetEnforcer composes the two enforcement layers per subject: the
// authoritative Pod-netns forward hook (pkg/fleetnft) plus the per-sandbox
// netns OUTPUT defense in depth (pkg/sandboxnft). It implements the
// fleetNftApplier surface, so the policy server and the DNS callback stay
// layer-agnostic. Pod first, sandbox second: the authoritative layer is
// always in place before the defense-in-depth layer, and a sandbox-layer
// failure fails the operation (the subject stays denying / on the old
// policy) instead of activating with a gap.
type fleetEnforcer struct {
pod *fleetnft.Applier
sandbox *sandboxnft.Applier
}

var _ fleetNftApplier = (*fleetEnforcer)(nil)

func (e *fleetEnforcer) ApplyDenyFirst(ctx context.Context, s subject.Subject, slot slotsource.Slot) error {
if err := e.pod.ApplyDenyFirst(ctx, s, slot); err != nil {
return err
}
return e.sandbox.ApplyDenyFirst(ctx, s, slot)
}

func (e *fleetEnforcer) ApplyPolicy(ctx context.Context, s subject.Subject, pol *policy.NetworkPolicy) error {
if err := e.pod.ApplyPolicy(ctx, s, pol); err != nil {
return err
}
return e.sandbox.ApplyPolicy(ctx, s, pol)
}

// ApplyDispatchUpdate is Pod-netns dispatch plus the sandbox-layer
// reconciliation: an unchanged-fencing slot update that moved the netns path
// or gateway must reinstall the subject's sandbox table (with its current
// policy) so the defense-in-depth layer stays aligned.
func (e *fleetEnforcer) ApplyDispatchUpdate(ctx context.Context, s subject.Subject, slot slotsource.Slot) error {
if err := e.pod.ApplyDispatchUpdate(ctx, s, slot); err != nil {
return err
}
return e.sandbox.ApplySlotUpdate(ctx, s, slot)
}

func (e *fleetEnforcer) Remove(ctx context.Context, s subject.Subject) error {
podErr := e.pod.Remove(ctx, s)
// Best effort: the sandbox rules die with the netns; a gone netns is
// expected and must never fail the unload.
_ = e.sandbox.Remove(ctx, s)
return podErr
}

// AddResolvedIPs mirrors DNS-learned leases into both layers.
func (e *fleetEnforcer) AddResolvedIPs(ctx context.Context, s subject.Subject, ips []nftables.ResolvedIP) error {
if err := e.pod.AddResolvedIPs(ctx, s, ips); err != nil {
return err
}
return e.sandbox.AddResolvedIPs(ctx, s, ips)
Comment thread
Pangjiping marked this conversation as resolved.
}

// StartConnectionRefresh launches the per-subject refresh loop; the sandbox
// layer is refreshed in lockstep through the mirror callback.
func (e *fleetEnforcer) StartConnectionRefresh(ctx context.Context) {
e.pod.StartConnectionRefresh(ctx, e.sandbox.AddResolvedIPs)
}

// wipeSandboxTables deletes the sandbox-layer table in every netns the
// previous egress generation could have installed into: the slot store is the
// authoritative list, and the shared netns mount dir covers slots whose files
// are already gone. Best effort — a missing netns or table is expected.
func wipeSandboxTables(ctx context.Context, src slotsource.Source, sandboxNft *sandboxnft.Applier) {
var paths []string
if slots, err := src.List(ctx); err == nil {
for _, slot := range slots {
paths = append(paths, slot.HostNetnsPath)
}
} else {
log.Warnf("slot store unreadable during recovery (slot-driven sandbox wipe skipped): %v", err)
}
paths = append(paths, netnsMountEntries()...)
sandboxNft.Reset(ctx, paths)
log.Infof("fleet sandbox tables reset (%d netns path(s))", len(paths))
}

// netnsMountEntries lists the shared netns mount dir (OSEP-0022 deployment
// precondition: /var/run/netns or equivalent).
func netnsMountEntries() []string {
entries, err := os.ReadDir(constants.DefaultNetnsMountDir)
if err != nil {
return nil
}
out := make([]string, 0, len(entries))
for _, e := range entries {
out = append(out, filepath.Join(constants.DefaultNetnsMountDir, e.Name()))
}
return out
}
56 changes: 32 additions & 24 deletions components/egress/nft.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,36 +70,44 @@ func setupNft(ctx context.Context, nftMgr nftApplier, initialPolicy *policy.Netw
nftMgr.StartConnectionRefresh(ctx)
}

// parseDoHBlocklist parses the comma-separated OPENSANDBOX_EGRESS_DOH_BLOCKLIST
// value (IP or CIDR entries) into v4/v6 lists. Invalid entries are logged and
// skipped. Shared by the sidecar and fleet profiles so both enforce the same
// DoH-443 semantics.
func parseDoHBlocklist(raw string) (v4, v6 []string) {
for _, p := range strings.Split(raw, ",") {
target := strings.TrimSpace(p)
if target == "" {
continue
}
if addr, err := netip.ParseAddr(target); err == nil {
if addr.Is4() {
v4 = append(v4, target)
} else if addr.Is6() {
v6 = append(v6, target)
}
continue
}
if prefix, err := netip.ParsePrefix(target); err == nil {
if prefix.Addr().Is4() {
v4 = append(v4, target)
} else if prefix.Addr().Is6() {
v6 = append(v6, target)
}
continue
}
log.Warnf("ignoring invalid DoH blocklist entry: %s", target)
}
return v4, v6
}

func parseNftOptions() nftables.Options {
opts := nftables.Options{BlockDoT: true}
if constants.IsTruthy(os.Getenv(constants.EnvBlockDoH443)) {
opts.BlockDoH443 = true
}
if raw := os.Getenv(constants.EnvDoHBlocklist); strings.TrimSpace(raw) != "" {
parts := strings.Split(raw, ",")
for _, p := range parts {
target := strings.TrimSpace(p)
if target == "" {
continue
}
if addr, err := netip.ParseAddr(target); err == nil {
if addr.Is4() {
opts.DoHBlocklistV4 = append(opts.DoHBlocklistV4, target)
} else if addr.Is6() {
opts.DoHBlocklistV6 = append(opts.DoHBlocklistV6, target)
}
continue
}
if prefix, err := netip.ParsePrefix(target); err == nil {
if prefix.Addr().Is4() {
opts.DoHBlocklistV4 = append(opts.DoHBlocklistV4, target)
} else if prefix.Addr().Is6() {
opts.DoHBlocklistV6 = append(opts.DoHBlocklistV6, target)
}
continue
}
log.Warnf("ignoring invalid DoH blocklist entry: %s", target)
}
opts.DoHBlocklistV4, opts.DoHBlocklistV6 = parseDoHBlocklist(raw)
}
return opts
}
4 changes: 4 additions & 0 deletions components/egress/pkg/constants/configuration.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,10 @@ const (
DefaultSlotStoreDir = "/run/fast-sandbox/network"
DefaultPendingPushTTL = 30
DefaultSlotPollIntervalSeconds = 1
// DefaultNetnsMountDir is where per-sandbox netns paths are mounted for
// host-domain consumers (egress runs nsenter --net=<path> against them);
// the deployment precondition of OSEP-0022.
DefaultNetnsMountDir = "/var/run/netns"
)

const (
Expand Down
Loading
Loading