diff --git a/cmd/proxy/serve.go b/cmd/proxy/serve.go index 4e537e0..a927ae5 100644 --- a/cmd/proxy/serve.go +++ b/cmd/proxy/serve.go @@ -33,8 +33,6 @@ const shutdownTimeout = 30 * time.Second // fxLogger swallows fx's event stream except for Started/Stopped failures, // which it forwards to the app logger. Configuration errors are caught // separately via app.Err() so they are visible at startup. -// -// The net effect is fx.NopLogger plus visible lifecycle hook failures. type fxLogger struct { log logger.Logger } diff --git a/internal/api/auth.go b/internal/api/auth.go index 46670cf..9d5ed76 100644 --- a/internal/api/auth.go +++ b/internal/api/auth.go @@ -18,16 +18,10 @@ import ( ) // Auth authenticates an inbound stream by delegating the decision to an -// extension server implementing api.auth.v1.AuthService. It is the escape hatch -// for identity systems the built-in authenticators do not cover: the proxy -// asks, the operator's server decides. -// -// The headers name the metadata carrying the caller's credentials. They are -// declared rather than discovered because a verdict reports only admit-or-deny -// and says nothing about which headers mattered, and the proxy needs to know two -// things: which values to lift into the request, and which to report as -// [Auth.SecureHeaders] so they are stripped from the stream before it reaches an -// upstream, where a caller credential would collide with the proxy's own. +// extension server implementing api.auth.v1.AuthService. Its headers name the +// metadata carrying the caller's credentials: their values are lifted into the +// request, and they are reported as [Auth.SecureHeaders] so they are stripped +// from the stream before it reaches an upstream. type Auth struct { client auth.AuthServiceClient headers []string @@ -50,31 +44,16 @@ func NewAuth(cc grpc.ClientConnInterface, secureHeaders []string) *Auth { } // Authenticate asks the extension server whether the caller may proceed. Only an -// explicit DECISION_ALLOW admits the stream: an error, a denial, and an answer -// carrying no verdict all deny it, so a server that is down, misconfigured, or -// newer than this build fails the request closed rather than opening the gateway -// to everyone for as long as it is that way. -// -// A denial reaches the caller as an [rpc.Reject], whose message is generic while -// the provider's reason becomes the server-side detail: a provider writes that -// reason for whoever operates it, not for the caller it just turned away. An -// error keeps the provider's status code, since it tells a worker whether to fix -// its credential or retry; a decision carries no code, so the proxy supplies one. +// explicit DECISION_ALLOW admits the stream; an error, a denial, or an answer +// carrying no verdict denies it, so the request fails closed. Every rejection is +// an [rpc.Reject] with a generic message and the provider's reason as the +// server-side detail. An error keeps the provider's status code; a denial is +// PermissionDenied and a missing verdict is Internal. // // The declared credential headers are lifted into the request and withheld from -// the forwarded metadata, so each credential reaches the server in exactly one -// place. That separation is what lets the proxy hold a credential of its own to -// this server: metadata carries the proxy's, the request carries the caller's, -// and neither has to be told apart from the other on a shared header. It also -// puts the caller's credential out of reach of the interceptor [outbound.DialOptions] -// installs, which deletes the proxy's credential header from forwarded metadata -// and cannot tell that on this one call that header is the subject of the request -// rather than incidental cargo. -// -// The caller's remaining metadata is forwarded so the server can weigh context -// such as the method being invoked. gRPC drops reserved keys (":authority", -// "user-agent", "content-type", "grpc-*") when writing the request, so a caller -// cannot reach the extension server's transport this way. +// the forwarded metadata. The caller's remaining metadata is forwarded, except +// the reserved keys gRPC drops (":authority", "user-agent", "content-type", +// "grpc-*"), so a caller cannot reach the extension server's transport this way. func (a *Auth) Authenticate(ctx context.Context, target meta.Target, md metadata.MD) error { req := &auth.AuthRequest{ Target: &auth.Target{FullName: target.FullName, Namespace: target.Namespace}, diff --git a/internal/api/fx.go b/internal/api/fx.go index f134089..5846844 100644 --- a/internal/api/fx.go +++ b/internal/api/fx.go @@ -13,14 +13,10 @@ import ( ) // extensionKeyPrefix namespaces extension-server entries in the shared -// connection pool. -// -// A static resolver uses its dial address as the pool key, and Pool.ConnOrCreate -// returns the existing connection for a key while ignoring the options passed -// with it. Upstream hostPorts are unique among themselves, but nothing stops an -// extension server from sitting on the same host:port as an upstream, so without -// a distinct key the two would collapse onto whichever was dialed first and -// silently inherit its TLS settings and credentials. +// connection pool. A static resolver keys the pool by dial address and +// Pool.ConnOrCreate ignores the options for an existing key, so without the +// prefix an extension server on an upstream's host:port would share that +// connection and silently inherit its TLS settings and credentials. const extensionKeyPrefix = "extension:" // Module provides the pooled connection for every configured extension server, @@ -86,24 +82,17 @@ type ( // Connections maps an extension server name to a connection to that server. // It carries no lifecycle: closing a connection is the owner's - // responsibility, not the caller's. - // - // Callers get connections rather than finished clients because the two do not - // correspond one-to-one: several keys may live on one extension server, so a - // caller builds one [KMS] per key over the shared connection. + // responsibility, not the caller's. Several keys may live on one extension + // server, so a caller builds one [KMS] per key over the shared connection. Connections map[string]grpc.ClientConnInterface ) // extensionConn builds the pooled connection for a single extension server. // Config rejects a templated hostPort, so the target is always static: the // resolver is fixed and [connect.NewConn] creates the connection here, which the -// module then opens on start. -// -// This is deliberately narrower than the equivalent upstream path in -// internal/proxy. There is no namespace translation, because an extension -// server is not a Temporal service and has no namespaces to rewrite, and no -// payload encryption interceptor, because an extension server is the thing that -// wraps DEKs; sealing its traffic with the vault it backs would be circular. +// module then opens on start. Unlike an upstream connection it installs no +// namespace translation or payload encryption interceptor; an extension server +// wraps DEKs, so sealing its traffic with the vault it backs would be circular. func extensionConn(pool *connect.Pool, s *config.ExtensionServer) (*connect.Conn, error) { var opts []grpc.DialOption diff --git a/internal/auth/authenticator.go b/internal/auth/authenticator.go index b9de5f9..d77c664 100644 --- a/internal/auth/authenticator.go +++ b/internal/auth/authenticator.go @@ -24,9 +24,8 @@ type ( // status error to reject it. SecureHeaders reports the metadata headers the // authenticator consumes, so the proxy can strip the caller's credentials // before forwarding upstream; it returns nil when the authenticator consumes - // no header. An authenticator may name more than one because it need not own - // the header it reads: an external one is told which headers its server - // consumes. + // no header and may name more than one (an external authenticator reports + // the headers its server consumes). // // The target is what the gateway resolved for this stream. Its Namespace is // empty for a request that named none, so an implementation weighing it must diff --git a/internal/auth/jwks.go b/internal/auth/jwks.go index 476cb3e..836c75c 100644 --- a/internal/auth/jwks.go +++ b/internal/auth/jwks.go @@ -57,20 +57,9 @@ type JWKSAuthenticator struct { } // NewJWKSAuthenticator builds a JWKSAuthenticator that resolves signing keys -// from the JWKS at rawURL. -// -// keyfunc.NewDefaultOverrideCtx performs its first key fetch synchronously -// (see github.com/MicahParks/jwkset's NewStorageFromHTTP), so calling it -// directly here would block construction for up to its HTTP timeout if the -// IdP is unreachable. To keep construction non-blocking, the fetch is kicked -// off in a background goroutine (via deferredKeyfunc); until it completes, -// Authenticate reports codes.Unavailable (fail closed, retryable) instead of -// blocking startup. -// -// rawURL is validated synchronously before the goroutine is started, so a -// malformed configuration (bad scheme/host, not a transient IdP outage) fails -// construction immediately instead of surfacing as a permanent -// codes.Unavailable for every future request. +// from the JWKS at rawURL. It returns an error if rawURL is malformed but does +// not block on the IdP: the first key fetch runs in the background, and until +// it completes Authenticate fails closed with codes.Unavailable. func NewJWKSAuthenticator(rawURL string, audiences []string, issuer, header, scheme string) (*JWKSAuthenticator, error) { if err := validateJWKSURL(rawURL); err != nil { return nil, err @@ -233,14 +222,9 @@ func deferredKeyfunc(load func() (jwt.Keyfunc, error)) (jwt.Keyfunc, *atomic.Poi // availability problem (errKeysUnavailable -> codes.Unavailable), not a bad // token. Any other resolver error is treated as an availability problem. // -// keysPresent reports whether the keyset currently holds any keys. resolve is -// the underlying key resolver (in production, keyfunc.Keyfunc's method value). -// -// When keysPresent cannot positively confirm keys exist (e.g. its underlying -// read errors), it should return false; the failure then maps to -// codes.Unavailable (retryable) rather than masking as a bad token - the -// conservative fail-closed-toward-retryable choice for an infra-level read -// error. +// keysPresent reports whether the keyset currently holds any keys; it must +// return false when it cannot confirm keys exist (e.g. its read errors), so +// the failure maps to codes.Unavailable rather than a bad token. func wrapKeyfunc(resolve jwt.Keyfunc, keysPresent func() bool) jwt.Keyfunc { return func(t *jwt.Token) (any, error) { key, err := resolve(t) diff --git a/internal/cloud/endpoint.go b/internal/cloud/endpoint.go index 268407f..8b99862 100644 --- a/internal/cloud/endpoint.go +++ b/internal/cloud/endpoint.go @@ -13,13 +13,11 @@ const endpointSuffix = ".tmprl.cloud" // same address the Cloud SDK dials by default. const APIHostPort = "saas-api" + endpointSuffix + ":443" -// IsEndpoint reports whether hostPort addresses Temporal Cloud. The port is -// optional, and a template action in place of the host is tolerated, since a -// templated address is rendered per request but keeps its domain. -// -// This recognizes per-namespace and regional endpoints. Private-link endpoints -// use per-VPC hostnames that carry no Cloud domain, so those have to be declared -// rather than detected. +// IsEndpoint reports whether hostPort addresses Temporal Cloud through a +// per-namespace or regional endpoint. The port is optional, and a template +// action in place of the host is tolerated. Private-link endpoints use per-VPC +// hostnames with no Cloud domain, so those have to be declared rather than +// detected. func IsEndpoint(hostPort string) bool { host := hostPort if h, _, err := net.SplitHostPort(host); err == nil { diff --git a/internal/cloud/namespace.go b/internal/cloud/namespace.go index 1afca89..8d9ad14 100644 --- a/internal/cloud/namespace.go +++ b/internal/cloud/namespace.go @@ -31,7 +31,7 @@ func ValidateAccountID(id string) error { // ValidateNamespace checks that ns is a well-formed Temporal Cloud namespace // identifier, meaning ".". Every broken rule is reported, not -// just the first, so a caller can show the whole story at once. +// just the first. // // See: https://docs.temporal.io/cloud/namespaces for details func ValidateNamespace(ns string) error { diff --git a/internal/cloud/translation/interceptor.go b/internal/cloud/translation/interceptor.go index 07ff951..ac8389d 100644 --- a/internal/cloud/translation/interceptor.go +++ b/internal/cloud/translation/interceptor.go @@ -15,14 +15,9 @@ type options struct { } // Via sends a translated call over cc instead of continuing down the chain to -// the connection the interceptor is installed on. -// -// The upstream method belongs to a different service, which the connection that -// received the call does not serve: the caller asked a Temporal Service for -// ListNamespaces, and only Temporal Cloud's control plane can answer it. Where -// that service lives is as fixed as the conversions themselves, so the -// translation carries the connection rather than the request being routed to it, -// and a request reaching any upstream is answered the same way. +// the connection the interceptor is installed on, for an upstream method that +// connection does not serve, such as one only Temporal Cloud's control plane +// answers. func Via(cc grpc.ClientConnInterface) Option { return func(o *options) { o.via = cc } } @@ -32,13 +27,6 @@ func Via(cc grpc.ClientConnInterface) Option { // connection, last, so translation is the innermost interceptor: every other // interceptor on the chain then sees the method and message types the caller // asked for rather than the substitute sent upstream. -// -// The result is a slice though it holds a single option today. A [Translation] -// substitutes a unary method, so a unary interceptor is all there is to install; -// translating a streaming method would add a stream interceptor beside it, the -// way the namespace and Cloud-namespace helpers in internal/proxy already pair -// the two. Keeping the slice means that arrives without changing this signature -// or the call sites, which already spread the result. func DialOptions(r *Registry, opts ...Option) []grpc.DialOption { return []grpc.DialOption{grpc.WithChainUnaryInterceptor(unaryClientInterceptor(r, opts...))} } diff --git a/internal/cloud/translation/namespaces.go b/internal/cloud/translation/namespaces.go index 0697b1b..cf51654 100644 --- a/internal/cloud/translation/namespaces.go +++ b/internal/cloud/translation/namespaces.go @@ -40,13 +40,10 @@ const ( ) // listNamespaces translates WorkflowService.ListNamespaces onto -// CloudService.GetNamespaces. -// -// The Cloud API version header is required - without it GetNamespaces fails with -// InvalidArgument - and is pinned to the version the SDK this package compiles -// against defaults to, since that is the same module the message types and their -// versioned fields come from. Bumping go.temporal.io/cloud-sdk moves the -// conversions and the version they were written against together. +// CloudService.GetNamespaces. GetNamespaces fails with InvalidArgument without +// the Cloud API version header, so it is pinned to the default of the +// go.temporal.io/cloud-sdk the message types come from; bumping that module +// moves the conversions and the version together. func listNamespaces() *Translation { return Adapt(listNamespacesMethod, getNamespacesMethod, listNamespacesRequest, listNamespacesResponse). WithHeader(cloudclient.TemporalCloudAPIVersionHeader(), cloudclient.DefaultAPIVersion()) @@ -97,23 +94,14 @@ func listNamespacesResponse( // response ListNamespaces returns per namespace, with state already mapped by // the caller so the deleted filter and this conversion agree on it. // -// Only fields Cloud actually reports carry a value. Cloud has no namespace UUID, -// description, or owner email to give, and it describes replication as regional -// replicas rather than as the clusters ReplicationConfig names, so -// IsGlobalNamespace is derived from how many replicas there are while -// ReplicationConfig is left empty rather than filled with region ids a client -// would read as cluster names. FailoverVersion and FailoverHistory have no -// Cloud equivalent at all. +// Only fields Cloud reports carry a value. Cloud has no namespace UUID, +// description, owner email, or failover data, and its regional replicas are not +// the clusters ReplicationConfig names, so IsGlobalNamespace is derived from the +// replica count and ReplicationConfig is left empty. // -// Empty is not the same as absent, though, and the difference is load-bearing: -// every sub-message a Temporal Service would populate is allocated here even -// when there is nothing to put in it. A frontend builds NamespaceInfo, Config -// and ReplicationConfig unconditionally on every path (the server funnels them -// all through namespaceHandler.createResponse), so clients are written against a -// reply where they are always present - the temporal CLI reads -// resp.ReplicationConfig.ActiveClusterName with no nil check and dies on a nil -// one. An empty message says "Cloud did not report this" just as well as an -// absent one, without breaking a client that has never had to handle absence. +// Every sub-message a Temporal Service populates is allocated even when empty: +// clients assume they are present, and the temporal CLI reads +// resp.ReplicationConfig.ActiveClusterName with no nil check. func describeNamespace( ns *cloudnamespace.Namespace, state enumspb.NamespaceState, diff --git a/internal/cloud/translation/translation.go b/internal/cloud/translation/translation.go index d37f56d..b3f2222 100644 --- a/internal/cloud/translation/translation.go +++ b/internal/cloud/translation/translation.go @@ -114,16 +114,11 @@ func Answer[Req, Resp proto.Message](from string, answer func(Req, Resp) error) } } -// WithHeader stamps key: value on the substituted call and returns t, so a -// mapping can declare the dialect the upstream method needs alongside the -// conversions themselves. It replaces any value the caller sent rather than -// adding to it: the caller did not ask for this upstream method and cannot know -// what its API expects, so its own header is not intent worth preserving. A -// caller invoking that API directly is forwarded untranslated and keeps its -// header. -// -// Headers travel only on a call this translation substituted; a method the -// registry does not translate is untouched. +// WithHeader stamps key: value on the substituted call and returns t. It +// replaces any value the caller sent rather than adding to it. Headers travel +// only on a call this translation substituted: a method the registry does not +// translate is untouched, and a caller invoking the upstream API directly is +// forwarded untranslated and keeps its header. func (t *Translation) WithHeader(key, value string) *Translation { if t.headers == nil { t.headers = make(map[string]string, 1) diff --git a/internal/codecserver/doc.go b/internal/codecserver/doc.go index 12517b5..30f9988 100644 --- a/internal/codecserver/doc.go +++ b/internal/codecserver/doc.go @@ -5,7 +5,7 @@ // # Why It Exists // // A worker or client connecting through the gateway never sees ciphertext: the -// per-upstream proxy opens inbound payloads whether or not sealing is enabled. +// upstream's forwarder opens inbound payloads whether or not sealing is enabled. // Two callers do not have that path. The Temporal Cloud UI reaches Cloud // directly and calls a codec endpoint from the operator's browser, and the // Temporal CLI run with --codec-endpoint reaches a Temporal Service directly. @@ -33,7 +33,7 @@ // // The handler applies a [Codecs], in practice // [github.com/temporalio/temporal-proxy/internal/proxy.Codecs], which is the -// same value the per-upstream proxies install as a gRPC client interceptor. A +// same value each upstream's forwarder installs as a gRPC client interceptor. A // payload therefore transforms identically whichever path it travelled, // because the chain has one construction site rather than two free to drift // apart. diff --git a/internal/codecserver/fx.go b/internal/codecserver/fx.go index 74498b7..90d0d71 100644 --- a/internal/codecserver/fx.go +++ b/internal/codecserver/fx.go @@ -13,26 +13,16 @@ import ( ) // Module provides the codec server and forces its construction, since nothing -// else depends on it and fx would otherwise never build it. Include it -// unconditionally: a disabled codecServer block yields a nil [Server] and no -// lifecycle hook, so the module is inert rather than something the call site -// has to gate. -// -// It depends on the dataplane's payload codec chain rather than building its -// own, which is both what makes a payload transform identically on either path -// and a necessity, since the encryption collectors register once per registry. +// else depends on it. Include it unconditionally: a disabled codecServer block +// yields a nil [Server] and no lifecycle hook, so the module is inert. var Module = fx.Options( fx.Provide(newFromParams), fx.Invoke(func(*Server) {}), ) // Params collects the fx-provided dependencies the codec server needs. Every -// field is required; the codec server has no optional dependency, because a -// missing one would leave it either unauthenticated or unreported. -// -// Codecs comes from the dataplane rather than being built here, so the chain -// has one construction site. Conns is needed only to resolve an -// extension-server authenticator, and is harmlessly empty otherwise. +// field is required. Conns is used only to resolve an extension-server +// authenticator and may be empty otherwise. type Params struct { fx.In Shutdowner fx.Shutdowner @@ -49,11 +39,9 @@ type Params struct { // the fx lifecycle, or returns nil when the codec server is disabled. It warns // rather than fails for the two configurations that are legal but probably // unintended: no authentication, which config only permits on a loopback bind, -// and no encryption keys, which makes both routes identity transforms. -// -// Returns an error when the namespace override mapping is ambiguous, when the -// authenticator cannot be built, or when the TLS material will not load. Each -// is a startup failure rather than a runtime one. +// and no encryption keys, which makes both routes identity transforms. Returns +// an error when the namespace override mapping is ambiguous, when the +// authenticator cannot be built, or when the TLS material will not load. func newFromParams(p Params) (*Server, error) { cfg := &p.Config.CodecServer if !cfg.Enabled { diff --git a/internal/codecserver/handler.go b/internal/codecserver/handler.go index ad1dfb0..82ebae5 100644 --- a/internal/codecserver/handler.go +++ b/internal/codecserver/handler.go @@ -47,9 +47,9 @@ const ( type ( // Codecs transforms payloads on their way to and from an upstream. It is the // subset of [github.com/temporalio/temporal-proxy/internal/proxy.Codecs] the - // handler needs, and an implementation must be the same value the - // per-upstream proxies apply, or a payload will transform differently - // depending on which path it travelled. + // handler needs, and an implementation must be the same value the upstream + // forwarders apply, or a payload will transform differently depending on + // which path it travelled. // // Implementations must be safe for concurrent use, must return the same // number of payloads they were given in the same order, and must not mutate @@ -80,10 +80,7 @@ type ( // needs no translation. Namespaces interface { // Local returns the local namespace name for remote, or remote unchanged - // when it maps no override. Returning the input is deliberate rather - // than a fallback: the vault resolves an unknown namespace to the - // default key policy, which is the correct policy for a namespace that - // has none of its own. + // when it maps no override. Local(remote string) string } @@ -127,19 +124,12 @@ type ( ) // Handler returns the codec server's routes, each served both bare and under a -// namespace path segment because a caller may name the namespace either way. -// Use it when a client reaches a Temporal Service without passing through the -// gateway; a client that goes through the gateway needs no codec server at all. +// namespace path segment. Use it when a client reaches a Temporal Service +// without passing through the gateway. An unknown path answers 404, and a known +// path used with any method but POST answers 405 with an Allow header, where +// the SDK's own handler answers 404. // -// Returns an http.Handler wrapping an http.ServeMux. An unknown path answers -// 404; a known path used with any method but POST answers 405 with an Allow -// header, which diverges from the SDK's own handler answering 404 there. -// -// Panics if r is nil. A missing reporter is a wiring mistake, and failing at -// construction beats a nil dereference on the first request, which would answer -// 500 on a surface that otherwise only ever answers 4xx. -// -// The returned handler is safe for concurrent use. +// Panics if r is nil. The returned handler is safe for concurrent use. func Handler(c Codecs, n Namespaces, r *Reporter, opts ...Option) http.Handler { if r == nil { panic("codecserver: Handler requires a non-nil reporter") @@ -172,13 +162,10 @@ func Handler(c Codecs, n Namespaces, r *Reporter, opts ...Option) http.Handler { return cors(o.origins, o.credentials, mux) } -// WithAuth requires every request to satisfy a, which answers 401 for any -// request a denies. Omitting it serves every request unauthenticated, which is -// only defensible on a loopback bind, and is why configuration rejects an -// enabled codec server bound beyond loopback with no auth block. -// -// A preflight is answered before a reaches it, since a browser sends no -// credentials on one. +// WithAuth requires every request to satisfy a, answering 401 for any request +// a denies. A preflight is answered before a reaches it. Omitting it serves +// every request unauthenticated, which configuration only permits on a +// loopback bind. func WithAuth(a Authenticator) Option { return func(o *options) { o.auth = a } } @@ -202,11 +189,7 @@ func WithLogger(l logger.Logger) Option { } // WithMaxBodyBytes bounds how much of a request body is read, answering 400 -// once a caller exceeds it. The cap matters more here than on an ordinary -// endpoint because this surface unwraps a DEK on demand, so an unbounded body -// is an unbounded amount of work for an unauthenticated-at-the-edge caller. -// -// A value of zero or less keeps the default of 4 MiB. +// once a caller exceeds it. A value of zero or less keeps the default of 4 MiB. func WithMaxBodyBytes(n int64) Option { return func(o *options) { if n > 0 { @@ -216,13 +199,9 @@ func WithMaxBodyBytes(n int64) Option { } // WithNamespaceRequired makes /encode answer 400 when a request names no -// namespace. Set it when a per-namespace key policy is configured: sealing a -// payload that cannot be attributed to a namespace would silently take the -// default policy rather than the one the operator wrote, which is a policy -// violation that no later error reveals. -// -// It never constrains /decode, which does not need a namespace to open a -// payload. +// namespace; it never constrains /decode. Set it when a per-namespace key +// policy is configured, or a payload naming no namespace is silently sealed +// under the default policy. func WithNamespaceRequired(required bool) Option { return func(o *options) { o.namespaceRequired = required } } diff --git a/internal/codecserver/middleware.go b/internal/codecserver/middleware.go index 08ece6b..a3b8b3b 100644 --- a/internal/codecserver/middleware.go +++ b/internal/codecserver/middleware.go @@ -11,18 +11,11 @@ import ( const allowedHeaders = "Content-Type, X-Namespace, Authorization, Authorization-Extras" // cors answers preflight requests and stamps the cross-origin headers onto -// every other response. -// -// It wraps the router rather than sitting inside a route, for two reasons. A -// preflight arrives as OPTIONS, which matches no route and would otherwise -// answer 405, and it has to be answered before authentication because a -// browser sends no credentials on one. -// -// The request's own origin is echoed back, and only when it is on the allowed -// list. A wildcard is never emitted, since a browser rejects one whenever -// credentials are included. -// -// Returns next unchanged when origins is empty, otherwise a handler wrapping it. +// every other response. It must wrap the router: a preflight arrives as +// OPTIONS, matches no route, and carries no credentials, so it is answered +// before routing and authentication. An allowed origin is echoed back; a +// wildcard is never emitted, since a browser rejects one whenever credentials +// are included. Returns next unchanged when origins is empty. func cors(origins []string, credentials bool, next http.Handler) http.Handler { if len(origins) == 0 { return next diff --git a/internal/codecserver/namespaces.go b/internal/codecserver/namespaces.go index 6a2f796..f4a1d91 100644 --- a/internal/codecserver/namespaces.go +++ b/internal/codecserver/namespaces.go @@ -9,34 +9,22 @@ import ( ) // OverrideMap maps a remote namespace name to the local name a per-namespace -// codec policy is keyed by. It satisfies [Namespaces]. -// -// It holds only the namespaces where the answer can differ, which is those -// carrying a per-namespace policy under an upstream that translates names. -// Every other name is left alone, so the map is small and usually empty. +// codec policy is keyed by, holding only namespaces with a policy under an +// upstream that translates names. It satisfies [Namespaces]. // // A nil OverrideMap is usable and translates nothing. An OverrideMap is // read-only after construction and safe for concurrent use. type OverrideMap map[string]string -// NewOverrideMap builds the remote-to-local namespace mapping from -// configuration, once at startup. Call it before serving; it reads cfg and -// keeps no reference to it. -// -// For each upstream that translates namespaces, it walks the namespaces -// carrying a per-namespace codec policy and records the remote name that -// upstream would have produced. Deriving the mapping forwards, through the -// upstream's own rules, is what makes it exact: it honours an explicit -// local-to-remote override, where inverting a translation would be a lossy -// suffix trim. Identity entries are skipped, so an upstream whose rules do not -// change a name contributes nothing. +// NewOverrideMap builds the remote-to-local namespace mapping from cfg, +// recording the remote name each translating upstream produces for every +// namespace carrying a per-namespace codec policy. Identity entries are +// skipped, so an upstream whose rules do not change a name contributes nothing. +// It keeps no reference to cfg. // -// Returns a mapping that may be nil when no namespace carries a policy, which -// callers may use directly. +// Returns a nil mapping, which is usable, when no namespace carries a policy. // Returns an error when a remote name is also a policy key, or when two local -// namespaces produce the same remote name. Both are ambiguous, and an -// ambiguous mapping would seal one namespace's payloads under another's key -// policy, so it fails at startup rather than resolving arbitrarily per request. +// namespaces produce the same remote name, since either mapping is ambiguous. func NewOverrideMap(cfg *config.Config) (OverrideMap, error) { locals := policyNamespaces(cfg) if len(locals) == 0 { @@ -85,15 +73,8 @@ func NewOverrideMap(cfg *config.Config) (OverrideMap, error) { } // Local returns the local namespace name for remote, or remote unchanged when -// no override matches. It implements [Namespaces] and has no error path. -// -// Passing an unknown name through is the correct answer rather than a -// fallback. The vault resolves a namespace it holds no key for to the default -// key policy, which is what a namespace with no override should get, and it -// also means a caller that already speaks local names reaches its own override -// directly without the mapping having to recognise it. -// -// Safe for concurrent use. Safe on a nil receiver. +// no override matches. It implements [Namespaces]. Safe for concurrent use and +// on a nil receiver. func (m OverrideMap) Local(remote string) string { if local, ok := m[remote]; ok { return local diff --git a/internal/codecserver/reporter.go b/internal/codecserver/reporter.go index bddea66..d42991f 100644 --- a/internal/codecserver/reporter.go +++ b/internal/codecserver/reporter.go @@ -9,24 +9,17 @@ import ( ) // Reporter records codec server request telemetry to Prometheus: one count and -// one duration sample per request, on every path including the failures. -// -// A Reporter is safe for concurrent use. Handles are resolved per call through -// WithLabelValues rather than pre-computed, since the label set is small and -// fixed. +// one duration sample per request, on every path including the failures. A +// Reporter is safe for concurrent use. type Reporter struct { requests *prometheus.CounterVec duration *prometheus.HistogramVec } // NewReporter builds the Prometheus-backed request Reporter and registers its -// collectors. Build one per registry; Prometheus rejects a duplicate -// registration, and the factory panics rather than erring on one. f must -// already be scoped to the "codec_server" subsystem, which is what produces the -// published tmprl_proxy_codec_server_* names. -// -// Returns a Reporter publishing requests_total, labelled by route and code, and -// request_duration_seconds, labelled by route. +// requests_total (by route and code) and request_duration_seconds (by route) +// collectors. f must already be scoped to the "codec_server" subsystem. Build +// one per registry: the factory panics on a duplicate registration. func NewReporter(f *metrics.Factory) *Reporter { return &Reporter{ requests: f.NewCounter(prometheus.CounterOpts{ @@ -41,10 +34,8 @@ func NewReporter(f *metrics.Factory) *Reporter { } // Request records one served request, counting it and observing its duration. -// Call it once per request, on every path including the failures, so an error -// rate is derivable from the code label alone. -// -// Safe for concurrent use. +// Call it once per request, on every path including the failures. Safe for +// concurrent use. func (r *Reporter) Request(route string, code int, seconds float64) { r.requests.WithLabelValues(route, strconv.Itoa(code)).Inc() r.duration.WithLabelValues(route).Observe(seconds) diff --git a/internal/codecserver/server.go b/internal/codecserver/server.go index bc34833..777b6e0 100644 --- a/internal/codecserver/server.go +++ b/internal/codecserver/server.go @@ -55,11 +55,8 @@ func NewServer( } } -// Addr is the address the server is accepting on, or nil before [Server.Start] -// returns successfully. It is readable at all because the listener is bound -// explicitly rather than by ListenAndServe, which never reports the port it -// chose; that is what lets a caller bind port zero and still find the server. -// +// Addr is the address the server is accepting on, including the port chosen +// for a port-zero hostPort, or nil before [Server.Start] returns successfully. // Safe for concurrent use. func (s *Server) Addr() net.Addr { s.mu.Lock() @@ -68,13 +65,10 @@ func (s *Server) Addr() net.Addr { return s.addr } -// Start binds the listener and serves in a background goroutine, returning as -// soon as the listener is accepting. Serving continues until [Server.Stop], -// so Start does not block. -// -// Returns an error only if the bind fails, typically an address already in use. -// A failure after Start returns cannot be reported through it, so an unexpected -// stop reaches the abort function given to [NewServer] instead. +// Start binds the listener and serves in a background goroutine until +// [Server.Stop], returning once the listener is accepting. It returns an error +// only if the bind fails; a serving failure after that reaches the abort +// function given to [NewServer] instead. // // Call it at most once. A Server is not restartable after [Server.Stop]. func (s *Server) Start(ctx context.Context) error { @@ -108,11 +102,9 @@ func (s *Server) Start(ctx context.Context) error { return nil } -// Stop closes the listener and waits for in-flight requests to finish. A -// graceful stop is not an error, so the serving goroutine's abort function is -// not called. -// -// Returns ctx's error if the drain does not finish in time, nil otherwise. +// Stop closes the listener and waits for in-flight requests to finish, without +// calling the abort function. Returns ctx's error if the drain does not finish +// in time, nil otherwise. func (s *Server) Stop(ctx context.Context) error { s.logger.Info("Shutting down the codec server") diff --git a/internal/config/auth.go b/internal/config/auth.go index f5c478a..5d850d4 100644 --- a/internal/config/auth.go +++ b/internal/config/auth.go @@ -40,17 +40,14 @@ type ( // Name selects which configured extension server to ask. CredentialHeaders // names the metadata headers carrying the caller's credentials, which the // proxy lifts into the request it sends that server and removes from the - // stream it forwards upstream. It has to be declared because a verdict - // reports only admit-or-deny, so nothing in the exchange reveals which - // headers mattered. + // stream it forwards upstream. // - // Leaving it empty does not hide the caller's credentials from the server. The - // proxy forwards the caller's metadata on the call either way, so the server - // still sees whatever headers the caller sent; what it loses is the request - // field naming them, so it has to know which metadata to read and cannot tell - // a header this proxy vouches for from any other. Nothing is stripped before - // proxying upstream either, so the caller's credential continues to the - // upstream alongside any credential configured for it. + // Leaving CredentialHeaders empty does not hide the caller's credentials + // from the server: it still receives the caller's metadata, but no request + // field names the credential headers, so it must know which metadata to + // read. Nothing is stripped before forwarding either, so the caller's + // credential continues to the upstream alongside any credential configured + // for it. ExternalAuthConfig struct { Name string `yaml:"name"` CredentialHeaders []string `yaml:"credentialHeaders"` @@ -100,13 +97,10 @@ func (a *AuthConfig) Validate() error { // referentialRules checks that external authentication names a configured // extension server, given the set of known names. A failure is stamped with the -// referring field's YAML path so it lands on "auth.external"/"name". -// -// The rule is appended at the Config level rather than composed under Validate -// because it needs the full set of extension server names, which is only known -// there. A nil receiver or a blank name yields nothing: the former means no auth -// block at all, and the latter is already reported as required by -// [ExternalAuthConfig.Validate], which leaves this rule no server to name. +// referring field's YAML path so it lands on "auth.external"/"name". It is +// appended at the Config level, where the full set of names is known. A nil +// receiver or a blank name yields nothing; a blank name is already reported as +// required by [ExternalAuthConfig.Validate]. func (a *AuthConfig) referentialRules(known map[string]struct{}) []validation.Rule { if a == nil || a.External == nil || a.External.Name == "" { return nil diff --git a/internal/config/cloudapi.go b/internal/config/cloudapi.go index d308482..06cfd00 100644 --- a/internal/config/cloudapi.go +++ b/internal/config/cloudapi.go @@ -6,76 +6,46 @@ import ( ) // APITranslations configures rewriting a method an upstream does not serve into -// the one that does. It is optional and usually absent, and it carries no switch: -// whether a method is translated is derived from the rest of the configuration -// rather than declared. -// -// What derives it is [Upstream.IsCloud]: every Cloud upstream gets translation, -// and a method only fires where routing sends it. A namespace-less method lands -// on the upstream serving namespace-less requests, so an operator who wants the -// untranslated failure back for those routes them at a Temporal Service that -// serves them, which is the same statement made where it belongs. -// -// This block exists for the two things detection cannot know: which Cloud -// environment the control plane lives in, and the API key an mTLS upstream has -// none of to inherit. -// -// The zero value is the block an operator did not write, which is what almost -// every configuration has, and it answers for the default - so nothing here is a -// pointer and no caller has to check before asking. The same holds for [CloudAPI]. +// the one that does. It is optional and carries no switch: every upstream +// [Upstream.IsCloud] recognizes gets translation, applied only where routing +// sends the method, so a namespace-less method is translated only when the +// upstream serving namespace-less requests is a Cloud one. The block holds what +// detection cannot know: which Cloud environment the control plane lives in, +// and an API key for an mTLS upstream. The zero value is the absent block and +// yields the defaults; the same holds for [CloudAPI]. type APITranslations struct { CloudAPI CloudAPI `yaml:"cloudApi"` } // CloudAPI overrides how the proxy reaches Temporal Cloud's control plane, which -// answers the methods Cloud does not serve on a namespace frontend. -// -// The block is optional and usually absent. An upstream that [Upstream.IsCloud] -// recognizes gets method translation on its own, over a connection to -// [cloud.APIHostPort] carrying that upstream's credentials - the same API key -// authorizes both, so there is nothing more to say. +// answers the methods Cloud does not serve on a namespace frontend. It is +// optional: an upstream [Upstream.IsCloud] recognizes gets method translation +// over a connection to [cloud.APIHostPort] carrying that upstream's +// credentials. // -// It is required in one case. The Cloud Ops API accepts an API key only; unlike a -// namespace frontend it does not accept mTLS. An upstream authenticating with a -// client certificate therefore has no credential to inherit, and must name an API -// key here or its translated methods are refused. Beyond that, configure this -// only to reach a different Cloud environment. +// It is required when that upstream authenticates with a client certificate: +// the Cloud Ops API accepts only an API key, so such an upstream must name one +// here or its translated methods are refused. Otherwise configure it only to +// reach a different Cloud environment. See https://docs.temporal.io/ops. // -// See https://docs.temporal.io/ops. -// -// It is deliberately not an entry in Upstreams: an upstream is a forwarder and -// a routing destination, and the control plane is only ever a client -// connection. Declaring it there would give it a forwarder it cannot use and -// routability it should not have. +// It is not an entry in upstreams: the control plane is only ever a client +// connection, never a forwarder or a routing destination. type CloudAPI struct { Listen ListenConfig `yaml:",inline"` Credentials *CredentialConfig `yaml:"credentials"` } // Upstream renders the control plane as an [Upstream] for the Cloud upstream -// src, so its connection is dialled by the same resolver, TLS, and credential -// machinery as any other rather than by a second code path. The zero value is -// the unconfigured case and yields the inherited defaults, so callers need not +// src, dialled by the same resolver, TLS, and credential machinery as any +// other. The zero value yields the inherited defaults, so callers need not // branch on whether the block is present. // -// Credentials are inherited from src because a Temporal Cloud API key authorizes -// the control plane as well as the frontend. TLS is not: the control plane is a -// different host, so src's server name or client certificate would not apply to -// it, and the dial default stands instead - verification against the system root -// pool, which is what the real control plane presents. The connection settings -// are inherited too, so one block tunes both connections to the same account, -// though the dataplane keeps this one to a single connection whatever -// maxConnections says, since it carries no long polls. -// -// When this block is present its tls and insecure are authoritative, the same way -// they are on an upstream: an absent tls still verifies against the system roots, -// and plaintext has to be asked for. That is only reachable for a control plane -// with no credentials, since Validate rejects credentials on an insecure hop, so -// a key still cannot be sent in the clear. -// -// The name is derived from src rather than fixed, so two Cloud upstreams with -// different credentials get distinct connections instead of sharing whichever -// was dialled first. +// Credentials and connection settings are inherited from src, though the +// dataplane keeps this to a single connection whatever maxConnections says. +// TLS is not inherited: the block's tls and insecure are authoritative, an +// absent tls verifies against the system roots, and Validate rejects +// credentials on an insecure hop. The name is derived from src, so Cloud +// upstreams with different credentials get distinct connections. func (c CloudAPI) Upstream(src *Upstream) *Upstream { up := &Upstream{ Name: src.Name + "/cloud-api", @@ -111,17 +81,11 @@ func (c CloudAPI) IsSaasAPI() bool { return cloud.IsEndpoint(c.Listen.HostPort) } -// Validate checks the control plane as it will actually be dialled, by -// validating the [Upstream] it renders to. That covers the same ground as any -// upstream - dial target, outbound TLS, credentials, and credentials requiring -// TLS - without restating the rules, and checks the effective configuration -// (the defaulted address included) rather than only the fields an operator -// supplied. -// -// An address that is not a Cloud endpoint is not rejected. Nothing but a Cloud -// deployment serves CloudService, but a test double or a private environment -// legitimately does not carry the Cloud domain, and the proxy has no way to tell -// that apart from a typo. It is logged as a warning at startup instead. +// Validate checks the control plane as it will actually be dialled, defaulted +// address included, by validating the [Upstream] it renders to. An address +// that is not a Cloud endpoint is not rejected, since a test double or private +// environment may legitimately use one; it is logged as a warning at startup +// instead. func (c CloudAPI) Validate() error { return c.Upstream(&Upstream{Name: "cloudApi"}).Validate() } diff --git a/internal/config/encryption.go b/internal/config/encryption.go index 971391d..20e07f9 100644 --- a/internal/config/encryption.go +++ b/internal/config/encryption.go @@ -50,13 +50,6 @@ type ( // DEKCacheSize is the DEK cache size to apply: the configured size when there is // one, and [crypto.DefaultCacheSize] when the field is absent. Zero disables the // cache, so every Open unwraps its DEK through the KEK. -// -// The distinction is why the field is a pointer. Zero is a meaningful value here -// and a plain int cannot tell an operator who wrote nothing from one who wrote -// zero - so an absent field would read as "disable the cache" and silently -// override the vault's own default, turning every payload the proxy opens into a -// KMS round trip. Absent means "no opinion", and disabling the cache has to be -// written down. func (e *Encryption) DEKCacheSize() int { if e == nil || e.CacheSize == nil { return crypto.DefaultCacheSize @@ -114,10 +107,8 @@ func (p *KeyPolicy) Validate() error { // the referring policy's YAML path so it lands on the right key (e.g. // "encryption.default"/"uri" or "encryption.overrides[payments]"/"decryptURIs[1]"). // Non-extension URIs are skipped; their scheme is already checked by validKeyURI. -// -// The rules are appended at the Config level rather than composed under -// Encryption.Validate because they need the full set of extension server names, -// which is only known there. +// The rules are appended at the Config level, where the full set of names is +// known. func (e *Encryption) referentialRules(known map[string]struct{}) []validation.Rule { var rules []validation.Rule @@ -174,11 +165,8 @@ func validKeyURI() validation.Check[url.URL] { // validKeyURIRef rejects a key URI whose scheme is not one of the supported KMS // providers, and an extension URI that names no server. The scheme match is -// case-insensitive. -// -// The empty-host case is checked here rather than in referentialRules because -// that rule reports a host that matches no configured server, and a missing host -// gives it no name to report. +// case-insensitive. The empty-host case is checked here because +// referentialRules reports unknown server names and a missing host has none. func validKeyURIRef() validation.Check[*url.URL] { return func(u *url.URL) error { if !slices.Contains(validKeySchemes, strings.ToLower(u.Scheme)) { diff --git a/internal/config/extensions.go b/internal/config/extensions.go index 3e202b6..f6a3c8b 100644 --- a/internal/config/extensions.go +++ b/internal/config/extensions.go @@ -15,14 +15,11 @@ type ( // a backend the proxy does not support natively, such as an on-prem HSM, an // internal key service, or an in-house identity system. // - // Name identifies the server within the configuration so other blocks can - // reference it, and must be unique across the list. Credentials, when set, - // attach per-request credentials to the outbound calls and require TLS, - // since sending them over a plaintext connection would expose them on the wire. - // - // Unlike Upstream, an extension server is dialed at a fixed address rather - // than resolved per request, so a templated hostPort is rejected outright - // instead of being deferred to request time. + // Name identifies the server so other blocks can reference it, and must be + // unique across the list. Credentials, when set, attach per-request + // credentials to the outbound calls and require TLS. Unlike an Upstream, an + // extension server is dialed at a fixed address, so a templated hostPort is + // rejected. ExtensionServer struct { Name string `yaml:"name"` Listen ListenConfig `yaml:",inline"` diff --git a/internal/config/routing.go b/internal/config/routing.go index c5c7416..0cae169 100644 --- a/internal/config/routing.go +++ b/internal/config/routing.go @@ -50,10 +50,8 @@ func (r *Routing) Validate() error { // YAML path so it lands on the right key (e.g. "routing.rules[0]"/"upstream"). // Empty references are skipped: default and system are optional, and a rule's // missing upstream is already reported as required by RoutingRule.Validate. -// -// The rules are appended at the Config level rather than composed under -// Routing.Validate because they need the full set of upstream names, which is -// only known there. +// The rules are appended at the Config level, where the full set of names is +// known. func (r *Routing) referentialRules(known map[string]struct{}) []validation.Rule { check := knownUpstream(known) ref := func(subject, field, name string) validation.Rule { diff --git a/internal/dataplane/dataplane.go b/internal/dataplane/dataplane.go index c64afa1..c32c648 100644 --- a/internal/dataplane/dataplane.go +++ b/internal/dataplane/dataplane.go @@ -432,12 +432,9 @@ func translates(cfg *config.Config) bool { // cloudAPIConn builds the connection translated methods for up are answered // over, along with the translations that use it. The connection is dialled by // the same resolver and credential machinery as an upstream but is not -// registered with the router: nothing routes to it, and it is reached only by a -// translation. -// -// It is not opened eagerly. Translation is incidental to an upstream's normal -// traffic, so a control plane that is unreachable must not stop the proxy -// serving everything else. +// registered with the router; only a translation reaches it. It is not opened +// eagerly, so an unreachable control plane does not stop the proxy serving +// everything else. func cloudAPIConn(cfg *config.Config, o *options, up *config.Upstream) (*translation.Registry, *connect.Conn, error) { reg, err := translation.Default() if err != nil { diff --git a/internal/kms/fx.go b/internal/kms/fx.go index 05bc4fb..67dad2f 100644 --- a/internal/kms/fx.go +++ b/internal/kms/fx.go @@ -148,13 +148,8 @@ func runRotation(ctx context.Context, v vaultRefresher, interval time.Duration, // createVault builds a vault from the registry, applying the configured cache // size and, when a default key policy is set, its DEK duration and renewal lead -// time. -// -// A disabled cache is logged rather than left to be inferred. It is a legitimate -// choice, but an expensive one - every Open becomes a KEK round trip - and the -// cache metrics cannot report it: hits and misses both sit at zero whether the -// cache is off or merely idle, so this line is the only thing that distinguishes -// them. +// time. A disabled cache is logged, since every Open then costs a KEK round +// trip and the cache metrics read zero whether the cache is off or merely idle. func createVault( c *config.Config, r *crypto.KEKRegistry, diff --git a/internal/kms/reporter.go b/internal/kms/reporter.go index 76423e4..1c8e2fa 100644 --- a/internal/kms/reporter.go +++ b/internal/kms/reporter.go @@ -56,12 +56,10 @@ type ( // NewReporter builds the Prometheus-backed encryption Reporter, pre-resolving // the meaningful KEK label combinations so every series starts at zero. f must -// already be scoped to the "encryption" subsystem by the caller. -// -// Prometheus panics rather than erring on a collector it will not accept, so -// recover and return an error: a configured fixed label can name one of these -// series' own labels, and config cannot refuse that without knowing every -// collector's label set. +// already be scoped to the "encryption" subsystem by the caller. Returns an +// error rather than panicking when a collector cannot be registered, such as a +// second reporter on the registry or a configured fixed label colliding with +// one of a collector's own labels. func NewReporter(f *metrics.Factory) (rep *Reporter, err error) { defer func() { rec := recover() @@ -197,17 +195,12 @@ func (r *Reporter) Observe(e crypto.Event) { // envelopeOp records the AES-256-GCM portion of one envelope operation: its // duration and, via dek_ops_total, its own result. // -// Total and Namespace are deliberately unused. internal/proxy already records -// the end-to-end duration and operation counts, labeled by namespace, around -// its own Seal and Open calls as vault_ops_duration_seconds and vault_ops_total; -// recording them here would duplicate those series and collide with them on -// the shared "encryption" subsystem. Err is likewise unused here for a -// reason, not an oversight: the envelope result already lives on -// internal/proxy's vault_ops_total. The result label instead comes from -// CryptoErr, the AES step's own outcome, deliberately not Err: a Seal that -// encrypts successfully and then fails to wrap its DEK is a KEK failure that -// kek_ops_total already reports, and counting it here would blame the wrong -// actor. +// Total, Namespace, and Err are deliberately unused: internal/proxy records the +// end-to-end duration and result as vault_ops_duration_seconds and +// vault_ops_total, and recording them here would duplicate those series on the +// shared "encryption" subsystem. The result label comes from CryptoErr, so a +// Seal whose DEK wrap fails is reported by kek_ops_total alone rather than +// blamed on the AES step. func (r *Reporter) envelopeOp(e crypto.EnvelopeEvent) { // No AES step, no DEK operation to record. CryptoAttempted is the only sound // test for that: a zero Crypto cannot distinguish a step that never ran from diff --git a/internal/proxy/version.go b/internal/proxy/version.go index 54dbe6b..32349dc 100644 --- a/internal/proxy/version.go +++ b/internal/proxy/version.go @@ -11,11 +11,8 @@ import ( // VersionDialOptions returns the dial options that stamp the proxy's own build // version on every outbound request as meta.VersionHeader. Callers fold them // into the dial options for the upstream connection, and only for an upstream -// that is Temporal Cloud: Cloud reads the header to tell which proxy build a -// request came from, and no other upstream has asked for it. -// -// An empty version installs nothing, so a build with no version to report sends -// no header rather than an empty one. +// that is Temporal Cloud. An empty version installs nothing, so no header is +// sent rather than an empty one. func VersionDialOptions(version string) []grpc.DialOption { if version == "" { return nil diff --git a/internal/router/reporter.go b/internal/router/reporter.go index 8720b1c..d03eb96 100644 --- a/internal/router/reporter.go +++ b/internal/router/reporter.go @@ -21,15 +21,12 @@ const ( type ( // Reporter records router telemetry to Prometheus: routing decisions and the // forwarding failures the router itself originates. It pre-resolves a counter - // for every meaningful (upstream, outcome) and (upstream, reason) combination - // so the emit path is a lock-free map read; an unexpected label combination - // falls back to CounterVec.WithLabelValues. A Reporter is safe for concurrent - // use. + // for every meaningful (upstream, outcome) and (upstream, reason) combination, + // so each series starts at zero. A Reporter is safe for concurrent use. // - // Configured metadata labels suppress that pre-resolution, because their - // values arrive with a request and cannot be enumerated at startup. Every - // emit then takes the fallback, and no series starts at zero, so a query for - // a counter that has not been incremented yet finds nothing rather than 0. + // Configured metadata labels suppress that pre-resolution, since their values + // arrive with a request. No series then starts at zero, so a query for a + // counter that has not been incremented yet finds nothing rather than 0. Reporter struct { decisions *prometheus.CounterVec errors *prometheus.CounterVec @@ -53,12 +50,7 @@ type ( // with the factory's registry and pre-resolving the meaningful label // combinations so every series starts at zero. upstreams is the configured // upstream name list. labels are the configured metadata labels, and may be -// the zero value. -// -// Nothing is pre-resolved when metadata labels are configured: a handle would -// have to pin their values, which only a request carries. Leaving the maps -// empty routes every emit through the fallback the maps exist to avoid, rather -// than adding a second path that could drift from it. +// the zero value; when any are configured, nothing is pre-resolved. func NewReporter(f *metrics.Factory, upstreams []string, labels metrics.MetadataLabels) *Reporter { decisions := f.NewCounter(prometheus.CounterOpts{ Name: "decisions_total", diff --git a/internal/server/loopback.go b/internal/server/loopback.go index 48b908a..95a469f 100644 --- a/internal/server/loopback.go +++ b/internal/server/loopback.go @@ -35,25 +35,13 @@ type ( // grpc.health.v1.Health/Watch on the server it belongs to, reads one message, // and reports whether the exchange was answered at all. // - // Watch rather than Check, and that is the whole point: a server assembled - // here carries stream interceptors and no unary ones, so the unary Check a - // probe runs reaches the health handler without traversing the chain every - // forwarded request goes through. Watch is a locally registered streaming - // method, so a call to it runs that chain, and a wedge in it stops being - // invisible. - // - // The call never leaves the process. It goes over an in-process listener - // served by [newLoopbackServer]'s twin of the real server: the same - // interceptor chain and the same health service, but with no address to - // learn, no handshake to complete and no credential to present. That is what - // lets the check behave the same whether the server is plaintext, TLS, or - // mutual TLS, which last would otherwise leave it dialling a listener that - // demands a client certificate it has no way to hold. - // - // What it does not cover is anything below the chain: accepting on the real - // listener, and the handshake above that. Those belong to whatever probes the - // server from outside, over a tcpSocket or grpc handler, rather than to - // this. + // It must use Watch, not Check: the server carries stream interceptors and + // no unary ones, so only a streaming call runs the chain every forwarded + // request goes through. The call goes over an in-process listener served by + // [newLoopbackServer], with no handshake or credential, so it behaves the + // same whether the real server is plaintext, TLS, or mutual TLS. It does not + // cover accepting on the real listener or the handshake; those belong to an + // external probe. loopbackCheck struct { interval time.Duration timeout time.Duration @@ -64,13 +52,10 @@ type ( // newLoopbackServer builds the server behind the in-process listener, from the // same options as the real one so that it runs the same interceptor chain, and -// registering the same health service so both answer from one status. -// -// It is a second server rather than a second listener on the first because -// transport credentials belong to a [grpc.Server]: one server cannot accept on a -// listener that requires client certificates and on another that requires -// nothing. Nothing outside the process holds the listener, so the credential-free -// half is reachable only by the check. +// registering the same health service so both answer from one status. It is a +// separate server because transport credentials belong to a [grpc.Server], and +// this one requires none; nothing outside the process holds its listener, so +// only the check can reach it. func newLoopbackServer(o *options, hc *health.Server) *grpc.Server { svr := grpc.NewServer(o.serverOptions(grpc.Creds(insecure.NewCredentials()))...) grpc_health_v1.RegisterHealthServer(svr, hc) @@ -95,15 +80,10 @@ func newLoopbackCheck(lis *bufconn.Listener, interval, timeout time.Duration, lo func (c *loopbackCheck) Interval() time.Duration { return c.interval } // Status runs one exchange and maps its outcome. Any answer, including a -// rejection, means the chain ran end to end and reports SERVING: that is what -// keeps the check free of credentials, since it never needs to authenticate, -// only to be answered. A deadline means the chain did not answer, and a listener -// that is not accepting or a failure carrying no status at all means the server -// is not answering either; both report NOT_SERVING. -// -// No hysteresis is applied here. A probe's failureThreshold already debounces, -// and a second threshold underneath it makes the real detection latency hard to -// reason about. +// rejection, means the chain ran end to end and reports SERVING, so the check +// never needs to authenticate. A deadline, a listener that is not accepting, or +// a failure carrying no status at all reports NOT_SERVING. No hysteresis is +// applied; debouncing is left to the probe's failureThreshold. func (c *loopbackCheck) Status(ctx context.Context) grpc_health_v1.HealthCheckResponse_ServingStatus { ctx, cancel := context.WithTimeout(ctx, c.timeout) defer cancel() @@ -123,14 +103,9 @@ func (c *loopbackCheck) Status(ctx context.Context) grpc_health_v1.HealthCheckRe } // probe opens a connection over the in-process listener and runs one exchange -// across it. Dialling is the only thing it adds over watchOnce, which is what -// keeps watchOnce drivable from a test with no server behind it. -// -// The connection is built and closed per run rather than held open. A check runs -// every 30 seconds by default, and an in-process connection costs far less than -// one kept correct across the server's whole lifetime; dialling afresh also -// covers accepting, which a reused connection would skip, so a Serve goroutine -// that had stopped would otherwise go unnoticed. +// across it. The connection is built and closed per run so that each run also +// covers accepting; a reused connection would not notice a stopped Serve +// goroutine. func (c *loopbackCheck) probe(ctx context.Context) error { conn, err := grpc.NewClient( loopbackName, @@ -148,13 +123,9 @@ func (c *loopbackCheck) probe(ctx context.Context) error { } // watchOnce opens a Watch stream, reads one message, and returns what the -// exchange ended with. -// -// The stream is cancelled as soon as that message arrives: health.Server -// registers a watcher per Watch stream, so one left open would leak an entry -// every interval. It takes the client rather than a listener so a test can drive -// it with a fake and assert that cancellation, which is otherwise unobservable -// from outside the process. +// exchange ended with. The stream is cancelled as soon as that message arrives: +// health.Server registers a watcher per Watch stream, so one left open would +// leak an entry every interval. func watchOnce(ctx context.Context, client grpc_health_v1.HealthClient) error { ctx, cancel := context.WithCancel(ctx) defer cancel() diff --git a/internal/server/reporter.go b/internal/server/reporter.go index 8f71449..0d93ac0 100644 --- a/internal/server/reporter.go +++ b/internal/server/reporter.go @@ -18,19 +18,15 @@ import ( var durationBuckets = []float64{.005, .01, .025, .05, .1, .25, .5, 1, 2.5, 5, 10, 30, 60, 120} // Reporter records server-layer telemetry to Prometheus: per-RPC latency and -// completed-request counts by gRPC status code, both labeled by method. The -// method label set is not known at startup, so handles are resolved per call -// via WithLabelValues rather than pre-resolved. A Reporter is safe for -// concurrent use. +// completed-request counts by gRPC status code, both labeled by method. A +// Reporter is safe for concurrent use. // -// Cardinality assumption: method comes from the request line, and the proxy -// serves every request through a catch-all handler, so any distinct method -// string a client sends becomes a new series. This is bounded only for trusted -// callers (real Temporal SDK clients use a fixed method set); a client sending -// arbitrary method paths can grow the series set without bound. The proxy -// therefore assumes trusted callers and must not be exposed directly to -// untrusted clients without first bounding this label. namespace is never a -// label for the same reason. +// The method label comes from the request line, and every request goes through +// a catch-all handler, so each distinct method a client sends becomes a new +// series. Temporal SDK clients use a fixed set, but arbitrary method paths grow +// the series set without bound: the proxy must not be exposed to untrusted +// clients without first bounding this label. Namespace is never a label, for +// the same reason. type Reporter struct { duration *prometheus.HistogramVec requests *prometheus.CounterVec diff --git a/internal/server/server.go b/internal/server/server.go index 2768d5a..5f914ad 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -181,12 +181,9 @@ func WithHealthCheck(hc HealthCheck) Option { // WithLoopbackHealthCheck drives the serving status from a call the server makes // to its own Health/Watch, so the status reports whether a request can still // travel the stream interceptor chain rather than only whether the process -// accepts connections. See [loopbackCheck] for why it is Watch and why the call -// does not leave the process. -// -// interval is how often the check runs and timeout bounds one run; neither is -// defaulted here. It takes precedence over [WithHealthCheck] regardless of the -// order the two are supplied in. +// accepts connections. interval is how often the check runs and timeout bounds +// one run; neither is defaulted here. It takes precedence over +// [WithHealthCheck] regardless of the order the two are supplied in. func WithLoopbackHealthCheck(interval, timeout time.Duration) Option { return optFunc(func(o *options) { o.loopback = &loopbackTimings{interval: interval, timeout: timeout} diff --git a/internal/transport/connect/conn.go b/internal/transport/connect/conn.go index d6b5f71..8d73dd0 100644 --- a/internal/transport/connect/conn.go +++ b/internal/transport/connect/conn.go @@ -115,16 +115,11 @@ func WithConnections(n int) ConnOption { } // WaitReady opens conns and blocks until each is ready or ctx is done, -// whichever comes first. They are waited on concurrently, so they share ctx's -// deadline rather than consuming it in turn, and every target that never came up -// is reported rather than only the first, so one unreachable address cannot mask -// another. -// -// When ctx has a deadline this stops just short of it. Callers are fx start -// hooks, and fx prefers its start context's error over what a hook returns -// (app.go, withTimeout), so a wait that runs to the deadline is reported as a -// bare "context deadline exceeded" and the target names are lost. Returning -// early is what keeps them. +// whichever comes first. They are waited on concurrently and share ctx's +// deadline, and every target that never came up is reported, not only the +// first. When ctx has a deadline this returns just short of it, so an fx start +// hook reports the target names rather than fx's bare "context deadline +// exceeded". func WaitReady(ctx context.Context, conns ...*Conn) error { if deadline, ok := ctx.Deadline(); ok && time.Until(deadline) > 2*readyMargin { var cancel context.CancelFunc @@ -173,19 +168,16 @@ func (c *Conn) NewStream( } // WaitReady opens the underlying connections and blocks until each is ready or -// ctx is done, whichever comes first. [NewConn] creates a static Conn's connection -// but grpc.NewClient only dials on demand, so nothing is open until this runs (or -// the first request arrives); this is what makes the connection real ahead of -// serving traffic. +// ctx is done, whichever comes first. A static Conn dials on demand, so nothing +// is open until this runs or the first request arrives. // // A refused connection is not on its own fatal: gRPC retries with backoff, so a // target still coming up passes as long as it answers before ctx expires. gRPC // keeps the underlying dial error private, so one that never answers is reported // by the state it was stuck in, wrapping ctx's error. // -// A dynamic Conn holds no connection until a request resolves one, so there is -// nothing to open and this does nothing. Callers can pass a mixed set of conns -// without sorting them first. +// A dynamic Conn holds no connection until a request resolves one, so this does +// nothing for it and callers can pass a mixed set of conns. func (c *Conn) WaitReady(ctx context.Context) error { if !c.resolver.IsStatic() { return nil diff --git a/internal/transport/connect/fx.go b/internal/transport/connect/fx.go index 7fff153..735cf64 100644 --- a/internal/transport/connect/fx.go +++ b/internal/transport/connect/fx.go @@ -3,13 +3,9 @@ package connect import "go.uber.org/fx" // Module provides a *Pool and binds its lifecycle to the application, closing -// every pooled connection on shutdown via an fx stop hook. -// -// The pool is deliberately not opened as a whole on start. It also holds -// connections that are lazy on purpose, such as a templated upstream's, created -// per request, and the Cloud control plane's, opened on first use; waiting on -// those here would hold up startup for connections nothing needs yet. Opening -// eager connections is the job of whoever owns one, through [WaitReady]. +// every pooled connection on shutdown via an fx stop hook. It does not open the +// pool on start; opening an eager connection is the job of whoever owns it, +// through [WaitReady]. var Module = fx.Options( fx.Provide(NewPool), fx.Invoke(func(p *Pool, lc fx.Lifecycle) { diff --git a/pkg/crypto/keys.go b/pkg/crypto/keys.go index 26fff4a..5c9c74c 100644 --- a/pkg/crypto/keys.go +++ b/pkg/crypto/keys.go @@ -23,12 +23,8 @@ type ( // KeyFactory opens [KEK]s from key URIs, choosing an opener by URI scheme. // It handles the cloud KMS schemes out of the box (see [NewKeyFactory]) and // takes additional or replacement schemes through - // [WithKeyFactoryFuncForScheme], so a caller can serve keys from its own key - // store without the code consuming those KEKs knowing where they come from. - // - // Schemes are registered during construction only, so a KeyFactory never - // changes afterwards and one may be shared by any number of goroutines opening - // keys at once. + // [WithKeyFactoryFuncForScheme]. Schemes are registered during construction + // only, so a KeyFactory is safe for concurrent use. KeyFactory struct { funcs map[string]KeyFactoryFunc } @@ -45,34 +41,25 @@ type ( keyFactoryOpt func(*KeyFactory) - // CloudKey is a [KEK] backed by a cloud KMS key. The embedded - // [secrets.Keeper] supplies Decrypt and Close; CloudKey adds the ID and the - // namespace-aware Encrypt that KEK requires. - // - // A Keeper addresses a single fixed key, so one CloudKey wraps DEKs for every - // namespace it is handed. Open one CloudKey per KMS key and let a - // [KEKRegistry] decide which namespaces map onto which key. + // CloudKey is a [KEK] backed by a single cloud KMS key, with Decrypt and + // Close supplied by the embedded [secrets.Keeper]. It wraps DEKs for every + // namespace it is handed with that one key; use a [KEKRegistry] to map + // namespaces onto keys. CloudKey struct { *secrets.Keeper id string } ) -// NewKeyFactory returns a KeyFactory that opens cloud KMS keys, then applies -// opts in order. The schemes registered up front are [DefaultSchemes], all -// served by [NewCloudKey]: +// NewKeyFactory returns a KeyFactory that serves [DefaultSchemes] with +// [NewCloudKey], then applies opts in order, so [WithKeyFactoryFuncForScheme] +// can replace any default as well as add schemes. Importing this package links +// in the driver for each default scheme: // // awskms:// AWS KMS // azurekeyvault:// Azure Key Vault // gcpkms:// Google Cloud KMS // testing:// a local in-process key, for tests and local runs only -// -// The driver behind each scheme is linked in by importing this package, so no -// further imports are needed to use them. -// -// opts are applied after those defaults, which means -// [WithKeyFactoryFuncForScheme] can replace any of them as well as add schemes -// of its own. func NewKeyFactory(opts ...KeyFactoryOption) *KeyFactory { schemes := DefaultSchemes() funcs := make(map[string]KeyFactoryFunc, len(schemes)) @@ -90,19 +77,14 @@ func NewKeyFactory(opts ...KeyFactoryOption) *KeyFactory { // NewCloudKey opens the cloud KMS key addressed by uri, which must use one of // the schemes listed in [NewKeyFactory]. Close the returned key when it is no -// longer needed; a [KEKRegistry] does that for the keys it holds. +// longer needed; a [KEKRegistry] does that for the keys it holds. The key's ID +// is uri, except for the testing scheme. // -// The "testing://" scheme is rewritten to gocloud's "base64key://" local keeper -// so tests and local runs need no cloud KMS at all. Everything after that scheme -// is the base64-encoded 32-byte key; pass a bare "testing://" to get a random -// one. Key material is kept out of the errors this function returns, but it -// still reaches the key's ID, and therefore every DEK the key wraps, which is -// one more reason to keep the scheme away from production. -// -// For every other scheme the ID is just uri, so it is stable across processes -// and identifies the key again on the decrypt path. Schemes are matched without -// regard to case, so the ID of a testing key is the same however its scheme was -// spelled. +// "testing://", matched without regard to case, opens gocloud's local +// "base64key://" keeper with the ID "base64key://" plus everything after the +// scheme: the base64-encoded 32-byte key, or nothing for a random one. Key +// material is kept out of returned errors but reaches the ID, and therefore +// every DEK the key wraps, so keep the scheme away from production. func NewCloudKey(ctx context.Context, uri string) (KEK, error) { open, material := uri, "" if after, ok := cutPrefixFold(uri, testingScheme); ok { @@ -143,13 +125,9 @@ func WithKeyFactoryFuncForScheme(scheme string, fn KeyFactoryFunc) KeyFactoryOpt } // DefaultSchemes lists the URI schemes a [KeyFactory] serves with [NewCloudKey] -// before any option is applied, lowercased as [KeyFactory.Create] matches them. -// It is useful for validating a key URI ahead of opening it, so a typo in a -// scheme can be reported alongside the rest of a config rather than at the point -// the key is first needed. -// -// Each call returns a fresh slice; the caller may sort or filter it freely -// without disturbing the factory. +// before any option is applied, lowercased as [KeyFactory.Create] matches them, +// for validating a key URI ahead of opening it. Each call returns a fresh slice +// the caller may modify. func DefaultSchemes() []string { return []string{ "awskms", diff --git a/pkg/crypto/vault.go b/pkg/crypto/vault.go index 409f1aa..233e515 100644 --- a/pkg/crypto/vault.go +++ b/pkg/crypto/vault.go @@ -75,12 +75,8 @@ type ( // cryptoStep records whether the AES-256-GCM step ran, what it cost, and // whether it failed. Its zero value means the step was never reached, which - // is what an early failure returns. - // - // ran is tracked explicitly rather than inferred from dur being nonzero, - // because a small payload can complete inside the clock's resolution: a step - // that did run can measure as zero, so dur cannot distinguish "never ran" - // from "ran very fast". + // is what an early failure returns. ran is tracked explicitly because a step + // that did run can measure a zero dur within the clock's resolution. cryptoStep struct { ran bool dur time.Duration @@ -252,9 +248,8 @@ func (v *Vault) Seal(ctx context.Context, ns string, data []byte) (*Message, err // using the KEK identified by the material carried in msg, served from the // decrypted-DEK cache when it is enabled. // -// Exactly one [EnvelopeEvent] is reported to the Observer, on every path -// including failures. It carries no namespace: the KEK is selected by ID from -// the material, so Open never learns one. +// Exactly one [EnvelopeEvent], with no namespace, is reported to the Observer +// on every path including failures. func (v *Vault) Open(ctx context.Context, msg *Message) ([]byte, error) { start := time.Now() pt, step, err := v.open(ctx, msg) @@ -272,11 +267,9 @@ func (v *Vault) Open(ctx context.Context, msg *Message) ([]byte, error) { } // Refresh rotates every namespace DEK that has reached its renewal threshold. -// It is meant to be called periodically. Seal also rotates an expired DEK on -// demand, so Refresh is an optimization that keeps rotation off the request -// path rather than a correctness requirement. -// -// One [RotationEvent] with [RotationScheduled] is reported per key rotated. +// It is meant to be called periodically to keep rotation off the request path; +// Seal rotates an expired DEK on demand, so calling it is optional. One +// [RotationEvent] with [RotationScheduled] is reported per key rotated. func (v *Vault) Refresh() error { // Find expired keys without acquiring a write lock. v.mu.RLock() diff --git a/pkg/ext/auth.go b/pkg/ext/auth.go index 7880a55..40f4d71 100644 --- a/pkg/ext/auth.go +++ b/pkg/ext/auth.go @@ -23,15 +23,11 @@ var ( // spelled out, so it cannot drift from what was registered. healthPrefix = "/" + grpc_health_v1.Health_ServiceDesc.ServiceName + "/" - // healthCheckMethods is the set [IsHealthCheckMethod] reports on, and concerns - // the other end from healthPrefix above: these are methods callers of the proxy - // invoke, not methods of this server. - // - // GetSystemInfo is not a health check. It is here because it is the first call an - // SDK client makes on connect, and Temporal's own authorizer groups it with the - // health checks for that reason. Health_Watch is here because the proxy serves - // it: it registers gRPC's standard health server, so a caller that watches rather - // than polls would otherwise be refused. + // healthCheckMethods is the set [IsHealthCheckMethod] reports on: methods + // callers of the proxy invoke, not methods of this server. GetSystemInfo is + // not a health check but is an SDK client's first call on connect, and + // Health_Watch is included because the proxy serves gRPC's standard health + // server. healthCheckMethods = map[string]struct{}{ grpc_health_v1.Health_Check_FullMethodName: {}, grpc_health_v1.Health_Watch_FullMethodName: {}, @@ -50,23 +46,20 @@ type ( // Auth decides whether an inbound caller of the proxy may proceed. Register an // implementation with [WithAuth]. // - // The request carries what the proxy knows about the call. Its credentials are - // what the proxy lifted from the caller's stream, one entry per configured - // credential header the caller actually sent, so an empty slice means it - // presented none; its target is what the call is addressing. The proxy's own - // credential to this server is not among the credentials and stays in the - // request metadata, alongside the caller's other metadata. + // The request's credentials hold one entry per configured credential header + // the caller sent, so an empty slice means it presented none; its target is + // what the call is addressing. The proxy's own credential to this server is + // not among the credentials and stays in the request metadata. // - // Answer with a response whose Decision is set: only DECISION_ALLOW admits, so - // an unset decision denies rather than admits by accident. Reason is for - // whoever operates this server, and the proxy keeps it out of what the rejected - // caller is told. Return an error only when no verdict was reached, such as an + // Only a Decision of DECISION_ALLOW admits; an unset decision denies. Reason + // is for whoever operates this server and is withheld from the rejected + // caller. Return an error only when no verdict was reached, such as an // unreachable backend; the proxy denies either way, but an error keeps its // status code, so [google.golang.org/grpc/codes.Unavailable] tells a worker to // retry where a denial does not. // // Implementations must be safe for concurrent use and must not block - // indefinitely, since a caller is waiting and the proxy denies on timeout. + // indefinitely; the proxy denies on timeout. Auth interface { Authenticate(context.Context, *auth.AuthRequest) (*auth.AuthResponse, error) } @@ -84,12 +77,9 @@ func Allow() *auth.AuthResponse { } // Deny returns the response that refuses a caller. The reason is recorded by the -// proxy and withheld from the caller, so write it for whoever operates this -// server: it may name subjects and internal systems. -// -// Use this for a caller judged and found wanting, and an error for a verdict never -// reached, such as an unreachable backend. Both deny, but an error keeps its -// status code, which is what tells a worker whether retrying could help. +// proxy and withheld from the caller, so it may name subjects and internal +// systems. Return an error instead when no verdict was reached; both deny, but +// an error keeps its status code, which tells a worker whether to retry. func Deny(reason string) *auth.AuthResponse { return &auth.AuthResponse{Decision: auth.AuthResponse_DECISION_DENY, Reason: reason} } @@ -99,11 +89,8 @@ func Deny(reason string) *auth.AuthResponse { // // Every failure is an [google.golang.org/grpc/codes.Unauthenticated] status error, // ready to return from [Auth.Authenticate]: no credential on that header, more -// than one value, or a value carrying some other scheme. A repeated value is -// refused rather than resolved by taking the first, since choosing among -// credentials a caller sent is how a check gets bypassed. -// -// An implementation that would rather answer [Deny], or that accepts a credential +// than one value, or a value carrying some other scheme or no token. An +// implementation that would rather answer [Deny], or that accepts a credential // with no scheme at all, should read req.GetCredentials() directly. func BearerToken(req *auth.AuthRequest, header string) (string, error) { hdr := strings.ToLower(header) @@ -139,13 +126,10 @@ func BearerToken(req *auth.AuthRequest, header string) (string, error) { // IsHealthCheckMethod reports whether full, a gRPC full method name as it arrives // in [api.auth.v1.Target], is one an implementation will usually admit without a -// credential: the gRPC health methods, and GetSystemInfo, which is the first call -// an SDK client makes on connect and so decides whether it can connect at all. -// -// Whether to admit them is policy and stays with the implementation, which is why -// this reports rather than decides. Refusing them is a defensible choice; it makes -// the proxy look unhealthy to anything probing it, and makes an unauthenticated -// client fail at dial instead of on its first real call. +// credential: the gRPC health methods, and GetSystemInfo, the first call an SDK +// client makes on connect. Whether to admit them is the implementation's call; +// refusing them makes the proxy look unhealthy to probes and makes an +// unauthenticated client fail at dial instead of on its first real call. func IsHealthCheckMethod(full string) bool { _, ok := healthCheckMethods[full] @@ -178,13 +162,8 @@ func (a *authService) Auth(ctx context.Context, req *auth.AuthRequest) (*auth.Au // unaryGuard returns the interceptor [WithServerAuth] installs, which documents // the contract. [Serve] installs it only for a non-nil check. Rejections are // Unauthenticated, and the message separates an absent credential from a rejected -// one: both ends here are operator-run, so telling "wrong header" from "wrong -// value" is worth more than withholding it. -// -// The health service is exempt. Guarding it would break every probe, which has no -// credential to present and in Kubernetes' native gRPC prober cannot send -// metadata at all, and would withhold nothing in exchange: Watch reports the same -// status over a stream, which this interceptor does not see. +// one. The health service is exempt, since probes have no credential to present +// and Watch, a stream, is not seen by this interceptor anyway. func unaryGuard(hdr string, check CredentialCheck) grpc.UnaryServerInterceptor { return func( ctx context.Context, diff --git a/pkg/ext/keywrap.go b/pkg/ext/keywrap.go index 518893b..783d6fd 100644 --- a/pkg/ext/keywrap.go +++ b/pkg/ext/keywrap.go @@ -38,16 +38,12 @@ type ( Version string } - // KeyLookup supplies the wrapping keys [NewKeyWrapper] seals with. It is the - // only thing NewKeyWrapper cannot supply for itself. - // - // A lookup must answer for every version it ever reported, not only the - // current one: forgetting a version destroys every payload sealed under it. - // Returning a [google.golang.org/grpc/status] error passes its code through to - // the proxy, which is how an unreachable key store is distinguished from a - // version that will never resolve. - // - // A lookup must be safe for concurrent use. + // KeyLookup supplies the wrapping keys [NewKeyWrapper] seals with. A lookup + // must be safe for concurrent use and must answer for every version it ever + // reported, not only the current one: forgetting a version destroys every + // payload sealed under it. A [google.golang.org/grpc/status] error passes its + // code through to the proxy, which is how an unreachable key store is + // distinguished from a version that will never resolve. KeyLookup func(context.Context, KeyRequest) (Key, error) // KeyWrapperOption configures a key wrapper during construction. @@ -67,17 +63,11 @@ type ( ) // NewKeyWrapper returns a [KMS] that seals DEKs with an AEAD over keys from -// lookup and frames them as [ext.KeyMaterial], so an extension server supplies -// key material and nothing else. -// -// New material is sealed with AES-256-GCM unless [WithCipher] says otherwise, -// and carries both the cipher that sealed it and the key version lookup reported -// at the time. Opening reads those from the material rather than from the -// configuration, so changing cipher, or a key store rotating underneath, leaves -// everything already sealed readable. -// -// Ciphers are registered during construction only, so the returned KMS never -// changes afterwards and may be shared by any number of goroutines. +// lookup and frames them as [ext.KeyMaterial]. New material is sealed with +// AES-256-GCM unless [WithCipher] says otherwise, and records the cipher and key +// version used; opening reads both from the material, so changing cipher or +// rotating keys leaves everything already sealed readable. The returned KMS is +// safe for concurrent use. func NewKeyWrapper(lookup KeyLookup, opts ...KeyWrapperOption) (KMS, error) { if lookup == nil { return nil, errors.New("a key lookup is required") @@ -114,14 +104,11 @@ func WithCipher(id CipherID) KeyWrapperOption { } // WithCipherFunc registers fn as the constructor for id, replacing whatever was -// registered before, including a built-in. -// -// Ids from 128 up are reserved for exactly this and will never be assigned by -// [ext.KeyMaterial_Cipher], so a cipher registered there cannot collide with -// one added later; [MustCipherID] builds one. An id below that is accepted, -// since replacing a built-in with a stricter construction of the same cipher is -// reasonable, but reusing a built-in id for a different cipher makes material -// that other servers will misread. An id of zero or less is rejected. +// registered before, including a built-in. Ids from 128 up are reserved for +// this and will never be assigned by [ext.KeyMaterial_Cipher]; [MustCipherID] +// builds one. A lower id is accepted, but reusing a built-in id for a different +// cipher makes material that other servers will misread. An id of zero or less +// is rejected. func WithCipherFunc(id CipherID, fn CipherFunc) KeyWrapperOption { return keyWrapperOpt(func(w *keyWrapper) error { if fn == nil { diff --git a/pkg/ext/kms.go b/pkg/ext/kms.go index 7428be2..81e38fa 100644 --- a/pkg/ext/kms.go +++ b/pkg/ext/kms.go @@ -14,22 +14,19 @@ type ( // KMS wraps and unwraps the proxy's data encryption keys. Register an // implementation with [WithKMS]. // - // Only key material crosses the wire. The plaintext handed to Wrap is a DEK, - // never a payload, so an implementation is free to make each call a round trip - // to an HSM; the proxy caches the DEK and does the bulk encryption itself. + // The plaintext handed to Wrap is a DEK, never a payload; the proxy caches + // the DEK and does the bulk encryption itself, so each call may be a round + // trip to an HSM. // // Wrap receives the namespace, so an implementation may hold a distinct key per - // namespace. Unwrap does not: it gets only the ciphertext, so whatever - // identifies the key has to be inside what Wrap returned, usually an opaque - // header framed around it. That makes Unwrap's input a durable format worth - // versioning, and retiring a key destroys every payload it wrapped. + // namespace. Unwrap gets only the ciphertext, so whatever identifies the key + // has to be inside what Wrap returned, usually an opaque header framed around + // it. That makes Unwrap's input a durable format worth versioning, and retiring + // a key destroys every payload it wrapped. // // A [google.golang.org/grpc/status] error is passed through with its code - // intact, so an implementation that can tell a bad ciphertext from an - // unreachable backend can say which it was; any other error takes the code - // documented on the method that called it. - // - // Implementations must be safe for concurrent use. + // intact; any other error takes the code documented on the method that called + // it. Implementations must be safe for concurrent use. KMS interface { Wrap(context.Context, string, []byte) ([]byte, error) Unwrap(context.Context, []byte) ([]byte, error) @@ -44,14 +41,10 @@ type ( } ) -// Encrypt wraps a DEK for the namespace named in the request. -// -// An empty namespace is refused rather than defaulted, because it selects the key: -// an implementation keyed by namespace would otherwise wrap under whatever its -// zero value picks, and the ciphertext would be unrecoverable once the mistake was -// found. A failure from the implementation is Internal, since the fault is this -// server's rather than the request's, unless the implementation chose a code of -// its own. +// Encrypt wraps a DEK for the namespace named in the request. An empty namespace +// is refused with InvalidArgument rather than defaulted, since it selects the +// key. A failure from the implementation is Internal unless the implementation +// chose a code of its own. func (s *kmsService) Encrypt(ctx context.Context, req *kms.EncryptRequest) (*kms.EncryptResponse, error) { if s.kms == nil { return s.UnimplementedEncryptionServiceServer.Encrypt(ctx, req) @@ -70,13 +63,10 @@ func (s *kmsService) Encrypt(ctx context.Context, req *kms.EncryptRequest) (*kms return &kms.EncryptResponse{Ciphertext: ct}, nil } -// Decrypt unwraps a DEK previously produced by Encrypt. -// -// A bare failure is InvalidArgument rather than Internal because the likely -// cause is the ciphertext: wrapped by another server, under a retired key, or in -// a format this build no longer reads, none of which improve on a retry. An -// unreachable backend does improve on one, and an implementation that can tell -// the two apart says so by returning Unavailable itself. +// Decrypt unwraps a DEK previously produced by Encrypt. A failure from the +// implementation is InvalidArgument, marking the ciphertext as the likely cause, +// unless the implementation chose a code of its own, such as Unavailable for an +// unreachable backend. func (s *kmsService) Decrypt(ctx context.Context, req *kms.DecryptRequest) (*kms.DecryptResponse, error) { if s.kms == nil { return s.UnimplementedEncryptionServiceServer.Decrypt(ctx, req) @@ -92,15 +82,10 @@ func (s *kmsService) Decrypt(ctx context.Context, req *kms.DecryptRequest) (*kms } // implError reports err from a [KMS] implementation, keeping its code when it -// picked one so the proxy can tell a retry from a dead end, and falling back to -// code otherwise. Both handlers route through here because the two must not -// drift: a code preserved on one path and discarded on the other is a contract -// an implementation cannot write against. -// -// A status carrying codes.OK falls back as well, which [status.Error] cannot -// produce but a hand-rolled GRPCStatus can. gRPC writes it as a call that -// succeeded without a response, so the proxy sees a cardinality violation rather -// than whatever the implementation was reporting. +// picked one and falling back to code otherwise. Both handlers route through +// here so the two cannot drift. A status carrying codes.OK, which a hand-rolled +// GRPCStatus can produce, falls back as well: gRPC would write it as a success +// without a response. func implError(err error, code codes.Code, msg string) error { if s, ok := status.FromError(err); ok && s.Code() != codes.OK { return err diff --git a/pkg/ext/plaintext.go b/pkg/ext/plaintext.go index 3959f36..19e1dd6 100644 --- a/pkg/ext/plaintext.go +++ b/pkg/ext/plaintext.go @@ -17,13 +17,9 @@ import ( const plaintextMessage = "Serving in plaintext. Supply credentials via WithServerOption for production use." // plaintextWarning returns interceptors that log [plaintextMessage] once if calls -// are arriving over an unencrypted connection. -// -// It reads the connection rather than the configuration because a -// [grpc.ServerOption] is opaque: what [WithServerOption] was handed cannot be read -// back, so whether this server ended up serving TLS is only knowable from a call -// that actually arrived. The cost is that a server nobody ever calls stays quiet, -// which the health service makes unlikely. +// are arriving over an unencrypted connection. It inspects arriving calls +// because what [WithServerOption] was handed cannot be read back, so a server +// nobody calls stays quiet. func plaintextWarning(log logger.Logger) (grpc.UnaryServerInterceptor, grpc.StreamServerInterceptor) { var once sync.Once diff --git a/pkg/ext/sealwrap.go b/pkg/ext/sealwrap.go index b825e09..046bbb8 100644 --- a/pkg/ext/sealwrap.go +++ b/pkg/ext/sealwrap.go @@ -50,13 +50,11 @@ type ( // KeySealer seals and opens DEKs through a key service that will not hand over // its keys. Hand one to [NewSealWrapper]. // - // The wrapper authenticates nothing, and cannot: it holds no key, and the - // version and opaque bytes it would bind are chosen by Seal, so they do not - // exist until the call it would bind them to has already happened. An - // implementation that ignores [BindingContext] produces material whose - // namespace, version, and opaque can be swapped by anyone able to write a - // payload's metadata. Passing those bytes to the key service as an encryption - // context is what makes relabelled material fail to open instead. + // The wrapper authenticates nothing. An implementation that ignores + // [BindingContext] produces material whose namespace, version, and opaque can + // be swapped by anyone able to write a payload's metadata; pass those bytes to + // the key service as an encryption context so relabelled material fails to + // open. // // A [google.golang.org/grpc/status] error is passed through with its code // intact. Implementations must be safe for concurrent use. diff --git a/pkg/ext/server.go b/pkg/ext/server.go index 2ac6280..bc57dee 100644 --- a/pkg/ext/server.go +++ b/pkg/ext/server.go @@ -169,21 +169,14 @@ func WithLogger(l logger.Logger) Option { // WithServerAuth guards this server's unary methods with check, which receives // the single value of the header metadata. It authenticates the proxy to this // server, unlike the [Auth] service, which is how the proxy asks about somebody -// else. Unary covers both generated services, since every method on them is unary; -// adding a streaming method to either proto means pairing this with a -// [grpc.StreamServerInterceptor]. -// -// The health service is exempt, and deliberately: a probe has no credential to -// give, and guarding it would withhold nothing anyway, because Watch reports the -// same status over a stream that no unary interceptor sees. +// else. Every method on both generated services is unary; the health service is +// exempt. // // A call is rejected with Unauthenticated unless the header is present exactly -// once and check accepts its value. Repeats are refused rather than searched, so -// a caller cannot spray guesses in one call. Compare in constant time -// ([crypto/subtle.ConstantTimeCompare]) for a shared secret. A nil check installs -// nothing and leaves the server open, while a non-nil one with an empty header -// would reject everything, so [Serve] treats it as a configuration error and -// declines to start. +// once and check accepts its value. Compare in constant time +// ([crypto/subtle.ConstantTimeCompare]) for a shared secret. A nil check +// installs nothing and leaves the server open; a non-nil one with an empty +// header makes [Serve] return an error without starting. func WithServerAuth(header string, check CredentialCheck) Option { return func(o *options) { o.authHeader, o.authCheck = header, check } } @@ -193,29 +186,25 @@ func WithServerAuth(header string, check CredentialCheck) Option { // limit. // // Options accumulate, and [Serve] starts with insecure credentials, 128 concurrent -// streams, a 1MiB receive limit sized for key material rather than payloads, and -// keepalive settings that floor client ping intervals. An option given here is -// applied after those, so it wins for any setting gRPC resolves to one value: -// serving TLS is grpc.Creds(credentials.NewTLS(...)), with nothing to clear first. +// streams, a 1MiB receive limit sized for key material, and keepalive settings +// that floor client ping intervals. An option given here is applied after those, +// so it wins for any setting gRPC resolves to one value: serving TLS is +// grpc.Creds(credentials.NewTLS(...)), with nothing to clear first. // -// Interceptor order is visible. The guard [WithServerAuth] installs is chained -// ahead of anything added here, so a [grpc.ChainUnaryInterceptor] sees only -// admitted calls. [grpc.UnaryInterceptor] is gRPC's own exception: prepended ahead -// of the whole chain, it observes calls about to be rejected. +// The guard [WithServerAuth] installs is chained ahead of anything added here, so +// a [grpc.ChainUnaryInterceptor] sees only admitted calls. A +// [grpc.UnaryInterceptor] runs ahead of the whole chain and also sees calls about +// to be rejected. func WithServerOption(opts ...grpc.ServerOption) Option { return func(o *options) { o.serverOptions = append(o.serverOptions, opts...) } } // WithShutdownTimeout bounds how long shutdown waits for in-flight calls before // dropping connections. It defaults to five seconds and is clamped to a 50ms -// floor, so a zero or negative value still lets an already-answered call flush. -// Set it below the grace period of whatever supervises the process; overrunning -// that trades the graceful shutdown for a SIGKILL. -// -// A client watching the health service holds the drain open until this expires, -// because a Watch ends with its stream rather than with the status going -// NOT_SERVING. Expect the bound to be reached, not merely available, wherever -// something watches. +// floor. Set it below the grace period of whatever supervises the process, or +// shutdown ends in a SIGKILL. A client watching the health service holds the +// drain open until this expires, so expect the full bound wherever something +// watches. func WithShutdownTimeout(t time.Duration) Option { return func(o *options) { o.shutdownTimeout = max(50*time.Millisecond, t) } } diff --git a/pkg/validation/certs/certs.go b/pkg/validation/certs/certs.go index c32d008..f2126f1 100644 --- a/pkg/validation/certs/certs.go +++ b/pkg/validation/certs/certs.go @@ -192,22 +192,11 @@ func IsCA() Check { // it additionally rejects certificates whose public key type is incompatible // with every suite in that list. // -// Self-issued certificates (RawIssuer == RawSubject) are exempt from the -// signature-algorithm check. This is a superset of self-signed: it also -// admits CA key-rollover certs, which are self-issued but signed by a -// different (older or newer) key than the one they certify. The exemption -// holds regardless: a self-issued cert's own signature is never consulted -// during chain verification; the cert is trusted (or not) based on its -// presence in the trust store, or on its role elsewhere in the chain, not on -// its self-attestation. Many still-valid public roots, used to sign SHA-256 -// chains today, carry legacy SHA-1 self-signatures; rejecting them would -// make the system CA bundle unusable as a trust anchor. -// -// SecureAlgorithm is used both for trust-anchor validation and for -// certificates presented in a peer's chain. The exemption applies in both -// cases: a self-issued cert anywhere in a presented chain skips the -// weak-signature check, for the same reason. The key-type check still runs -// unconditionally. +// Self-issued certificates (RawIssuer == RawSubject), a superset of self-signed +// that includes CA key-rollover certs, are exempt from the signature-algorithm +// check wherever they appear, in a trust store or in a presented chain, since a +// self-issued cert's own signature is never consulted during chain +// verification. The key-type check runs unconditionally. func SecureAlgorithm(allowedSuites ...uint16) Check { return func(cert *x509.Certificate) error { selfIssued := bytes.Equal(cert.RawIssuer, cert.RawSubject)