feat(auth,sts,egress): task-scoped Biscuit attenuation, two-token STS, and egress brokers - #578
Conversation
There was a problem hiding this comment.
Code Review
This pull request removes the nano-init sandbox agent and introduces Task-Scoped Authorization (TAR) along with a Two-Token Model, allowing the exchange of platform JWTs for delegated Biscuits and the minting of border JWTs for cloud federation. It also adds pluggable credential brokering, content inspection (via Google Cloud Model Armor and Envoy ext_proc), and a new sam-node forward command. Feedback on the changes highlights a security issue in internal/controlplane/sts.go where peer IDs must be canonicalized before performing ban lookups, and a correctness bug in internal/node/egress_inspect.go where JSON-escaped PII can bypass redaction when using direct byte replacement.
| for _, peerID := range []string{claims.NodePeerID, claims.ActorNodePeerID, claims.ClientPeerID} { | ||
| if peerID == "" { | ||
| continue | ||
| } | ||
| banned, err := s.store.IsNodeBanned(ctx, peerID) | ||
| if err != nil { | ||
| return nil, http.StatusInternalServerError, fmt.Errorf("failed to check node ban: %w", err) | ||
| } | ||
| if banned { | ||
| return nil, http.StatusForbidden, fmt.Errorf("peer %s is banned", peerID) | ||
| } | ||
| } |
There was a problem hiding this comment.
Rule 1 Violation: Peer IDs must be canonicalized at the boundary\n\nThe peer IDs extracted from the Biscuit (claims.NodePeerID, claims.ActorNodePeerID, and claims.ClientPeerID) are raw strings from the wire and are not guaranteed to be in canonical base58 form. Passing them directly to s.store.IsNodeBanned is a security vulnerability: if an attacker presents a token with a non-canonical peer ID encoding (e.g., CIDv1 base32), the ban check will silently miss.\n\nDecode each peer ID using peer.Decode and use the canonical .String() representation for the ban lookup.
\tfor _, rawPeerID := range []string{claims.NodePeerID, claims.ActorNodePeerID, claims.ClientPeerID} {\n\t\tif rawPeerID == \"\" {\n\t\t\tcontinue\n\t\t}\n\t\tpID, err := peer.Decode(rawPeerID)\n\t\tif err != nil {\n\t\t\treturn nil, http.StatusForbidden, fmt.Errorf(\"invalid peer ID %q: %w\", rawPeerID, err)\n\t\t}\n\t\tcanonical := pID.String()\n\t\tbanned, err := s.store.IsNodeBanned(ctx, canonical)\n\t\tif err != nil {\n\t\t\treturn nil, http.StatusInternalServerError, fmt.Errorf(\"failed to check node ban: %w\", err)\n\t\t}\n\t\tif banned {\n\t\t\treturn nil, http.StatusForbidden, fmt.Errorf(\"peer %s is banned\", canonical)\n\t\t}\n\t}References
- Rule 1: Peer IDs must be canonicalized at the boundary. Every cache, ban set, storage row, and map[string] in SAM is keyed on that canonical form. A raw string passed to a lookup keyed on the canonical form (e.g., store.IsNodeBanned) silently misses if non-canonical. (link)
- Peer IDs must be canonicalized at the boundary before performing lookups keyed on the canonical form (such as IsNodeBanned).
| func (s *EgressService) inspectModelArmorRequest(ctx context.Context, cfg *api.ModelArmor, body []byte, callerCtx egressCallerContext) ([]byte, bool, error) { | ||
| if len(body) == 0 { | ||
| return body, false, nil | ||
| } | ||
| promptText := extractInspectableText(body) | ||
| if promptText == "" { | ||
| return body, false, nil | ||
| } | ||
| blocked, replacement, err := s.callModelArmorAPI(ctx, cfg, "sanitizeUserPrompt", "userPromptData", promptText, callerCtx) | ||
| if err != nil || blocked { | ||
| return body, blocked, err | ||
| } | ||
| if replacement != "" && replacement != promptText { | ||
| body = bytes.ReplaceAll(body, []byte(promptText), []byte(replacement)) | ||
| } | ||
| return body, false, nil | ||
| } | ||
|
|
||
| func (s *EgressService) inspectModelArmorResponse(ctx context.Context, cfg *api.ModelArmor, body []byte, callerCtx egressCallerContext) ([]byte, bool, error) { | ||
| if len(body) == 0 { | ||
| return body, false, nil | ||
| } | ||
| respText := extractInspectableText(body) | ||
| if respText == "" { | ||
| return body, false, nil | ||
| } | ||
| blocked, replacement, err := s.callModelArmorAPI(ctx, cfg, "sanitizeModelResponse", "modelResponseData", respText, callerCtx) | ||
| if err != nil || blocked { | ||
| return body, blocked, err | ||
| } | ||
| if replacement != "" && replacement != respText { | ||
| body = bytes.ReplaceAll(body, []byte(respText), []byte(replacement)) | ||
| } | ||
| return body, false, nil | ||
| } |
There was a problem hiding this comment.
Correctness & Security Bug: PII Redaction Bypass via JSON Escaping\n\nUsing bytes.ReplaceAll to replace promptText (or respText) directly in the raw JSON body will fail if the text contains any JSON-escaped characters (such as newlines \n, tabs \t, backslashes \\, or quotes \").\n\nFor example, if the prompt contains a literal newline, extractInspectableText parses the JSON and returns the unescaped string containing 0x0a. Model Armor inspects it and returns a redacted string also containing 0x0a. However, in the raw JSON body, the newline is represented as the two-character sequence \n (0x5c 0x6e). As a result, bytes.ReplaceAll will fail to find a match, and the sensitive/unredacted PII will be silently forwarded upstream.\n\nTo fix this, parse the JSON, modify the unescaped string values in-place within the parsed map/struct, and then re-marshal the JSON.
func (s *EgressService) inspectModelArmorRequest(ctx context.Context, cfg *api.ModelArmor, body []byte, callerCtx egressCallerContext) ([]byte, bool, error) {\n\tif len(body) == 0 {\n\t\treturn body, false, nil\n\t}\n\tvar doc map[string]any\n\tif err := json.Unmarshal(body, &doc); err != nil {\n\t\tpromptText := strings.TrimSpace(string(body))\n\t\tif promptText == \"\" {\n\t\t return body, false, nil\n\t\t}\n\t\tblocked, replacement, err := s.callModelArmorAPI(ctx, cfg, \"sanitizeUserPrompt\", \"userPromptData\", promptText, callerCtx)\n\t\tif err != nil || blocked {\n\t\t return body, blocked, err\n\t\t}\n\t\tif replacement != \"\" && replacement != promptText {\n\t\t return []byte(replacement), false, nil\n\t\t}\n\t\treturn body, false, nil\n\t}\n\n\tmodified := false\n\tif msgs, ok := doc[\"messages\"].([]any); ok {\n\t\tfor _, m := range msgs {\n\t\t if mm, ok := m.(map[string]any); ok {\n\t\t if c, ok := mm[\"content\"].(string); ok && c != \"\" {\n\t\t blocked, replacement, err := s.callModelArmorAPI(ctx, cfg, \"sanitizeUserPrompt\", \"userPromptData\", c, callerCtx)\n\t\t if err != nil || blocked {\n\t\t return body, blocked, err\n\t\t }\n\t\t if replacement != \"\" && replacement != c {\n\t\t mm[\"content\"] = replacement\n\t\t modified = true\n\t\t }\n\t\t }\n\t\t }\n\t\t}\n\t}\n\n\tif modified {\n\t\tnewBody, err := json.Marshal(doc)\n\t\tif err != nil {\n\t\t return body, false, err\n\t\t}\n\t\treturn newBody, false, nil\n\t}\n\treturn body, false, nil\n}\n\nfunc (s *EgressService) inspectModelArmorResponse(ctx context.Context, cfg *api.ModelArmor, body []byte, callerCtx egressCallerContext) ([]byte, bool, error) {\n\tif len(body) == 0 {\n\t\treturn body, false, nil\n\t}\n\tvar doc map[string]any\n\tif err := json.Unmarshal(body, &doc); err != nil {\n\t\trespText := strings.TrimSpace(string(body))\n\t\tif respText == \"\" {\n\t\t return body, false, nil\n\t\t}\n\t\tblocked, replacement, err := s.callModelArmorAPI(ctx, cfg, \"sanitizeModelResponse\", \"modelResponseData\", respText, callerCtx)\n\t\tif err != nil || blocked {\n\t\t return body, blocked, err\n\t\t}\n\t\tif replacement != \"\" && replacement != respText {\n\t\t return []byte(replacement), false, nil\n\t\t}\n\t\treturn body, false, nil\n\t}\n\n\tmodified := false\n\tif choices, ok := doc[\"choices\"].([]any); ok {\n\t\tfor _, ch := range choices {\n\t\t if cm, ok := ch.(map[string]any); ok {\n\t\t if msg, ok := cm[\"message\"].(map[string]any); ok {\n\t\t if c, ok := msg[\"content\"].(string); ok && c != \"\" {\n\t\t blocked, replacement, err := s.callModelArmorAPI(ctx, cfg, \"sanitizeModelResponse\", \"modelResponseData\", c, callerCtx)\n\t\t if err != nil || blocked {\n\t\t return body, blocked, err\n\t\t }\n\t\t if replacement != \"\" && replacement != c {\n\t\t msg[\"content\"] = replacement\n\t\t modified = true\n\t\t }\n\t\t }\n\t\t }\n\t\t }\n\t\t}\n\t}\n\n\tif modified {\n\t\tnewBody, err := json.Marshal(doc)\n\t\tif err != nil {\n\t\t return body, false, err\n\t\t}\n\t\treturn newBody, false, nil\n\t}\n\treturn body, false, nil\n}|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces Task-Scoped Authorization (TAR) for the SAM mesh, enabling fine-grained control over credential attenuation and egress inspection. Key changes include the introduction of api.TaskAuthorizationRule and api.FactTARBlock for Biscuit attenuation, the implementation of OIDCFederationExchanger and AWSAssumeRoleExchanger for cloud STS integration, and the addition of Envoy ext_proc support for egress inspection. The PR also updates the node's token renewal logic to support continuous platform attestation via fresh JWTs and refines the peer ID canonicalization and revocation logic. My review identified a critical security vulnerability in NarrowOIDCScopes where empty TAR rules could lead to a bypass of scope narrowing, and an efficiency issue regarding the recreation of HTTP clients in the egress inspection chain.
| func (s *EgressService) runExtProcRequestPhase(r *http.Request, cfg *api.ExtProc, reqBody []byte, callerCtx egressCallerContext) (*extProcClientStream, *extprocv3http.ProcessingMode, *extprocv3.ImmediateResponse, []byte, error) { | ||
| msgTimeout := defaultExtProcMessageTimeout | ||
| if cfg.GetMessageTimeout().IsValid() && cfg.GetMessageTimeout().AsDuration() > 0 { | ||
| msgTimeout = cfg.GetMessageTimeout().AsDuration() | ||
| } | ||
| client, endpoint, err := buildExtProcHTTPClient(cfg, s.secretsDir) | ||
| if err != nil { | ||
| return nil, nil, nil, reqBody, err | ||
| } | ||
| stream, err := dialExtProcStream(r.Context(), client, endpoint, msgTimeout*4) |
There was a problem hiding this comment.
Efficiency Issue: HTTP Connection Pool Disabling
buildExtProcHTTPClient is called on every request, which recreates the http.Client and http.Transport on every single invocation.
This completely disables HTTP connection pooling (keep-alive), leading to a new TCP connection and TLS handshake for every external processor call. Under load, this will cause severe latency overhead and port exhaustion.
We should initialize and cache the client/transport once per EgressService during initialization.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request replaces the legacy sandboxed agent architecture with a task-scoped authorization model utilizing Biscuit attenuation (tar_block) and a two-token border STS model. It introduces pluggable credential brokers, L7 content inspection (Model Armor and Envoy ext_proc), and secure L4 TCP tunneling with TLS ClientHello SNI verification, alongside support for workload OIDC issuers, Envoy ext_authz integration, and an STS benchmarking suite. The review feedback highlights critical performance and maintainability improvements in the control plane's STS implementation, specifically addressing a database query bottleneck in the revocation check path, redundant parsing of static Datalog rules on every request, and duplicated dead code that should be removed.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces Task-Scoped Authorization (TAR) and a two-token model to support stateless token exchange and border STS. It also adds pluggable credential brokers, content inspection (Model Armor and Envoy ext_proc), and unencrypted TLS CONNECT tunneling to egress destinations, while refactoring sandboxed agents and adding workload OIDC issuer support. The review feedback highlights three critical issues: first, peer ID strings extracted from Biscuits must be canonicalized at the boundary to prevent security bypasses in downstream lookups; second, response buffering during egress inspection must be bounded to prevent Out-Of-Memory crashes on large upstream responses; and third, AWS session policy parsing must robustly handle both single objects and arrays for the Statement field.
| } | ||
| ep := kind.ExtProc | ||
| stream, mode, imm, newBody, err := s.runExtProcRequestPhase(r, ep, reqBody, callerCtx) | ||
| if err != nil { | ||
| if !ep.GetFailureModeAllow() { |
There was a problem hiding this comment.
When response inspection is enabled, the entire upstream response is buffered in memory using httptest.NewRecorder(). However, the response size limit (max_buffered_bytes) is only checked after the response has been fully buffered. If an upstream destination returns an extremely large response (e.g., a large file or an infinite stream), the node will buffer it entirely, leading to high memory consumption or an Out-Of-Memory (OOM) crash.
To prevent this, wrap the recorder or limit the bytes written during proxy.ServeHTTP when response buffering is active.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces significant changes to the SAM control plane and node architecture, primarily focusing on task-scoped authorization (TAR) and a two-token model for egress gateways. Key changes include the introduction of the 'tar_block' Datalog predicate for Biscuit attenuation, the addition of an STS border JWT minting flow, and improved credential broker integration for egress destinations. The PR also updates the node's enrollment and refresh logic to support platform-native JWTs (e.g., GCE/K8s) and adds a new forward command for TCP tunneling. My feedback highlights a simplification for the handleForwardConn implementation to avoid manual buffer management, which is error-prone and redundant when using bufio.Reader.
| if buffered := br.Buffered(); buffered > 0 { | ||
| peeked, err := br.Peek(buffered) | ||
| if err != nil { | ||
| return err | ||
| } | ||
| if _, err := clientConn.Write(peeked); err != nil { | ||
| return err | ||
| } | ||
| _, _ = br.Discard(buffered) | ||
| } | ||
|
|
||
| var wg sync.WaitGroup | ||
| wg.Add(2) | ||
| go func() { | ||
| defer wg.Done() | ||
| _, _ = io.Copy(nodeConn, clientConn) | ||
| if tc, ok := nodeConn.(interface{ CloseWrite() error }); ok { | ||
| _ = tc.CloseWrite() | ||
| } | ||
| }() | ||
| go func() { | ||
| defer wg.Done() | ||
| _, _ = io.Copy(clientConn, nodeConn) | ||
| if tc, ok := clientConn.(interface{ CloseWrite() error }); ok { | ||
| _ = tc.CloseWrite() | ||
| } | ||
| }() | ||
| wg.Wait() |
There was a problem hiding this comment.
Section 5. API surfaces and secrets
Instead of manually peeking, writing, and discarding buffered bytes from br before copying from nodeConn directly, you can simplify the logic and guarantee correctness by copying directly from br (which wraps nodeConn). bufio.Reader.Read automatically drains any buffered bytes first before reading from the underlying connection, making the manual buffer-draining logic redundant and error-prone.
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
_, _ = io.Copy(nodeConn, clientConn)
if tc, ok := nodeConn.(interface{ CloseWrite() error }); ok {
_ = tc.CloseWrite()
}
}()
go func() {
defer wg.Done()
_, _ = io.Copy(clientConn, br)
if tc, ok := clientConn.(interface{ CloseWrite() error }); ok {
_ = tc.CloseWrite()
}
}()
wg.Wait()…uation and cross-language verification
…TCP tunnels, sam-bench sts, and SDK attenuation
…ation, STS, and egress brokers
…refresh re-attestation, and security architecture docs
…t_proc HTTP clients
…k skew via Date header
…atic STS Datalog rules
…nse buffering, and accept single AWS Statement
f206afc to
eddd69a
Compare
Summary
Implements the task-scoped authorization, two-token STS border federation, and egress credential broker architecture across
sam-control-plane,sam-node,sam-router, and the TypeScript/Python SDKs:cmd/sam-box,cmd/nano-init,internal/sambox, andX-Sam-Agentself-asserted header machinery in favor of cryptographic task credentials.tar_block) & single-pass verification:TaskAuthorizationRule(api/sam.proto) encoded astar_block("<base64url-proto>")in appended Biscuit blocks (0rules,0checks,1fact) with strict pre-authorizer structural bounds (MaxAttenuationBlocks = 8,MaxTARBytes = 4096) and intersection semantics across Go, TypeScript, and Python (sdk/testdata/tar_conformance.json).ext_authz, and revocations:POST /token/exchange(stateless JWT-to-Delegated-Biscuit exchange withactor_nodeandclient_peer_id),POST /sts/token(ES256 border JWT minting with/.well-known/openid-configurationand/jwks),GET /revocations,POST /oauth/revoke, MCP OAuth 2.1 (/.well-known/oauth-protected-resource,/oauth/authorize,/oauth/token), and Envoyext_authz(/ext_authz).CloudTokenExchanger(oidc_federation,aws_assume_role,platform_identity,static_secret), Model Armor (model_armor) and Envoyext_proc(ext_proc) inspectors, named TLS-SNI-verified TCPCONNECTtunnels (EGRESS_MODE_TCP,sam-node forward),sam-bench sts, and offline SDK.attenuate()/.seal().*prefix/suffix wildcards on claim-backedPolicyBinding.members,--workload-issuerand--workload-session-ttl(including Helm chart values), continuous JWT re-attestation onPOST /refreshacrosssam-node,sam-router, and the SDKs, andTokenSourcesupport for GCE/Cloud Run metadata (--cloud-provider=gcp|auto) and SDKjwtcallbacks.