mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
fixes
This commit is contained in:
@@ -7,7 +7,7 @@ Router has several responsibilities:
|
||||
Services in Kubernetes.
|
||||
With `--atenet-router=agentgateway`, the sidecar uses a static ConfigMap and
|
||||
atenet does not start an xDS server.
|
||||
* ext_proc server for the proxy. To make the deployment and debugging easier, we will run this component together
|
||||
* ext_proc server for the dataplane. To make the deployment and debugging easier, we will run this component together
|
||||
with the router, but this will be split later into its own component.
|
||||
* ext_proc will call into the ATE gRPC API to get the set of relevant backends (specific the worker IP) and
|
||||
route the traffic accordingly
|
||||
@@ -20,22 +20,22 @@ Router has several responsibilities:
|
||||
bounded wait elapses, instead of failing fast. See
|
||||
[docs/request-parking.md](../../../../../docs/request-parking.md).
|
||||
* Drains gracefully on SIGTERM: flips `/readyz` so the Service stops sending
|
||||
new connections, waits out endpoint propagation (`--drain-delay`), drives
|
||||
Envoy's admin API to drain established connections, gracefully stops the
|
||||
ext_proc server so parked requests finish normally (`--drain-timeout`,
|
||||
derived from the parking budget), then writes a drain-complete marker that
|
||||
releases the Envoy container's `preStop` hook. See `drain.go` and
|
||||
`envoydrain.go`.
|
||||
* Authenticates actor identity on egress: on every CONNECT, the egress Envoy's
|
||||
ext_proc handler re-verifies the actor's client certificate against the
|
||||
actor-identity CA, reads the `ActorIdentity` X.509 extension out of it, and
|
||||
checks the certified UID against the ATE API.
|
||||
new connections, waits out endpoint propagation (`--drain-delay`), drains the
|
||||
dataplane's established connections (Envoy only — driven over its admin API;
|
||||
agentgateway manages its own termination), gracefully stops the ext_proc
|
||||
server so parked requests finish normally (`--drain-timeout`, derived from
|
||||
the parking budget), then writes a drain-complete marker that releases the
|
||||
dataplane container's `preStop` hook. See `drain.go` and `envoydrain.go`.
|
||||
* Authenticates actor identity on egress: on every CONNECT, the egress
|
||||
gateway's ext_proc handler re-verifies the actor's client certificate against
|
||||
the actor-identity CA, reads the `ActorIdentity` X.509 extension out of it,
|
||||
and checks the certified UID against the ATE API.
|
||||
|
||||
## packages
|
||||
|
||||
The ext_proc server handles both traffic directions, and they apply opposite
|
||||
trust models — egress derives the actor identity from a client certificate
|
||||
Envoy verified against the actor-identity CA, ingress treats every request
|
||||
trust models — egress derives the actor identity from a client certificate the
|
||||
gateway verified against the actor-identity CA, ingress treats every request
|
||||
header as unauthenticated client input — so the two are kept in separate
|
||||
packages that cannot reach into each other:
|
||||
|
||||
@@ -48,8 +48,9 @@ packages that cannot reach into each other:
|
||||
* `egress` — certificate-based actor-identity authentication for outbound
|
||||
CONNECTs.
|
||||
|
||||
Direction is decided by the Envoy filter chain that accepted the request
|
||||
(`xds.filter_chain_name`), never by anything in the request itself, so a client
|
||||
Direction is decided by the filter chain the dataplane says accepted the
|
||||
request (`xds.filter_chain_name`, an Envoy attribute the egress gateway is
|
||||
configured to send), never by anything in the request itself, so a client
|
||||
cannot pick the egress path by crafting one. `router` itself does the wiring.
|
||||
|
||||
## modes
|
||||
@@ -67,9 +68,13 @@ than falling back to the other handler, which would run the request through the
|
||||
wrong trust model.
|
||||
|
||||
Ingress and egress are deployed separately today — `atenet-router` fronts the
|
||||
ingress Envoy, `atenet-egress` the egress one — because the two scale
|
||||
ingress dataplane, `atenet-egress` the egress gateway — because the two scale
|
||||
independently, not because they need separate binaries.
|
||||
|
||||
The `--atenet-router` choice only applies to the ingress dataplane. The egress
|
||||
gateway is its own Deployment with a statically configured Envoy, so
|
||||
`--atenet-router=agentgateway` leaves it untouched.
|
||||
|
||||
## status page
|
||||
|
||||
Serve a `/statusz` page on port 8080.
|
||||
|
||||
@@ -29,7 +29,7 @@ func NewRouterCmd() *cobra.Command {
|
||||
|
||||
cmd := &cobra.Command{
|
||||
Use: "router",
|
||||
Short: "Router components including xDS server and Envoy ExtProc gateway processing server",
|
||||
Short: "Router components including the Envoy xDS server and the ext_proc gateway processing server",
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
srv, err := NewRouterServer(cfg)
|
||||
if err != nil {
|
||||
@@ -41,7 +41,7 @@ func NewRouterCmd() *cobra.Command {
|
||||
},
|
||||
}
|
||||
|
||||
cmd.Flags().StringVar((*string)(&cfg.Mode), "mode", string(ModeAll), fmt.Sprintf("Traffic direction this instance serves: %q (also runs the xDS server and ActorTemplate controller for the ingress Envoy), %q (ext_proc only, needs no Kubernetes access), or %q for both. The ext_proc mux refuses a direction this instance was not started to serve rather than falling back to the other one", ModeIngress, ModeEgress, ModeAll))
|
||||
cmd.Flags().StringVar((*string)(&cfg.Mode), "mode", string(ModeAll), fmt.Sprintf("Traffic direction this instance serves: %q (also runs the ingress control plane — the xDS server and ActorTemplate controller — for an Envoy dataplane), %q (ext_proc only, needs no Kubernetes access), or %q for both. The ext_proc mux refuses a direction this instance was not started to serve rather than falling back to the other one", ModeIngress, ModeEgress, ModeAll))
|
||||
cmd.Flags().StringVar(&cfg.LogLevel, "log-level", "info", "Log level: debug, info, warn, error")
|
||||
cmd.Flags().StringVar(&cfg.MetricsAddr, "metrics-listen-addr", ":9090", "Address and port the prometheus metrics server should listen on.")
|
||||
cmd.Flags().BoolVar(&cfg.Standalone, "standalone", false, "Run in standalone mode, bypassing creation of managed deployment and services in Kubernetes cluster")
|
||||
@@ -51,13 +51,13 @@ func NewRouterCmd() *cobra.Command {
|
||||
cmd.Flags().StringVar(&cfg.AteapiAddr, "ateapi-address", "k8s:///api.ate-system.svc:443", "gRPC dial target for the cluster ateapi Control instance.")
|
||||
cmd.Flags().IntVar(&cfg.HttpPort, "port-http", 8080, "TCP port for workload traffic entering through the Envoy Router")
|
||||
cmd.Flags().IntVar(&cfg.XdsPort, "port-xds", 18000, "TCP port listening for the xDS dynamic Envoy connections")
|
||||
cmd.Flags().IntVar(&cfg.ExtprocPort, "port-extproc", 50051, "Listen port for the Envoy dynamic External Processing (ext_proc) server")
|
||||
cmd.Flags().StringVar(&cfg.ExtprocAddr, "extproc-address", "127.0.0.1", "Host IP or address of the Envoy External Processing (ext_proc) server")
|
||||
cmd.Flags().IntVar(&cfg.ExtprocPort, "port-extproc", 50051, "Listen port for the External Processing (ext_proc) server the dataplane calls")
|
||||
cmd.Flags().StringVar(&cfg.ExtprocAddr, "extproc-address", "127.0.0.1", "Host IP or address of the External Processing (ext_proc) server")
|
||||
cmd.Flags().StringVar(&cfg.EnvoyImage, "envoy-image", "envoyproxy/envoy:v1.30-latest", "Image URI used for dynamically launched router instances")
|
||||
cmd.Flags().StringVar(&cfg.TemplatesFile, "actor-templates-file", "", "Path to offline YAML configuration file listing ActorTemplates")
|
||||
cmd.Flags().IntVar(&cfg.StatusPort, "status-port", 4040, "Port to serve /statusz on (set <= 0 to disable serving status)")
|
||||
cmd.Flags().DurationVar(&cfg.HealthInterval, "health-interval", 1*time.Second, "Interval for checking health of dependent services")
|
||||
cmd.Flags().IntVar(&cfg.HttpsPort, "port-https", 8443, "TCP port for HTTPS workload traffic entering through the Envoy Router")
|
||||
cmd.Flags().IntVar(&cfg.HttpsPort, "port-https", 8443, "TCP port for HTTPS workload traffic entering through the router dataplane")
|
||||
cmd.Flags().StringVar(&cfg.EnvoyCertPath, "envoy-cert-path", "", "Path to the Envoy certificate file.")
|
||||
cmd.Flags().StringVar(&cfg.UpstreamCredentialBundlePath, "upstream-credential-bundle", "/run/podidentity.podcert.ate.dev/credential-bundle.pem", "PEM credential bundle (cert+key) the router presents as the client cert when dialing the actor's atunnel ingress server over mTLS. Empty disables upstream mTLS (legacy plaintext pod-IP:80).")
|
||||
cmd.Flags().StringVar(&cfg.UpstreamTrustBundlePath, "upstream-trust-bundle", "/run/podidentity.podcert.ate.dev/trust-bundle.pem", "PEM trust bundle used to validate the actor's atunnel ingress server certificate.")
|
||||
@@ -86,7 +86,7 @@ func NewRouterCmd() *cobra.Command {
|
||||
cmd.Flags().DurationVar(&cfg.DrainDelay, "drain-delay", 13*time.Second, "How long to keep serving after SIGTERM before starting the drain, covering readiness-probe detection and Service endpoint propagation")
|
||||
cmd.Flags().DurationVar(&cfg.DrainTimeout, "drain-timeout", 0, "Deadline for the ext_proc drain on shutdown; streams still open past it (parked requests included) are forcefully cancelled. 0 (the default) derives --parked-request-budget + the actor route timeout + margin so parked requests always finish normally. Explicit values must be >= --parked-request-budget")
|
||||
cmd.Flags().StringVar(&cfg.EnvoyAdminAddr, "envoy-admin-address", "127.0.0.1:9901", "Envoy admin interface the shutdown sequence drives to drain the sidecar (healthcheck/fail, drain_listeners, stats polling). Ignored with --atenet-router=agentgateway")
|
||||
cmd.Flags().StringVar(&cfg.DrainCompleteFile, "drain-complete-file", defaultDrainCompleteFile, "Marker file created (on a pod-shared emptyDir) once the shutdown drain completes; the Envoy container's preStop hook polls for it so Envoy exits as soon as — and no sooner than — the drain is done. Removed at startup to defuse stale markers. Empty disables the handshake")
|
||||
cmd.Flags().StringVar(&cfg.DrainCompleteFile, "drain-complete-file", defaultDrainCompleteFile, "Marker file created (on a pod-shared emptyDir) once the shutdown drain completes; the dataplane container's preStop hook polls for it so the proxy exits as soon as — and no sooner than — the drain is done. Removed at startup to defuse stale markers. Empty disables the handshake")
|
||||
|
||||
return cmd
|
||||
}
|
||||
|
||||
@@ -30,9 +30,9 @@ const (
|
||||
|
||||
// Mode selects which ext_proc directions an atenet instance serves. One binary
|
||||
// implements both, as two handlers behind the same ext_proc mux, but a
|
||||
// deployment usually fronts one Envoy and only needs the matching direction:
|
||||
// ingress and egress scale independently, so they run as separate Deployments
|
||||
// (atenet-router and atenet-egress).
|
||||
// deployment usually fronts one dataplane proxy and only needs the matching
|
||||
// direction: ingress and egress scale independently, so they run as separate
|
||||
// Deployments (atenet-router and atenet-egress).
|
||||
//
|
||||
// The mode is a floor on what an instance will answer, not just a hint. The mux
|
||||
// refuses a direction it has no handler for rather than falling back to the
|
||||
@@ -42,12 +42,15 @@ const (
|
||||
type Mode string
|
||||
|
||||
const (
|
||||
// ModeIngress serves actor-addressed traffic entering the mesh, and runs
|
||||
// the xDS server and ActorTemplate controller that configure its Envoy.
|
||||
// ModeIngress serves actor-addressed traffic arriving at the ingress
|
||||
// gateway, and runs the ingress control plane — the xDS server and
|
||||
// ActorTemplate controller — that configures its dataplane. Only the Envoy
|
||||
// dataplane takes configuration from it; agentgateway is statically
|
||||
// configured.
|
||||
ModeIngress Mode = "ingress"
|
||||
// ModeEgress serves actor CONNECTs leaving the mesh. Nothing else runs: the
|
||||
// egress Envoy is statically configured, so there is no xDS server, no
|
||||
// ActorTemplate controller, and no Kubernetes client.
|
||||
// ModeEgress serves actor CONNECTs leaving through the egress gateway.
|
||||
// Nothing else runs: the egress gateway is statically configured, so there
|
||||
// is no xDS server, no ActorTemplate controller, and no Kubernetes client.
|
||||
ModeEgress Mode = "egress"
|
||||
// ModeAll serves both directions from one instance. This is the default,
|
||||
// and what a single-gateway or local development setup wants.
|
||||
@@ -56,7 +59,7 @@ const (
|
||||
|
||||
// ServesIngress reports whether this mode answers ingress requests. It also
|
||||
// gates the ingress control plane: the xDS server that configures the ingress
|
||||
// Envoy and the ActorTemplate controller that feeds it.
|
||||
// dataplane (Envoy only) and the ActorTemplate controller that feeds it.
|
||||
func (m Mode) ServesIngress() bool { return m != ModeEgress }
|
||||
|
||||
// ServesEgress reports whether this mode answers egress CONNECTs.
|
||||
@@ -159,8 +162,8 @@ type routerConfig struct {
|
||||
EnvoyAdminAddr string
|
||||
|
||||
// DrainCompleteFile is the marker file written once shutdown completes,
|
||||
// releasing Envoy's preStop hook on the shared emptyDir. Removed at startup to
|
||||
// defuse stale markers. Empty disables the handshake.
|
||||
// releasing the dataplane container's preStop hook on the shared emptyDir.
|
||||
// Removed at startup to defuse stale markers. Empty disables the handshake.
|
||||
DrainCompleteFile string
|
||||
}
|
||||
|
||||
|
||||
@@ -237,7 +237,7 @@ func TestSetOtlpCollector(t *testing.T) {
|
||||
// No collector address may keep the router from starting. The address
|
||||
// defaults to OTEL_EXPORTER_OTLP_ENDPOINT, which also feeds the router's
|
||||
// own exporter and where https is perfectly valid; the router is the xDS
|
||||
// control plane for every Envoy in the mesh, so dropping Envoy's spans is
|
||||
// control plane for every ingress Envoy, so dropping Envoy's spans is
|
||||
// always the cheaper failure. setOtlpCollector returns nothing precisely so
|
||||
// this cannot regress into a startup error.
|
||||
tests := []struct {
|
||||
|
||||
@@ -12,7 +12,7 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
// Package egress implements the ext_proc handler for traffic leaving the mesh:
|
||||
// Package egress implements the ext_proc handler for outbound actor traffic:
|
||||
// it authenticates the actor behind an egress CONNECT before the gateway
|
||||
// tunnels it out.
|
||||
//
|
||||
@@ -88,18 +88,9 @@ func (h *Handler) Direction() extproc.Direction { return extproc.DirectionEgress
|
||||
// request metadata — contributes to the identity; the only inputs are the
|
||||
// certificate the actor-identity CA signed and the control plane's own view of
|
||||
// that actor.
|
||||
//
|
||||
// Authorization by destination and credential/token injection are deliberately
|
||||
// a TODO once we have SessionIdentity RPC service figured out.
|
||||
//
|
||||
// Egress never resumes an actor — it requires one already RUNNING — and picks
|
||||
// no upstream of its own, so the Result carries neither a resume outcome nor a
|
||||
// target.
|
||||
func (h *Handler) HandleRequestHeaders(ctx context.Context, md *extproc.RequestMetadata) (extproc.Result, error) {
|
||||
// Dispatch is by filter chain, so reaching here means the egress listener
|
||||
// accepted the request. That listener only routes CONNECT (its sole route
|
||||
// is a connect_matcher), so anything else is a config drift rather than a
|
||||
// client the gateway should tunnel for.
|
||||
// Sanity check that we were called on the Egress listener filter chain with
|
||||
// a CONNECT.
|
||||
if !strings.EqualFold(md.Method, "CONNECT") {
|
||||
return extproc.Result{}, extproc.NewReqError(envoy_type.StatusCode_MethodNotAllowed,
|
||||
"egress denied: expected CONNECT, got %q", md.Method)
|
||||
@@ -123,32 +114,62 @@ func (h *Handler) HandleRequestHeaders(ctx context.Context, md *extproc.RequestM
|
||||
"egress denied: invalid actor certificate")
|
||||
}
|
||||
|
||||
if err := validateIdentity(identity); err != nil {
|
||||
return extproc.Result{}, err
|
||||
}
|
||||
if err := h.validateActor(ctx, identity); err != nil {
|
||||
return extproc.Result{}, err
|
||||
}
|
||||
|
||||
slog.InfoContext(ctx, "egress identity authenticated",
|
||||
slog.String("atespace", identity.Atespace),
|
||||
slog.String("actor", identity.ActorName),
|
||||
slog.String("actorUid", identity.ActorUid),
|
||||
// For a CONNECT the :authority is the actor's original destination
|
||||
// (IP:port).
|
||||
slog.String("destination", md.Host))
|
||||
|
||||
// Identity is authenticated; let the CONNECT proceed unchanged.
|
||||
return extproc.Result{
|
||||
Response: &extprocv3.HeadersResponse{
|
||||
Response: &extprocv3.CommonResponse{},
|
||||
},
|
||||
}, nil
|
||||
}
|
||||
|
||||
// validateIdentity checks that the identity a verified actor certificate
|
||||
// carries names an actor that could exist at all, before it is used as a
|
||||
// control-plane lookup key.
|
||||
func validateIdentity(identity *substratex509.ActorIdentity) error {
|
||||
// The CA only ever mints these from control-plane state, so a name that is
|
||||
// not a legal resource name means the CA or its inputs are compromised.
|
||||
if !resources.IsValidResourceName(identity.Atespace) || !resources.IsValidResourceName(identity.ActorName) {
|
||||
return extproc.NewReqError(envoy_type.StatusCode_Forbidden,
|
||||
"egress denied: invalid actor identity %q/%q", identity.Atespace, identity.ActorName)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// validateActor checks the identity a certificate certifies against the control
|
||||
// plane's current view of that actor: it still exists, it is the actor the
|
||||
// certificate was issued to, and it is running. Every error it returns is
|
||||
// already a client-facing ext_proc denial.
|
||||
func (h *Handler) validateActor(ctx context.Context, identity *substratex509.ActorIdentity) error {
|
||||
atespace := identity.Atespace
|
||||
actorName := identity.ActorName
|
||||
actorUID := identity.ActorUid
|
||||
// For a CONNECT the :authority is the actor's original destination (IP:port).
|
||||
destination := md.Host
|
||||
|
||||
// The CA only ever mints these from control-plane state, so a name that is
|
||||
// not a legal resource name means the CA or its inputs are compromised.
|
||||
if !resources.IsValidResourceName(atespace) || !resources.IsValidResourceName(actorName) {
|
||||
return extproc.Result{}, extproc.NewReqError(envoy_type.StatusCode_Forbidden,
|
||||
"egress denied: invalid actor identity %q/%q", atespace, actorName)
|
||||
}
|
||||
|
||||
// Confirm the certified actor still exists. The name is only a lookup key
|
||||
// here; the UID below is what actually authorizes.
|
||||
// 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{
|
||||
Actor: &ateapipb.ObjectRef{Atespace: atespace, Name: actorName},
|
||||
})
|
||||
if err != nil {
|
||||
return extproc.Result{}, mapEgressIdentityError(atespace, actorName, err)
|
||||
return mapEgressIdentityError(atespace, actorName, err)
|
||||
}
|
||||
|
||||
// Authorize on the UID, not the name. Names are reused: delete an actor and
|
||||
// recreate it under the same atespace/name and it is a different actor with
|
||||
// a different UID. A certificate outliving its actor must not carry over to
|
||||
// the successor, so the UID the CA certified has to match the UID the
|
||||
// Authorize on the UID, not the name. The UID the CA certified has to match the UID the
|
||||
// control plane holds right now.
|
||||
if uid := actor.GetMetadata().GetUid(); uid != actorUID {
|
||||
slog.WarnContext(ctx, "egress denied: actor UID mismatch",
|
||||
@@ -156,31 +177,16 @@ func (h *Handler) HandleRequestHeaders(ctx context.Context, md *extproc.RequestM
|
||||
slog.String("actor", actorName),
|
||||
slog.String("certificateActorUid", actorUID),
|
||||
slog.String("currentActorUid", uid))
|
||||
return extproc.Result{}, extproc.NewReqError(envoy_type.StatusCode_Forbidden,
|
||||
return extproc.NewReqError(envoy_type.StatusCode_Forbidden,
|
||||
"egress denied: actor %q/%q is not the actor this certificate was issued to", atespace, actorName)
|
||||
}
|
||||
|
||||
// The actor performing egress must actually be running.
|
||||
if actor.GetStatus() != ateapipb.Actor_STATUS_RUNNING {
|
||||
return extproc.Result{}, extproc.NewReqError(envoy_type.StatusCode_Forbidden,
|
||||
return extproc.NewReqError(envoy_type.StatusCode_Forbidden,
|
||||
"egress denied: actor %q/%q is %s, not running", atespace, actorName, actor.GetStatus())
|
||||
}
|
||||
|
||||
slog.InfoContext(ctx, "egress identity authenticated",
|
||||
slog.String("atespace", atespace),
|
||||
slog.String("actor", actorName),
|
||||
slog.String("actorUid", actorUID),
|
||||
slog.String("destination", destination),
|
||||
slog.String("status", actor.GetStatus().String()))
|
||||
|
||||
// Identity is authenticated; let the CONNECT proceed unchanged. Milestone 2
|
||||
// would additionally authorize `destination` and inject upstream credentials
|
||||
// here by returning a HeaderMutation.
|
||||
return extproc.Result{
|
||||
Response: &extprocv3.HeadersResponse{
|
||||
Response: &extprocv3.CommonResponse{},
|
||||
},
|
||||
}, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// authenticateActorCertificate turns the mTLS peer certificate Envoy recorded
|
||||
|
||||
@@ -18,15 +18,15 @@ import (
|
||||
extprocv3 "github.com/envoyproxy/go-control-plane/envoy/service/ext_proc/v3"
|
||||
)
|
||||
|
||||
// Direction is the side of the mesh a request arrived on. It selects the
|
||||
// handler, so it must come from something Envoy asserts rather than from the
|
||||
// Direction is the gateway a request arrived through. It selects the handler,
|
||||
// so it must come from something the dataplane asserts rather than from the
|
||||
// request.
|
||||
type Direction string
|
||||
|
||||
const (
|
||||
// DirectionIngress is traffic entering the mesh, addressed to an actor.
|
||||
// DirectionIngress is inbound traffic addressed to an actor.
|
||||
DirectionIngress Direction = "ingress"
|
||||
// DirectionEgress is traffic leaving the mesh, tunneled out of an actor.
|
||||
// DirectionEgress is outbound traffic tunneled out of an actor.
|
||||
DirectionEgress Direction = "egress"
|
||||
)
|
||||
|
||||
|
||||
@@ -54,8 +54,8 @@ func WrapReqError(code envoy_type.StatusCode, cause error, format string, args .
|
||||
}
|
||||
}
|
||||
|
||||
// ImmediateResponse tells Envoy to answer the request itself, without going
|
||||
// upstream.
|
||||
// ImmediateResponse tells the dataplane to answer the request itself, without
|
||||
// going upstream.
|
||||
func ImmediateResponse(statusCode envoy_type.StatusCode, message string) *extprocv3.ProcessingResponse {
|
||||
return &extprocv3.ProcessingResponse{
|
||||
Response: &extprocv3.ProcessingResponse_ImmediateResponse{
|
||||
|
||||
@@ -12,8 +12,8 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
// Package extproc implements the Envoy external processing (ext_proc) gRPC
|
||||
// server that the atenet router serves to its Envoy gateways.
|
||||
// Package extproc implements the external processing (ext_proc) gRPC server
|
||||
// that the atenet router serves to its dataplane gateways.
|
||||
//
|
||||
// This package is the multiplexer and nothing else: it terminates the
|
||||
// ext_proc stream, works out which direction — ingress or egress — a request
|
||||
@@ -42,8 +42,8 @@ import (
|
||||
// direction absent from the map is refused at the mux — see Server.Process.
|
||||
type Handlers map[Direction]Handler
|
||||
|
||||
// Server implements the Envoy external processing gRPC server, dispatching
|
||||
// each request to the Handler for the direction it arrived on.
|
||||
// Server implements the external processing gRPC server, dispatching each
|
||||
// request to the Handler for the direction it arrived on.
|
||||
type Server struct {
|
||||
port int
|
||||
handlers Handlers
|
||||
@@ -132,12 +132,13 @@ func (s *Server) processRequestHeaders(
|
||||
|
||||
// One atenet binary serves both directions, as two ext_proc handlers
|
||||
// selected here. They are deployed separately today — atenet-router fronts
|
||||
// the ingress Envoy, atenet-egress the egress one — because the two scale
|
||||
// independently, and --mode restricts an instance to the direction its
|
||||
// deployment fronts. Nothing stops a single instance from serving both.
|
||||
// the ingress dataplane, atenet-egress the egress gateway — because the two
|
||||
// scale independently, and --mode restricts an instance to the direction
|
||||
// its deployment fronts. Nothing stops a single instance from serving both.
|
||||
//
|
||||
// Which handler runs is decided by the Envoy filter chain that accepted the
|
||||
// request, never by anything in the request itself (see directionOf).
|
||||
// Which handler runs is decided by the filter chain the dataplane says
|
||||
// accepted the request, never by anything in the request itself (see
|
||||
// directionOf).
|
||||
dir := directionOf(req)
|
||||
|
||||
var res Result
|
||||
@@ -145,10 +146,10 @@ func (s *Server) processRequestHeaders(
|
||||
if handler, ok := s.handlers[dir]; ok {
|
||||
res, err = handler.HandleRequestHeaders(ctx, md)
|
||||
} else {
|
||||
// The Envoy in front of this instance is sending traffic the instance
|
||||
// was not started to serve. Refuse it rather than falling back to the
|
||||
// other direction's handler: the two apply opposite trust models, so a
|
||||
// fallback would run a request through the wrong one.
|
||||
// The dataplane in front of this instance is sending traffic the
|
||||
// instance was not started to serve. Refuse it rather than falling back
|
||||
// to the other direction's handler: the two apply opposite trust
|
||||
// models, so a fallback would run a request through the wrong one.
|
||||
err = NewReqError(envoy_type.StatusCode_NotFound,
|
||||
"this router does not serve %s traffic", dir)
|
||||
}
|
||||
|
||||
@@ -30,8 +30,8 @@ type Handler interface {
|
||||
// its dispatch table by it.
|
||||
Direction() Direction
|
||||
|
||||
// HandleRequestHeaders decides what Envoy should do with a request whose
|
||||
// headers have just arrived.
|
||||
// HandleRequestHeaders decides what the dataplane should do with a request
|
||||
// whose headers have just arrived.
|
||||
//
|
||||
// A returned error denies the request: a *ReqError carries the status code
|
||||
// and client-safe body to answer with, anything else becomes a 500. The
|
||||
@@ -43,7 +43,7 @@ type Handler interface {
|
||||
|
||||
// Result is what a handler tells the mux about a request it allowed.
|
||||
type Result struct {
|
||||
// Response is the header mutation Envoy applies before the request
|
||||
// Response is the header mutation the dataplane applies before the request
|
||||
// continues. Handlers that only authenticate return an empty CommonResponse.
|
||||
Response *extprocv3.HeadersResponse
|
||||
|
||||
|
||||
@@ -26,8 +26,8 @@ import (
|
||||
const AuthorityHeader = ":authority"
|
||||
|
||||
// RequestMetadata is the request the mux hands to a direction handler: the
|
||||
// HTTP headers Envoy sent, flattened and lowercased, with the pseudo-headers
|
||||
// every handler needs pulled out.
|
||||
// HTTP headers the dataplane sent, flattened and lowercased, with the
|
||||
// pseudo-headers every handler needs pulled out.
|
||||
type RequestMetadata struct {
|
||||
// Headers holds every header, keyed by lowercased name.
|
||||
Headers map[string]string
|
||||
|
||||
@@ -36,7 +36,7 @@ import (
|
||||
const ServiceName = "atenet-router"
|
||||
|
||||
// atenet.router.route.duration measures the latency from when the ext_proc handler receives a request
|
||||
// (Envoy -> EPP) until the target worker endpoint is resolved
|
||||
// (dataplane -> EPP) until the target worker endpoint is resolved
|
||||
const routeDurationMetricName = "atenet.router.route.duration"
|
||||
|
||||
// NewRouteDurationHistogram creates the atenet.router.route.duration histogram from
|
||||
|
||||
@@ -158,7 +158,7 @@ func updateComponentHealth(health *ComponentHealth, healthy bool, msg string, ch
|
||||
|
||||
func (rh *routerHealth) checkDataplane(ctx context.Context) (bool, string) {
|
||||
// The dataplane this polls is the *ingress* proxy sharing the router's pod.
|
||||
// The egress Envoy is a separate, statically configured proxy on its own
|
||||
// The egress gateway is a separate, statically configured proxy on its own
|
||||
// admin port; an egress-only router has none beside it, and probing this
|
||||
// address would report a permanently unhealthy dependency.
|
||||
if !rh.cfg.Mode.ServesIngress() {
|
||||
|
||||
@@ -12,10 +12,10 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
// Package ingress implements the ext_proc handler for traffic entering the
|
||||
// mesh: it resolves the actor a request is addressed to, resumes it through
|
||||
// the control plane (parking the request while the worker pool is saturated),
|
||||
// and points Envoy at the worker that ends up hosting it.
|
||||
// Package ingress implements the ext_proc handler for traffic arriving at the
|
||||
// ingress gateway: it resolves the actor a request is addressed to, resumes it
|
||||
// through the control plane (parking the request while the worker pool is
|
||||
// saturated), and points the dataplane at the worker that ends up hosting it.
|
||||
//
|
||||
// Everything reaching this handler is unauthenticated client input. The
|
||||
// opposite trust model — an actor identity carried by a CA-signed client
|
||||
@@ -72,10 +72,10 @@ func (h *Handler) ParkingStatus() ParkingStatus { return h.parking.status() }
|
||||
func (h *Handler) HandleRequestHeaders(ctx context.Context, md *extproc.RequestMetadata) (extproc.Result, error) {
|
||||
slog.InfoContext(ctx, "Request", slog.String("host", md.Host))
|
||||
|
||||
// Envoy doesn't propagate trace context into the ext_proc gRPC
|
||||
// The dataplane doesn't propagate trace context into the ext_proc gRPC
|
||||
// stream's metadata — the per-request traceparent arrives in the
|
||||
// HTTP headers carried inside the ProcessingRequest payload. Extract
|
||||
// from there so our span links to the Envoy ingress span.
|
||||
// from there so our span links to the gateway's ingress span.
|
||||
ctx = otel.GetTextMapPropagator().Extract(ctx, propagation.MapCarrier(md.Headers))
|
||||
ctx, span := otel.Tracer(extproc.ServiceName).Start(ctx, "ExtProc.RequestHeaders")
|
||||
defer span.End()
|
||||
|
||||
@@ -133,7 +133,7 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
// shutdownCtx signals SIGTERM/SIGINT; kept separate from the work context
|
||||
// so in-flight ext_proc streams (parked requests, most of all) are not
|
||||
// cancelled the moment the signal arrives. drainOnShutdown drives the
|
||||
// shutdown sequence: readiness flip → route-drain delay → Envoy drain →
|
||||
// shutdown sequence: readiness flip → route-drain delay → dataplane drain →
|
||||
// ext_proc drain → stop the rest.
|
||||
shutdownCtx, stopSignals := signal.NotifyContext(ctx, syscall.SIGINT, syscall.SIGTERM)
|
||||
defer stopSignals()
|
||||
@@ -152,9 +152,9 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
}
|
||||
parkCfg := s.cfg.ParkedRequest.Normalized()
|
||||
|
||||
// The drain-complete marker persists container restarts (emptyDir); a
|
||||
// stale one would release the Envoy preStop hook the moment a later drain
|
||||
// begins.
|
||||
// The drain-complete marker persists container restarts (emptyDir); a stale
|
||||
// one would release the dataplane container's preStop hook the moment a
|
||||
// later drain begins.
|
||||
removeStaleDrainMarker(ctx, s.cfg.DrainCompleteFile)
|
||||
|
||||
serverboot.InitLogger()
|
||||
@@ -263,8 +263,8 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
|
||||
s.health = newRouterHealth(s.cfg.HealthInterval, s.clientset, s.apiClient, s.cfg)
|
||||
|
||||
// The dataplane control plane — the xDS server and the ActorTemplate
|
||||
// controller — configures the *ingress* Envoy. The egress Envoy is
|
||||
// The ingress control plane — the xDS server and the ActorTemplate
|
||||
// controller — configures the *ingress* dataplane. The egress gateway is
|
||||
// statically configured, so an egress-only instance runs neither and needs
|
||||
// no Kubernetes access.
|
||||
if s.cfg.Mode.ServesIngress() {
|
||||
@@ -281,7 +281,7 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
})
|
||||
|
||||
// Start ExtProc Server. Driven by the drain sequence rather than context
|
||||
// cancel: ext_proc is failClosed, so it must outlive Envoy's drain.
|
||||
// cancel: ext_proc is failClosed, so it must outlive the dataplane's drain.
|
||||
extprocGRPC := s.extprocSrv.NewGRPCServer()
|
||||
g.Go(func() error {
|
||||
slog.InfoContext(ctx, "Starting ExtProc Server", slog.Int("port", s.cfg.ExtprocPort))
|
||||
@@ -356,8 +356,8 @@ func (s *RouterServer) Run(ctx context.Context) error {
|
||||
// OTEL_EXPORTER_OTLP_ENDPOINT, which the router's own exporter reads too and
|
||||
// which legitimately carries forms Envoy's plaintext tracer cluster cannot
|
||||
// reach — an https collector, most of all. Refusing to start would take the
|
||||
// xDS control plane for every Envoy in the mesh down over a tracing endpoint
|
||||
// that works fine for its other reader. Losing Envoy's spans is the smaller
|
||||
// xDS control plane for every ingress Envoy down over a tracing endpoint that
|
||||
// works fine for its other reader. Losing Envoy's spans is the smaller
|
||||
// failure, so take it and say so loudly.
|
||||
func setOtlpCollector(ctx context.Context, xdsSrv *XdsServer, addr string) {
|
||||
if err := xdsSrv.SetOtlpCollector(addr); err != nil {
|
||||
|
||||
@@ -47,7 +47,7 @@ intercepted and carried over mTLS to a gateway that verifies who is making the r
|
||||
pod as an ext_proc sidecar, the same binary that serves ingress, started with `--mode=egress`)
|
||||
re-verifies the chain, requires exactly one `ActorIdentity` extension with `purpose: atunnel`,
|
||||
and calls the ate API (`GetActor`). It returns **403** unless the certified **UID** matches a
|
||||
real, `RUNNING` actor. This mirrors the ingress gateway's Envoy + ext_proc co-location; a
|
||||
real, `RUNNING` actor. This mirrors the ingress gateway's dataplane + ext_proc co-location; a
|
||||
standalone/shared ext_proc is a future step.
|
||||
|
||||
## Components
|
||||
@@ -126,7 +126,7 @@ kubectl -n ate-system logs deploy/atenet-egress | grep '\[egress\]'
|
||||
|
||||
# The co-located ext_proc sidecar logs the identity decision, including the UID it authorized on:
|
||||
kubectl -n ate-system logs deploy/atenet-egress -c ext-proc | grep -i 'egress identity\|egress denied'
|
||||
# egress identity authenticated atespace=demo actor=egress-demo actorUid=… status=STATUS_RUNNING
|
||||
# egress identity authenticated atespace=demo actor=egress-demo actorUid=… destination=<TARGET_IP>:80
|
||||
```
|
||||
|
||||
The `whoami` body shows `RemoteAddr: <atenet-egress pod IP>` — proof the request egressed
|
||||
|
||||
@@ -474,7 +474,7 @@ advertises the 4318 HTTP endpoint.
|
||||
use.** Envoy's tracer cluster is plaintext h2c, so an `https://` endpoint is
|
||||
neither honored nor silently downgraded: the router logs a warning, turns
|
||||
Envoy-side tracing off, and starts normally. Its own spans are unaffected.
|
||||
Taking the xDS control plane down for every Envoy in the mesh over a tracing
|
||||
Taking the xDS control plane down for every ingress Envoy over a tracing
|
||||
endpoint that works fine for the router's own exporter would be the larger
|
||||
failure.
|
||||
|
||||
|
||||
+1
-1
@@ -95,7 +95,7 @@ Below is a collection of finer-grained efforts which we believe align with the a
|
||||
* Agent Development Kit (ADK) Native Support: Developing first-class bindings for ADK, allowing developers to build stateful agents that natively leverage Substrate’s lifecycle management and persistent working memory.
|
||||
* LangChain Remote Execution Provider: A dedicated provider for LangChain to run complex, long-running agent tools in durable, sandboxed environments.
|
||||
* Native MCP Server Hosting: Built-in support for deploying Model Context Protocol (MCP) servers as managed Substrate Actors, creating a secure tool ecosystem for any LLM.
|
||||
* Actor-to-Actor (A2A) Calling Model: Standardized protocol for actors to discover and call other actors within the mesh via the gateway.
|
||||
* Actor-to-Actor (A2A) Calling Model: Standardized protocol for actors to discover and call other actors within Substrate via the gateway.
|
||||
* Native MCP Tool Hosting: Ability to define and deploy standard Model Context Protocol (MCP) servers as managed Substrate Actors, providing a plug-and-play ecosystem for agentic tools.
|
||||
|
||||
### Operability
|
||||
|
||||
@@ -19,8 +19,7 @@
|
||||
# hack/install-ate-kind.sh --deploy-demo-egress
|
||||
#
|
||||
# The egress demo Actor accepts {"url":"..."} and performs an HTTP GET. With
|
||||
# egress turned on (ate-api-server --egress-gateway-address, which ateapi stamps
|
||||
# onto every atelet Run/Restore), the Actor's outbound TCP is nftables-REDIRECTed
|
||||
# egress turned on, the Actor's outbound TCP is nftables-REDIRECTed
|
||||
# into atunnel, wrapped in mTLS + HTTP CONNECT, and sent to the Envoy egress
|
||||
# gateway, which terminates CONNECT and tunnels to the real destination. This
|
||||
# script drives that path and shows the gateway's access log proving the actor's
|
||||
|
||||
@@ -72,7 +72,7 @@ func (c *RouterClient) Close() {
|
||||
c.stop()
|
||||
}
|
||||
|
||||
// Get issues GET path to actor through the router, setting the actor's mesh Host
|
||||
// Get issues GET path to actor through the router, setting the actor's DNS Host
|
||||
// so the router routes (and resumes) it. The caller must close the body.
|
||||
func (c *RouterClient) Get(ctx context.Context, actorRef resources.ActorRef, path string) (*http.Response, error) {
|
||||
return c.request(ctx, http.MethodGet, actorRef, path, nil)
|
||||
|
||||
@@ -47,7 +47,7 @@ func (r ActorRef) LogValue() slog.Value {
|
||||
)
|
||||
}
|
||||
|
||||
// DNSName returns the mesh DNS name the actor is reachable at.
|
||||
// DNSName returns the uniform DNS name the actor is reachable at.
|
||||
// This is: "<name>.<atespace>.actors.resources.substrate.ate.dev".
|
||||
func (r ActorRef) DNSName() string {
|
||||
return r.Name + "." + r.Atespace + "." + ActorDNSSuffix
|
||||
|
||||
@@ -76,13 +76,7 @@ data:
|
||||
codec_type: HTTP1
|
||||
upgrade_configs:
|
||||
- upgrade_type: CONNECT
|
||||
# SEE(lior): this pair is how the actor's certificate reaches the
|
||||
# ext_proc handler. ext_proc can request Envoy attributes, but none
|
||||
# of them carry a custom X.509 extension, so the only way to check
|
||||
# ActorIdentity in Go is to have Envoy hand over the raw chain.
|
||||
# SANITIZE_SET drops whatever x-forwarded-client-cert the client
|
||||
# sent and writes Envoy's own view of the verified peer, and
|
||||
# chain: true puts the full URL-encoded PEM chain in it.
|
||||
# TODO(liorlieberman): Can we make this cleaner?
|
||||
forward_client_cert_details: SANITIZE_SET
|
||||
set_current_client_cert_details:
|
||||
chain: true
|
||||
@@ -155,7 +149,7 @@ data:
|
||||
clusters:
|
||||
# ext_proc gRPC server = the atenet router, co-located in this pod as a
|
||||
# sidecar and called over localhost (same topology the ingress gateway uses
|
||||
# for its Envoy + ext_proc).
|
||||
# for its dataplane + ext_proc).
|
||||
- name: ext_proc_server
|
||||
type: STATIC
|
||||
lb_policy: ROUND_ROBIN
|
||||
@@ -210,6 +204,7 @@ spec:
|
||||
sysctls:
|
||||
- name: net.ipv4.ip_unprivileged_port_start
|
||||
value: "0"
|
||||
terminationGracePeriodSeconds: 60
|
||||
containers:
|
||||
- name: envoy
|
||||
image: envoyproxy/envoy:v1.34-latest
|
||||
@@ -227,6 +222,26 @@ spec:
|
||||
- atenet-egress
|
||||
- --service-cluster
|
||||
- atenet-egress
|
||||
# Prevents Envoy from fast-exiting on SIGTERM before the ext-proc
|
||||
# sidecar finishes its drain sequence. Polls for the drain-complete
|
||||
# marker the sidecar writes on the shared emptyDir, terminating Envoy
|
||||
# as soon as the drain completes (or at terminationGracePeriodSeconds
|
||||
# if the sidecar crashes).
|
||||
#
|
||||
# TODO(liorlieberman): decide the drain policy for long-lived CONNECT
|
||||
# tunnels. envoyDrainer polls downstream_cx_active until it reaches
|
||||
# zero, which suits ingress (short request/response connections) but
|
||||
# not egress: atunnel holds a tunnel open for the life of the actor's
|
||||
# connection, so an actor still streaming when the rollout starts keeps
|
||||
# the count above zero until the ~15s window expires. Every rollout
|
||||
# with active egress will therefore log "N downstream connections still
|
||||
# active at the drain deadline". Either accept that as the honest
|
||||
# outcome (we tried to drain, then cut), or give egress its own shorter
|
||||
# window so we stop pretending a tunnel will close on its own.
|
||||
lifecycle:
|
||||
preStop:
|
||||
exec:
|
||||
command: ["sh", "-c", "while [ ! -f /var/run/atenet/drain-complete ]; do sleep 0.5; done"]
|
||||
ports:
|
||||
- name: https
|
||||
containerPort: 443
|
||||
@@ -256,6 +271,9 @@ spec:
|
||||
- name: actor-id-ca-certs
|
||||
mountPath: /run/actor-id-ca-certs
|
||||
readOnly: true
|
||||
- name: drain-signal
|
||||
mountPath: /var/run/atenet
|
||||
readOnly: true
|
||||
# Co-located ext_proc server: the same atenet router binary the ingress
|
||||
# gateway runs, started with --mode=egress so it serves only the egress
|
||||
# ext_proc handler (no xDS server, no ActorTemplate controller, no
|
||||
@@ -282,6 +300,11 @@ spec:
|
||||
# denies every egress CONNECT with a 503.
|
||||
- --actor-identity-ca-file=/run/actor-id-ca-certs/ca.crt
|
||||
- --otlp-collector-address=
|
||||
# The egress Envoy's admin listener is on 15000, not the 9901 the flag
|
||||
# defaults to (that is the ingress gateway's port). Without this the
|
||||
# drain sequence dials a closed port, reads the connection refusal as
|
||||
# "Envoy already exited", and reports a drain it never performed.
|
||||
- --envoy-admin-address=127.0.0.1:15000
|
||||
env:
|
||||
- name: POD_NAME
|
||||
valueFrom:
|
||||
@@ -311,10 +334,14 @@ spec:
|
||||
- name: actor-id-ca-certs
|
||||
mountPath: /run/actor-id-ca-certs
|
||||
readOnly: true
|
||||
- name: drain-signal
|
||||
mountPath: /var/run/atenet
|
||||
volumes:
|
||||
- name: config
|
||||
configMap:
|
||||
name: atenet-egress
|
||||
- name: drain-signal
|
||||
emptyDir: {}
|
||||
- name: servicedns
|
||||
projected:
|
||||
sources:
|
||||
|
||||
@@ -161,9 +161,11 @@ spec:
|
||||
image: ko://github.com/agent-substrate/substrate/cmd/atenet
|
||||
args:
|
||||
- "router"
|
||||
# Serve the ingress direction only: the xDS server and ActorTemplate
|
||||
# controller for the Envoy in this pod, plus the ingress ext_proc
|
||||
# handler. Egress runs the same binary with --mode=egress in its own
|
||||
# Serve the ingress direction only: the ingress ext_proc handler, plus
|
||||
# the xDS server and ActorTemplate controller that configure the
|
||||
# dataplane in this pod when it is Envoy (the agentgateway overlay is
|
||||
# statically configured and starts neither).
|
||||
# Egress runs the same binary with --mode=egress in its own
|
||||
# Deployment (manifests/ate-install/atenet-egress.yaml), because the two
|
||||
# directions scale independently.
|
||||
- "--mode=ingress"
|
||||
|
||||
Reference in New Issue
Block a user