From b2553bec4725f9345794ba08184a06fa9ede8baa Mon Sep 17 00:00:00 2001 From: "David Muto (pseudomuto)" Date: Wed, 7 Oct 2026 11:47:15 -0400 Subject: [PATCH] [docs]: Trim godoc to caller-visible contract Many doc comments had grown into multi-paragraph essays covering design rationale, history, and why alternatives were rejected. That context is useful once, during review, but every later reader of the godoc has to scroll past it to find the part that matters to them. This cuts identifier doc comments down to what the thing is plus the constraints a caller can observe: errors and status codes, nil handling, ordering, lifecycle and close obligations, concurrency safety, and operator or security warnings. Rationale that explains a non-obvious guard in unexported code is kept in condensed form so the guard does not look removable. Package docs and the examples are left alone, since they are where an overview or a walkthrough belongs. While here, the codecserver docs still described per-upstream proxies, which no longer exist; they now refer to each upstream's forwarder. Only comments change. --- cmd/proxy/serve.go | 2 - internal/api/auth.go | 45 +++-------- internal/api/fx.go | 29 +++---- internal/auth/authenticator.go | 5 +- internal/auth/jwks.go | 28 ++----- internal/cloud/endpoint.go | 12 ++- internal/cloud/namespace.go | 2 +- internal/cloud/translation/interceptor.go | 18 +---- internal/cloud/translation/namespaces.go | 34 +++----- internal/cloud/translation/translation.go | 15 ++-- internal/codecserver/doc.go | 4 +- internal/codecserver/fx.go | 26 ++---- internal/codecserver/handler.go | 55 ++++--------- internal/codecserver/middleware.go | 17 ++-- internal/codecserver/namespaces.go | 41 +++------- internal/codecserver/reporter.go | 23 ++---- internal/codecserver/server.go | 26 +++--- internal/config/auth.go | 28 +++---- internal/config/cloudapi.go | 96 +++++++---------------- internal/config/encryption.go | 20 +---- internal/config/extensions.go | 13 ++- internal/config/routing.go | 6 +- internal/dataplane/dataplane.go | 9 +-- internal/kms/fx.go | 9 +-- internal/kms/reporter.go | 27 +++---- internal/proxy/version.go | 7 +- internal/router/reporter.go | 20 ++--- internal/server/loopback.go | 71 +++++------------ internal/server/reporter.go | 20 ++--- internal/server/server.go | 9 +-- internal/transport/connect/conn.go | 26 +++--- internal/transport/connect/fx.go | 10 +-- pkg/crypto/keys.go | 62 +++++---------- pkg/crypto/vault.go | 21 ++--- pkg/ext/auth.go | 69 ++++++---------- pkg/ext/keywrap.go | 45 ++++------- pkg/ext/kms.go | 57 +++++--------- pkg/ext/plaintext.go | 10 +-- pkg/ext/sealwrap.go | 12 ++- pkg/ext/server.go | 47 +++++------ pkg/validation/certs/certs.go | 21 ++--- 41 files changed, 350 insertions(+), 747 deletions(-) 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)