From 0244cfbf2ab7613d000cfe9eccddaf371e0d7ede Mon Sep 17 00:00:00 2001 From: "David Muto (pseudomuto)" Date: Wed, 7 Oct 2026 10:47:13 -0400 Subject: [PATCH] [docs]: Fix inaccurate and missing godoc comments A number of doc comments had drifted from the code they describe. Some named identifiers that no longer exist (rpc.Method, a CLI logger, perUpstream), many still described the per-upstream proxy and unix socket removed in #205, and a few stated the opposite of what the code does: GT/GTE/LT claimed NaN never passes when it always does, NamespacedVault.Open claimed to use its namespace, etc. This corrects those against the current code, fills in the doc comments flagged as missing, and documents caller-visible constraints that were only discoverable by reading the implementation, such as NewKEKRegistry rejecting duplicate key IDs and KeyConfig's bounds. It also drops a spec reference and the remaining em-dashes. Only comments change here. --- cmd/proxy/main.go | 4 ++++ e2e/doc.go | 2 +- examples/authz/server/authorizer.go | 13 +++++++---- examples/authz/server/claims.go | 5 ++-- examples/authz/token.go | 4 ++-- internal/api/fx.go | 6 ++--- internal/api/kms.go | 2 +- internal/auth/jwks.go | 14 +++++------ internal/auth/outbound/provider.go | 5 ++-- internal/auth/outbound/static_provider.go | 8 +++---- internal/cloud/doc.go | 9 ++++---- internal/cloud/namespace.go | 3 ++- internal/cloud/translation/translation.go | 3 ++- internal/codecserver/doc.go | 2 +- internal/codecserver/fx.go | 8 +++---- internal/codecserver/reporter.go | 8 +++---- internal/config/auth.go | 2 +- internal/config/cloudapi.go | 5 ++-- internal/config/codecserver.go | 3 ++- internal/config/config.go | 2 +- internal/config/encryption.go | 2 +- internal/config/extensions.go | 11 +++++---- internal/config/fx.go | 2 ++ internal/config/routing.go | 6 ++--- internal/config/upstream.go | 14 ++++++++--- internal/dataplane/dataplane.go | 4 ++-- .../dataplane/dataplanetest/dataplanetest.go | 7 +++--- internal/dataplane/dataplanetest/upstream.go | 6 ++--- internal/dataplane/doc.go | 2 +- internal/dataplane/fx.go | 6 ++--- internal/dataplane/lifecycle.go | 9 ++++---- internal/kms/doc.go | 10 ++++---- internal/kms/fx.go | 2 +- internal/metrics/fx.go | 6 ++--- internal/metrics/labels.go | 4 ++-- internal/protoutil/translate.go | 6 ++--- internal/proxy/codec.go | 3 ++- internal/proxy/forward.go | 4 ++-- internal/proxy/reporter.go | 23 +++++++++++-------- internal/proxy/resolver.go | 19 +++++++-------- internal/proxy/translation.go | 7 +++--- internal/router/codec.go | 7 +++--- internal/router/mux.go | 2 +- internal/router/reporter.go | 12 +++++----- internal/rpc/doc.go | 12 +++++----- internal/rpc/pump.go | 17 +++++++------- internal/server/loopback.go | 8 +++---- internal/server/reporter.go | 5 ++-- internal/server/server.go | 11 +++++---- internal/services/services.go | 8 +++---- internal/template/template.go | 3 +++ internal/transport/connect/doc.go | 8 +++---- internal/transport/connect/pool.go | 15 ++++++------ internal/transport/creds/doc.go | 5 ++-- internal/transport/creds/options.go | 2 +- internal/transport/meta/meta.go | 7 +++--- pkg/api/ext/v1/key_material.go | 9 +++++--- pkg/codec/chain.go | 4 ++-- pkg/codec/encryptor.go | 4 ++-- pkg/crypto/dek.go | 6 +++-- pkg/crypto/kek.go | 14 +++++++---- pkg/crypto/keys.go | 4 +++- pkg/crypto/vault.go | 9 +++++--- pkg/ext/binding.go | 4 +++- pkg/ext/keywrap.go | 3 ++- pkg/ext/server.go | 11 +++++---- pkg/logger/doc.go | 7 +++--- pkg/logger/logger.go | 3 ++- pkg/logger/noop.go | 10 ++++++-- pkg/logger/tag/tag.go | 6 +++-- pkg/logger/test.go | 14 +++++++++-- pkg/logger/zero.go | 8 +++++++ pkg/testutil/certs.go | 18 +++++++-------- pkg/validation/certs/certs.go | 19 ++++++++------- pkg/validation/certs/doc.go | 7 +++--- pkg/validation/checks.go | 13 +++++++---- pkg/validation/doc.go | 5 ++-- pkg/validation/validate.go | 6 ++--- 78 files changed, 329 insertions(+), 238 deletions(-) diff --git a/cmd/proxy/main.go b/cmd/proxy/main.go index b02e668..9250a93 100644 --- a/cmd/proxy/main.go +++ b/cmd/proxy/main.go @@ -1,3 +1,7 @@ +// Command proxy runs the Temporal proxy: a gRPC gateway between Temporal SDK +// clients, workers, and the Temporal UI and one or more upstream Temporal +// Services, handling namespace translation, TLS, and payload encryption. The +// serve subcommand starts it from a YAML config file. package main import ( diff --git a/e2e/doc.go b/e2e/doc.go index 2121ab4..f0b43ff 100644 --- a/e2e/doc.go +++ b/e2e/doc.go @@ -1,4 +1,4 @@ // Package e2e contains black-box, full-stack integration tests that drive the // proxy the way production does: a client through the gateway, router, -// per-upstream proxy, and out to a fake upstream. +// per-upstream forwarder, and out to a fake upstream. package e2e diff --git a/examples/authz/server/authorizer.go b/examples/authz/server/authorizer.go index fc2896f..f7a2bc7 100644 --- a/examples/authz/server/authorizer.go +++ b/examples/authz/server/authorizer.go @@ -26,15 +26,17 @@ const ( ) const ( - // The scope a method acts in, which decides whose roles are consulted: a - // namespace-scoped method reads the caller's roles in the namespace it names, - // while a cluster-scoped one has no namespace and reads only system roles. + // scopeCluster and scopeNamespace are the scope a method acts in, which + // decides whose roles are consulted: a namespace-scoped method reads the + // caller's roles in the namespace it names, while a cluster-scoped one has no + // namespace and reads only system roles. scopeCluster scope = iota + 1 scopeNamespace ) const ( - // How much authority a method calls for, which maps to a role in access.role. + // accessReadOnly, accessWrite, and accessAdmin are how much authority a + // method calls for, which maps to a role in access.role. accessReadOnly access = iota + 1 accessWrite accessAdmin @@ -87,7 +89,8 @@ var ( workflowService + "RegisterNamespace": {scopeNamespace, accessAdmin}, // Cluster scope. Nothing here names a namespace, so only system roles count. - // GetSystemInfo would belong here too, but alwaysAllowed answers it first. + // GetSystemInfo would belong here too, but Authenticate admits it first, as + // a method ext.IsHealthCheckMethod reports. workflowService + "GetClusterInfo": {scopeCluster, accessReadOnly}, workflowService + "ListNamespaces": {scopeCluster, accessReadOnly}, } diff --git a/examples/authz/server/claims.go b/examples/authz/server/claims.go index 9bde5ec..1d5cad4 100644 --- a/examples/authz/server/claims.go +++ b/examples/authz/server/claims.go @@ -13,8 +13,9 @@ import ( ) const ( - // The roles a subject can hold within a scope, as a bitmask so they combine: - // an effective role is the OR of every role granted at that scope. These are + // roleWorker through roleAdmin are the roles a subject can hold within a + // scope, as a bitmask so they combine: an effective role is the OR of every + // role granted at that scope. These are // go.temporal.io/server/common/authorization's Role values, with the same // numbering, because Authenticate compares them the way that package does. roleWorker = role(1 << iota) diff --git a/examples/authz/token.go b/examples/authz/token.go index b49d8fa..d7a0b65 100644 --- a/examples/authz/token.go +++ b/examples/authz/token.go @@ -19,8 +19,8 @@ const ( // between them is visible, rather than in the commands that would drift apart. PermissionsClaim = "https://acme.example/temporal" - // SystemScope names the cluster scope in the permissions claim, the authority a - // subject holds regardless of namespace. + // SystemScope is the scope name logged for a role in Permissions.System, the + // authority a subject holds regardless of namespace. SystemScope = "system" ) diff --git a/internal/api/fx.go b/internal/api/fx.go index 2b677ae..f134089 100644 --- a/internal/api/fx.go +++ b/internal/api/fx.go @@ -25,7 +25,7 @@ const extensionKeyPrefix = "extension:" // Module provides the pooled connection for every configured extension server, // built when the provider runs rather than on first use, so a bad dial target -// surfaces at construction instead of on the first encryption call, and opened on +// surfaces at construction instead of on the first call, and opened on // start so an unreachable server (or one whose certificate this proxy will not // accept) fails startup. Per-call credentials are still only exercised by a real // request. @@ -57,8 +57,8 @@ var Module = fx.Options( // Config rejects a templated extension server hostPort, so every one of // these is static and reachable now or not at all. An extension server backs - // payload encryption, so serving without one means failing encrypted - // traffic; fail startup instead. + // payload encryption or inbound auth, so serving without one means failing + // that traffic; fail startup instead. p.Lifecycle.Append(fx.StartHook(func(ctx context.Context) error { if err := connect.WaitReady(ctx, conns...); err != nil { return fmt.Errorf("extension server connection not ready: %w", err) diff --git a/internal/api/kms.go b/internal/api/kms.go index 22bab6c..7b548ff 100644 --- a/internal/api/kms.go +++ b/internal/api/kms.go @@ -37,7 +37,7 @@ func (k *KMS) Close() error { return nil } -// ID returns a unique ID for this KEK, e.g. a KMS ARN. +// ID returns a unique ID for this KEK: the extension key URI it was opened from. func (k *KMS) ID() string { return k.id } diff --git a/internal/auth/jwks.go b/internal/auth/jwks.go index abb2470..476cb3e 100644 --- a/internal/auth/jwks.go +++ b/internal/auth/jwks.go @@ -225,13 +225,13 @@ func deferredKeyfunc(load func() (jwt.Keyfunc, error)) (jwt.Keyfunc, *atomic.Poi return keyfn, &ready } -// wrapKeyfunc adapts a key resolver so the error taxonomy matches the spec's -// intent: a genuinely unknown key id on a POPULATED keyset is a verification -// failure (jwkset.ErrKeyNotFound -> codes.Unauthenticated), while an unknown -// key id on an EMPTY keyset means the keyset was never fetched (IdP unreachable -// at startup and on-demand refresh still failing), which is an availability -// problem (errKeysUnavailable -> codes.Unavailable), not a bad token. Any other -// resolver error is treated as an availability problem. +// wrapKeyfunc adapts a key resolver so the error taxonomy separates bad tokens +// from outages: a genuinely unknown key id on a populated keyset is a +// verification failure (jwkset.ErrKeyNotFound -> codes.Unauthenticated), while +// an unknown key id on an empty keyset means the keyset was never fetched (IdP +// unreachable at startup and on-demand refresh still failing), which is an +// 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). diff --git a/internal/auth/outbound/provider.go b/internal/auth/outbound/provider.go index 46db337..592fdc4 100644 --- a/internal/auth/outbound/provider.go +++ b/internal/auth/outbound/provider.go @@ -13,8 +13,9 @@ import ( ) // CredentialProvider supplies per-RPC metadata for outbound calls to an -// upstream. Header reports the metadata header it sets, so the proxy can strip -// any forwarded value on that header before the credential adds its own. +// upstream, an extension server, or the Cloud API. Header reports the metadata +// header it sets, so the proxy can strip any forwarded value on that header +// before the credential adds its own. type CredentialProvider interface { credentials.PerRPCCredentials Header() string diff --git a/internal/auth/outbound/static_provider.go b/internal/auth/outbound/static_provider.go index cd82642..9008c19 100644 --- a/internal/auth/outbound/static_provider.go +++ b/internal/auth/outbound/static_provider.go @@ -11,10 +11,10 @@ const ( defaultScheme = "Bearer" ) -// StaticCredentialProvider attaches a fixed bearer header to every outbound -// request to an upstream. It implements google.golang.org/grpc/credentials -// PerRPCCredentials and requires transport security, so gRPC refuses to send -// the credential over an insecure connection. +// StaticCredentialProvider attaches a fixed " " header to every +// outbound request on the connection it is dialed with. It implements +// google.golang.org/grpc/credentials PerRPCCredentials and requires transport +// security, so gRPC refuses to send the credential over an insecure connection. type StaticCredentialProvider struct { header string value string diff --git a/internal/cloud/doc.go b/internal/cloud/doc.go index 452fa62..bc0c2b2 100644 --- a/internal/cloud/doc.go +++ b/internal/cloud/doc.go @@ -1,8 +1,9 @@ // Package cloud holds the rules that are specific to Temporal Cloud rather than // to any Temporal Service. // -// Today that means namespace validation. Cloud identifies a namespace as -// ".", a shape self-hosted deployments do not impose, so -// [ValidateNamespace] checks a string against it before the proxy uses that -// string to address a Cloud upstream. +// That covers recognizing Cloud endpoints ([IsEndpoint], [APIHostPort]) and +// namespace validation. Cloud identifies a namespace as ".", +// a shape self-hosted deployments do not impose, so [ValidateNamespace] checks +// a string against it before the proxy uses that string to address a Cloud +// upstream. package cloud diff --git a/internal/cloud/namespace.go b/internal/cloud/namespace.go index 9aed9eb..1afca89 100644 --- a/internal/cloud/namespace.go +++ b/internal/cloud/namespace.go @@ -17,7 +17,8 @@ const ( accountIDMaxLen = 20 ) -// start with letter, end with letter/number, contain only a-z0-9- +// nsNameRegex matches a namespace name label: it starts with a letter, ends +// with a letter or digit, and contains only a-z, 0-9, and hyphens. var nsNameRegex = regexp.MustCompile(`^[a-z][a-z0-9-]*[a-z0-9]$`) // ValidateAccountID checks that id is shaped like the account-id label of a diff --git a/internal/cloud/translation/translation.go b/internal/cloud/translation/translation.go index 39d99c7..d37f56d 100644 --- a/internal/cloud/translation/translation.go +++ b/internal/cloud/translation/translation.go @@ -18,7 +18,8 @@ type ( // caller is waiting on. Build one with [Adapt] rather than by hand, so the // conversions are written against concrete message types, or with [Answer] // for a method the proxy replies to itself. A Translation holds no per-call - // state and is safe for concurrent use. + // state and is safe for concurrent use once passed to [NewRegistry], which, + // like [Translation.WithHeader], modifies it. Translation struct { from, to string call func(ctx context.Context, req, reply proto.Message, send sendFunc) error diff --git a/internal/codecserver/doc.go b/internal/codecserver/doc.go index 597d134..12517b5 100644 --- a/internal/codecserver/doc.go +++ b/internal/codecserver/doc.go @@ -59,7 +59,7 @@ // This surface holds KMS unwrap permission. A reachable /decode without // authentication is a decryption oracle, and /encode is a sealing oracle, // which is why configuration refuses an enabled codec server bound beyond -// loopback with no auth block, and refuses auth without TLS. +// loopback with no auth block, and refuses auth without TLS on such a bind. // // The built-in authenticators validate the caller's credential and ignore the // target namespace, so a token that passes authorizes every namespace's diff --git a/internal/codecserver/fx.go b/internal/codecserver/fx.go index 06ceccf..74498b7 100644 --- a/internal/codecserver/fx.go +++ b/internal/codecserver/fx.go @@ -46,10 +46,10 @@ type Params struct { } // newFromParams builds the Server the configuration describes and binds it to -// 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. +// 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 diff --git a/internal/codecserver/reporter.go b/internal/codecserver/reporter.go index 5402971..bddea66 100644 --- a/internal/codecserver/reporter.go +++ b/internal/codecserver/reporter.go @@ -21,11 +21,9 @@ type Reporter struct { // 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. -// -// Parameters: -// - f: must already be scoped to the "codec_server" subsystem by the caller, -// which is what produces the published tmprl_proxy_codec_server_* names. +// 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. diff --git a/internal/config/auth.go b/internal/config/auth.go index 5ab0b30..f5c478a 100644 --- a/internal/config/auth.go +++ b/internal/config/auth.go @@ -133,7 +133,7 @@ func (c *StaticTokenConfig) Validate() error { ) } -// Validate requires a syntactically valid absolute JWKS URL. +// Validate requires a syntactically valid absolute https JWKS URL. func (c *JWKSConfig) Validate() error { return validation.Validate( "", diff --git a/internal/config/cloudapi.go b/internal/config/cloudapi.go index 41756fb..d308482 100644 --- a/internal/config/cloudapi.go +++ b/internal/config/cloudapi.go @@ -54,7 +54,7 @@ type CloudAPI struct { // 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. A nil receiver is +// 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 // branch on whether the block is present. // @@ -121,8 +121,7 @@ func (c CloudAPI) IsSaasAPI() bool { // 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 reported at startup instead, which mirrors how a -// namespace Cloud would reject is handled for a templated upstream. +// that apart from a typo. 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/codecserver.go b/internal/config/codecserver.go index bcf6812..ce3d769 100644 --- a/internal/config/codecserver.go +++ b/internal/config/codecserver.go @@ -31,7 +31,8 @@ type ( // disabled block is not checked at all, so a half-written one can be left in // place. An enabled one reachable beyond loopback requires authentication, and // requires TLS once it has any, because a browser will not send a token over -// plaintext. +// plaintext. CORS credentials require explicit origins, and a "*" origin is +// rejected. func (c *CodecServer) Validate() error { if !c.Enabled { return nil diff --git a/internal/config/config.go b/internal/config/config.go index 24ff808..78a3d57 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -33,7 +33,7 @@ type ( // Load reads and parses the YAML config specified in the Reader. // Values of the form ${VAR} are replaced with the corresponding environment // variable. A config that names no allowed services gets the default set, and -// one that leaves a metrics field empty gets that field's default. +// an empty metrics hostPort or namespace gets its default. func Load(r io.Reader) (*Config, error) { data, err := io.ReadAll(r) if err != nil { diff --git a/internal/config/encryption.go b/internal/config/encryption.go index fef49ac..971391d 100644 --- a/internal/config/encryption.go +++ b/internal/config/encryption.go @@ -67,7 +67,7 @@ func (e *Encryption) DEKCacheSize() int { // Validate requires a non-negative cache size, a Default policy whenever // encryption is Enabled, and (when a Default is present at all) that the policy -// itself is valid. +// itself is valid. Every Overrides entry needs a namespace and a valid policy. func (e *Encryption) Validate() error { rules := []validation.Rule{ // Checked through the accessor rather than the field, so an absent size is diff --git a/internal/config/extensions.go b/internal/config/extensions.go index bec6acc..3e202b6 100644 --- a/internal/config/extensions.go +++ b/internal/config/extensions.go @@ -8,11 +8,12 @@ import ( type ( // ExtensionServer addresses an operator-run gRPC server implementing the - // extension APIs under api/, currently api.kms.v1.EncryptionService, the - // pluggable Key Encryption Key provider. The proxy has built-in KMS providers - // (awskms, azurekeyvault, gcpkms); an extension server is how an operator plugs - // in a backend the proxy does not support natively, such as an on-prem HSM or - // an internal key service. + // extension APIs under api/: api.kms.v1.EncryptionService, the pluggable Key + // Encryption Key provider, or api.auth.v1.AuthService, the pluggable inbound + // authenticator. The proxy has built-in KMS providers (awskms, azurekeyvault, + // gcpkms) and authenticators; an extension server is how an operator plugs in + // 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, diff --git a/internal/config/fx.go b/internal/config/fx.go index 25e93c7..46059c7 100644 --- a/internal/config/fx.go +++ b/internal/config/fx.go @@ -2,6 +2,8 @@ package config import "go.uber.org/fx" +// ConfigFileTag annotates a provided string as the named value "configFile", +// the config file path [Module] loads. var ConfigFileTag = fx.ResultTags(`name:"configFile"`) // Module is an fx module that provides *Config by loading the file path supplied diff --git a/internal/config/routing.go b/internal/config/routing.go index a095d0d..c5c7416 100644 --- a/internal/config/routing.go +++ b/internal/config/routing.go @@ -8,9 +8,9 @@ import ( type ( // Routing selects which upstream serves a request. DefaultUpstream is the - // fallback when no rule matches and SystemUpstream serves system-namespace - // traffic; both name an upstream and are optional. Rules are evaluated in - // order against the incoming request. + // fallback when no rule matches and SystemUpstream serves a request that + // carries no namespace and matches no rule; both name an upstream and are + // optional. Rules are evaluated in order against the incoming request. Routing struct { DefaultUpstream string `yaml:"default"` SystemUpstream string `yaml:"system"` diff --git a/internal/config/upstream.go b/internal/config/upstream.go index 2bf7cf5..a82d430 100644 --- a/internal/config/upstream.go +++ b/internal/config/upstream.go @@ -18,7 +18,8 @@ type ( // Cloud declares the upstream to be Temporal Cloud, which turns on // Cloud-specific namespace rules. It is only needed for an address // [cloud.IsEndpoint] does not recognize, such as a private-link hostname; a - // .tmprl.cloud address is detected without it. + // .tmprl.cloud address, or a TLS server name that is one, is detected + // without it. // // The proxy dials an upstream over TLS unless Listen says otherwise, so a // plaintext upstream must set its Insecure field. @@ -31,6 +32,8 @@ type ( Connection ConnectionConfig `yaml:"connection"` } + // UpstreamList is the configured set of upstreams, named so the checks that + // span the whole collection live alongside the per-entry checks. UpstreamList []Upstream // NamespaceConfig groups the namespace translation rules for an upstream. @@ -64,8 +67,9 @@ type ( } ) -// Validate checks the upstream name, dial target, namespace, and connection -// configuration. +// Validate checks the upstream name, dial target, outbound TLS, namespace, and +// connection configuration. Credentials require TLS, insecure conflicts with a +// tls block, and a Cloud upstream must use Cloud namespace names. // A templated hostPort (containing a text/template action) is resolved // per-request, so it is not checked as a literal host:port here; a static // hostPort still is. @@ -161,6 +165,8 @@ func (u *Upstream) cloudRules() []validation.Rule { } } +// Validate checks every upstream and requires names and hostPorts to be unique +// across the list. func (ul UpstreamList) Validate() error { names := make([]string, len(ul)) hostPorts := make([]string, len(ul)) @@ -209,6 +215,8 @@ func (r *NamespaceRules) Remote(localNS string) string { return fmt.Sprintf("%s%s%s", r.Prefix, localNS, r.Suffix) } +// UnmarshalYAML decodes the rules and builds the override lookup maps, so +// overrides only take effect on rules decoded from YAML. func (r *NamespaceRules) UnmarshalYAML(unmarshal func(any) error) error { type raw NamespaceRules diff --git a/internal/dataplane/dataplane.go b/internal/dataplane/dataplane.go index 9b3217a..c64afa1 100644 --- a/internal/dataplane/dataplane.go +++ b/internal/dataplane/dataplane.go @@ -421,8 +421,8 @@ func (r keyedResolver) Resolve(ctx context.Context) (string, string, []grpc.Dial // translates reports whether any upstream will have method translation // installed, which is any upstream being Temporal Cloud. It is the same question -// perUpstream asks of one upstream, so a configuration this answers false for -// installs nothing anywhere. +// newUpstreamForwarder asks of one upstream, so a configuration this answers +// false for installs nothing anywhere. func translates(cfg *config.Config) bool { return slices.ContainsFunc(cfg.Upstreams, func(up config.Upstream) bool { return up.IsCloud() diff --git a/internal/dataplane/dataplanetest/dataplanetest.go b/internal/dataplane/dataplanetest/dataplanetest.go index ed16dc2..e7d70e6 100644 --- a/internal/dataplane/dataplanetest/dataplanetest.go +++ b/internal/dataplane/dataplanetest/dataplanetest.go @@ -205,10 +205,9 @@ func StartApp(t *testing.T, cfg *config.Config) *Fixture { return newFixture(t, dp, reg, codecSvr) } -// Gatherer is the registry every collector in this plane registered with, and -// the one its /metrics handler serves. A test asserts against it directly -// because [applyDefaults] binds that handler to an ephemeral port nothing -// reports. +// Gatherer is the registry every collector in this plane registered with. A +// test asserts against it directly: [Start] serves no /metrics endpoint, and +// the one [StartApp] wires is bound to an ephemeral port nothing reports. func (f *Fixture) Gatherer() prometheus.Gatherer { return f.reg } // Addr is the address the gateway is accepting on. diff --git a/internal/dataplane/dataplanetest/upstream.go b/internal/dataplane/dataplanetest/upstream.go index a3c5ab5..7ca5b15 100644 --- a/internal/dataplane/dataplanetest/upstream.go +++ b/internal/dataplane/dataplanetest/upstream.go @@ -48,9 +48,9 @@ func NewUpstream(t *testing.T) *Upstream { return newUpstream(t, nil) } -// NewTLSUpstream starts a fake frontend over TLS. Its [Upstream.TLSConfig] -// carries the CA and client identity needed to dial it, which is the only way -// to exercise credentials that refuse to travel over an insecure transport. +// NewTLSUpstream starts a fake frontend over TLS. Its [Upstream.Listen] carries +// the CA and client identity needed to dial it, which is the only way to +// exercise credentials that refuse to travel over an insecure transport. func NewTLSUpstream(t *testing.T) *Upstream { t.Helper() diff --git a/internal/dataplane/doc.go b/internal/dataplane/doc.go index dba0ffa..df68e8e 100644 --- a/internal/dataplane/doc.go +++ b/internal/dataplane/doc.go @@ -1,5 +1,5 @@ // Package dataplane assembles the proxy's request path: one inbound gateway -// that routes by namespace, and one proxy per upstream that translates +// that routes by namespace, and one forwarder per upstream that translates // namespaces, attaches outbound credentials, and optionally encrypts payloads // before forwarding to a Temporal Service. package dataplane diff --git a/internal/dataplane/fx.go b/internal/dataplane/fx.go index be808a1..d0bb56b 100644 --- a/internal/dataplane/fx.go +++ b/internal/dataplane/fx.go @@ -16,10 +16,8 @@ import ( "github.com/temporalio/temporal-proxy/pkg/logger" ) -// Module provides a [Dataplane] from the assembled application and binds -// Start and Stop to the fx lifecycle. It replaces the router, proxy, and -// server modules: this is the only place in the graph that owns the -// gateway/proxy topology. +// Module provides a [Dataplane] and its [proxy.Codecs] from the assembled +// application and binds Start and Stop to the fx lifecycle. var Module = fx.Options( fx.Provide(newFromParams), fx.Provide(func(d *Dataplane) *proxy.Codecs { return d.Codecs() }), diff --git a/internal/dataplane/lifecycle.go b/internal/dataplane/lifecycle.go index d12dda5..c229564 100644 --- a/internal/dataplane/lifecycle.go +++ b/internal/dataplane/lifecycle.go @@ -11,10 +11,11 @@ import ( ) // Start opens every static upstream connection so an unreachable one fails -// startup, then binds and serves the gateway. It returns once the gateway is accepting. ctx bounds -// startup only and should carry a deadline, since it is what limits the wait for -// an upstream to answer; the serving goroutines get the Context passed to New -// instead. A failure part-way through stops whatever already started. +// startup, then binds and serves the gateway. It returns once the gateway is +// accepting. ctx bounds startup only and should carry a deadline, since it is +// what limits the wait for an upstream to answer; the gateway's serving +// goroutine gets the Context passed to New instead. A failure part-way through +// stops whatever already started. func (d *Dataplane) Start(ctx context.Context) error { // A static upstream's connection is created during New, but gRPC does not // open a socket until it is used, so open them here: an unreachable upstream diff --git a/internal/kms/doc.go b/internal/kms/doc.go index 94d5325..92ff1b3 100644 --- a/internal/kms/doc.go +++ b/internal/kms/doc.go @@ -14,10 +14,12 @@ // It also meters what it opens, so each key's wraps and unwraps are recorded // against its provider. // -// The package exposes a single [Module] for Uber fx. When encryption is -// disabled the module provides a nil *crypto.Vault and starts no background -// work; when enabled it also runs a goroutine that periodically refreshes the -// vault so DEKs rotate ahead of expiry. +// The package exposes a single [Module] for Uber fx. With no default key policy +// the module provides a nil *crypto.Vault and starts no background work. With +// one, it builds a vault even when encryption is disabled, so payloads sealed +// earlier can still be opened; when encryption is enabled it also runs a +// goroutine that periodically refreshes the vault so DEKs rotate ahead of +// expiry. // // The package also reports encryption telemetry to Prometheus under the // "encryption" subsystem. A [Reporter] implements [crypto.Observer] to record diff --git a/internal/kms/fx.go b/internal/kms/fx.go index 36e5866..05bc4fb 100644 --- a/internal/kms/fx.go +++ b/internal/kms/fx.go @@ -263,7 +263,7 @@ func keyPolicyRegistryOpts( return nil, err } - // NB: KeyConfig.URI is required and therefore this will never be out of bounds. + // NB: KeyPolicy.URI is required and therefore this will never be out of bounds. if asDefault { opts = append(opts, crypto.WithDefaultKey(keys[0])) } else { diff --git a/internal/metrics/fx.go b/internal/metrics/fx.go index 871503f..cbaef9c 100644 --- a/internal/metrics/fx.go +++ b/internal/metrics/fx.go @@ -21,10 +21,8 @@ import ( // config names. Any configured fixed labels are stamped onto the Factory's // registerer rather than onto each collector, so every collector declared // through it carries them and the runtime's own go_* and process_* series, -// which register directly, do not. Consumers inject the [Factory] to declare their collectors, -// which auto-register under the configured namespace, and should pre-resolve -// labeled handles once at setup rather than per request to keep the emit path -// lock-free and allocation-free. +// which register directly, do not. Consumers inject the [Factory] to declare +// their collectors, which auto-register under the configured namespace. // // The HTTP server is bound to the fx lifecycle: it starts in a background // goroutine on OnStart and shuts down gracefully on OnStop. If the server diff --git a/internal/metrics/labels.go b/internal/metrics/labels.go index 058587a..e2074a9 100644 --- a/internal/metrics/labels.go +++ b/internal/metrics/labels.go @@ -20,7 +20,7 @@ const maxLabelValueLen = 256 // MetadataLabels is the ordered set of inbound metadata headers reported as // extra labels on request-scoped collectors. The zero MetadataLabels carries // none, which is what a deployment that configures none runs, so a reporter -// holding one keeps its label set and its emit path exactly as they were. +// holding one adds no labels and reads no metadata. type MetadataLabels struct { names []string // Prometheus label names, in configured order. headers []string // Lowercased metadata keys, parallel to names. @@ -52,7 +52,7 @@ func NewMetadataLabels(cfg []config.MetricLabel) MetadataLabels { // through it, so a constant an operator configures once reaches every series the // proxy publishes rather than only the ones emitted while serving a request. It // returns r unchanged when there are none, so a deployment configuring no fixed -// labels registers exactly what it did before. +// labels registers collectors with only their own labels. func WithFixedLabels(r prometheus.Registerer, labels map[string]string) prometheus.Registerer { if len(labels) == 0 { return r diff --git a/internal/protoutil/translate.go b/internal/protoutil/translate.go index dad1dd8..271c7e1 100644 --- a/internal/protoutil/translate.go +++ b/internal/protoutil/translate.go @@ -48,9 +48,9 @@ func (t *Translator) WarmService(name protoreflect.FullName) error { return nil } -// Translate rewrites every namespace name in m using fn. It is a no-op when m is -// nil, invalid, or carries no namespace field. fn maps a namespace name to its -// translated form (local to remote, or remote to local). +// Translate rewrites every namespace name in m using fn. It is a no-op when m or +// fn is nil, or m is invalid or carries no namespace field. fn maps a namespace +// name to its translated form (local to remote, or remote to local). func (t *Translator) Translate(m proto.Message, fn func(string) string) { if m == nil || fn == nil { return diff --git a/internal/proxy/codec.go b/internal/proxy/codec.go index b2f485d..a8ba03a 100644 --- a/internal/proxy/codec.go +++ b/internal/proxy/codec.go @@ -46,7 +46,8 @@ type ( // NewCodecs returns the [Codecs] opts select. Encoding is gated per codec; // decoding is not, so a decoder recognizes its own output and passes anything -// else through. +// else through. It errors when Encrypt is set without a Vault, or a Vault is +// set without a Reporter. func NewCodecs(opts CodecOptions) (*Codecs, error) { if opts.Encrypt && opts.Vault == nil { return nil, errors.New("proxy: encryption requires a vault") diff --git a/internal/proxy/forward.go b/internal/proxy/forward.go index bd81427..6efbbe4 100644 --- a/internal/proxy/forward.go +++ b/internal/proxy/forward.go @@ -224,8 +224,8 @@ func (f *Forwarder) resolveMethod(fullMethod string) *methodInfo { // forwardContext turns the inbound request's metadata into outgoing metadata for // the upstream call, minus the transport headers, without displacing a value // already set on the outgoing context. Templated upstream resolution reads the -// router-stamped namespace from there, so this is load-bearing rather than -// merely polite. +// caller's headers for its Metadata from there, so this is load-bearing rather +// than merely polite. func forwardContext(ctx context.Context) context.Context { incoming, ok := metadata.FromIncomingContext(ctx) if !ok { diff --git a/internal/proxy/reporter.go b/internal/proxy/reporter.go index 2c32bb3..84d60dd 100644 --- a/internal/proxy/reporter.go +++ b/internal/proxy/reporter.go @@ -10,13 +10,14 @@ import ( type ( // Reporter records envelope-operation telemetry to Prometheus: each seal - // (encrypt) and open (decrypt) the encryption interceptor performs, timed end - // to end, including any KEK wrap or unwrap and any DEK cache lookup along the - // way. The AES-step duration alone is owned by internal/kms. The namespace - // label is always declared but only carries a value when namespace labels are - // enabled, since the set of namespaces is unbounded; handles are resolved per - // call via WithLabelValues rather than pre-computed. A Reporter is safe for - // concurrent use. + // (encrypt) and open (decrypt) the encryption codec performs, on both the + // gRPC path and the codec server's HTTP path, timed end to end, including + // any KEK wrap or unwrap and any DEK cache lookup along the way. The AES-step + // duration alone is owned by internal/kms. The namespace label is always + // declared but only carries a value when namespace labels are enabled, since + // the set of namespaces is unbounded; handles are resolved per call via + // WithLabelValues rather than pre-computed. A Reporter is safe for concurrent + // use. Reporter struct { ops *prometheus.CounterVec duration *prometheus.HistogramVec @@ -24,6 +25,7 @@ type ( labels metrics.MetadataLabels } + // ReporterOption configures a Reporter at construction. ReporterOption func(*Reporter) ) @@ -48,18 +50,21 @@ func NewReporter(f *metrics.Factory, opts ...ReporterOption) *Reporter { return r } +// WithNamespaceLabels sets whether the namespace label carries a value. Off by +// default, since the set of namespaces is unbounded. func WithNamespaceLabels(enabled bool) ReporterOption { return func(r *Reporter) { r.namespaceLabels = enabled } } +// WithMetadataLabels sets the request-metadata headers reported as extra labels. func WithMetadataLabels(labels metrics.MetadataLabels) ReporterOption { return func(r *Reporter) { r.labels = labels } } // VaultOp records a single envelope operation and its duration. ctx is the // request's, and supplies the configured metadata label values when there are -// any: on the per-upstream hop that is what the gateway forwarded rather than -// what it received, so a header the inbound authenticator consumed is gone, +// any: on the gRPC path that is the incoming metadata after the inbound +// authenticator stripped the headers it consumed, so those headers are gone, // and on the codec server's HTTP path there is no gRPC metadata at all, so // those labels come out blank there. func (r *Reporter) VaultOp(ctx context.Context, operation, result, namespace string, seconds float64) { diff --git a/internal/proxy/resolver.go b/internal/proxy/resolver.go index 7d0bebc..99df779 100644 --- a/internal/proxy/resolver.go +++ b/internal/proxy/resolver.go @@ -92,7 +92,8 @@ func WithRemoteNamespacer(f func(string) string) ResolverOption { } // WithOptionsFactory sets the function that produces the dial options for a -// resolved request. It receives the rendered host and server name via RouteData. +// resolved request. It receives the template context and the rendered server +// name via RouteData. func WithOptionsFactory(f func(RouteData) ([]grpc.DialOption, error)) ResolverOption { return func(r *DynamicResolver) { r.opts = f } } @@ -103,15 +104,15 @@ func WithResolverLogger(l logger.Logger) ResolverOption { return func(r *DynamicResolver) { r.logger = l } } -// ResolverFor builds the [connect.Resolver] for an upstream. When neither -// the hostPort nor the TLS server name is templated it returns a static -// resolver, whose connection is constructed while the graph is built, opened on -// start, and reused for every request; otherwise it returns a DynamicResolver -// that renders the target and server name, and rebuilds credentials, per request. +// ResolverFor builds the [connect.Resolver] for an upstream. When neither the +// hostPort nor the TLS server name is templated it returns a static resolver, +// whose connection is constructed while the graph is built, opened on start, +// and reused for every request; otherwise it returns a DynamicResolver that +// renders the target and server name, and rebuilds credentials, per request. // opts holds the request-independent dial options (namespace translation and -// outbound credentials); the upstream's connection settings and the srv:/// resolver are appended here. -// log, when non-nil, is threaded into the DynamicResolver for per-request debug -// entries. +// outbound credentials); the upstream's connection settings and the srv:/// +// resolver are appended here. log, when non-nil, is threaded into the +// DynamicResolver for per-request debug entries. func ResolverFor(upstream *config.Upstream, opts []grpc.DialOption, log logger.Logger) (connect.Resolver, error) { // One Dialer per upstream owns the TLS-mode decision and parses its // certificate material once, so a templated upstream reuses it across every diff --git a/internal/proxy/translation.go b/internal/proxy/translation.go index ae4bacb..aab741c 100644 --- a/internal/proxy/translation.go +++ b/internal/proxy/translation.go @@ -31,9 +31,10 @@ type translatingClientStream struct { // TranslationDialOptions returns the dial options that install namespace // translation on the outbound connection: t rewrites message bodies and typed -// error details, out maps local names to remote on the way out, and in maps -// remote names to local on the way back. Callers fold them into the dial -// options for the upstream connection. +// error details, out maps local names to remote on the way out (including the +// outgoing temporal-namespace header), and in maps remote names to local on the +// way back. Callers fold them into the dial options for the upstream +// connection. func TranslationDialOptions(t *protoutil.Translator, out, in func(string) string) []grpc.DialOption { return []grpc.DialOption{ grpc.WithChainUnaryInterceptor(unaryClientInterceptor(t, out, in)), diff --git a/internal/router/codec.go b/internal/router/codec.go index eb2a261..646012c 100644 --- a/internal/router/codec.go +++ b/internal/router/codec.go @@ -23,9 +23,10 @@ type ( } ) -// Codec returns the hybrid pass-through codec. It must be applied per-call via -// grpc.ForceServerCodecV2 / grpc.ForceCodecV2; it is deliberately not registered -// globally so it never shadows the real proto codec process-wide. +// Codec returns the hybrid pass-through codec. It must be installed on the +// gateway server via grpc.ForceServerCodecV2 (see server.WithServerCodec); it is +// deliberately not registered globally so it never shadows the real proto codec +// process-wide. func Codec() encoding.CodecV2 { return defaultCodec } diff --git a/internal/router/mux.go b/internal/router/mux.go index 6f1f0de..cf5ae1d 100644 --- a/internal/router/mux.go +++ b/internal/router/mux.go @@ -20,7 +20,7 @@ const ( type ( // Mux selects the upstream that serves a request by matching it against an // ordered list of rules. It holds upstream names only, not connections, so - // callers map the name Switch returns to a connection. A Mux is read-only + // callers map the name Switch returns to a stream handler. A Mux is read-only // after construction and safe for concurrent use. Mux struct { def string diff --git a/internal/router/reporter.go b/internal/router/reporter.go index ac04a5b..8720b1c 100644 --- a/internal/router/reporter.go +++ b/internal/router/reporter.go @@ -27,9 +27,9 @@ type ( // 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. + // 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. Reporter struct { decisions *prometheus.CounterVec errors *prometheus.CounterVec @@ -56,9 +56,9 @@ type ( // 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. +// 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. func NewReporter(f *metrics.Factory, upstreams []string, labels metrics.MetadataLabels) *Reporter { decisions := f.NewCounter(prometheus.CounterOpts{ Name: "decisions_total", diff --git a/internal/rpc/doc.go b/internal/rpc/doc.go index ace630f..c2722a8 100644 --- a/internal/rpc/doc.go +++ b/internal/rpc/doc.go @@ -1,8 +1,8 @@ -// Package rpc holds the gRPC stream plumbing the proxy's forwarding paths share. +// Package rpc holds the gRPC helpers the proxy's forwarding paths share: stream +// plumbing, method-name parsing, outgoing metadata, and caller-facing errors. // -// [ServiceMethod], [Service], and [Method] split a gRPC full method name into -// the halves an allowlist check or a descriptor lookup needs. [Pump.Forward] -// forwards one call between an inbound server stream and an outbound client -// stream, and [StatusError] maps a failure along the way to the status the caller -// sees. +// [ServiceMethod] and [Service] split a gRPC full method name into the halves +// an allowlist check or a descriptor lookup needs. [Pump.Forward] forwards one +// call between an inbound server stream and an outbound client stream, and +// [StatusError] maps a failure along the way to the status the caller sees. package rpc diff --git a/internal/rpc/pump.go b/internal/rpc/pump.go index d2b8802..9eb47dc 100644 --- a/internal/rpc/pump.go +++ b/internal/rpc/pump.go @@ -71,10 +71,10 @@ func (p *Pump) Forward(in, out Frame) error { return status.Error(codes.Internal, "forwarding ended without completion") } -// requests forwards request messages the caller sends to the upstream, one in per -// message, until either side fails, and reports that first failure on the -// returned channel. io.EOF means the caller half-closed cleanly, which is the caller's cue -// to close the upstream's send side rather than to fail the call. +// requests forwards request messages the caller sends to the upstream, one in +// per message, until either side fails, and reports that first failure on the +// returned channel. io.EOF means the caller half-closed cleanly, which is the +// caller's cue to close the upstream's send side rather than to fail the call. func (p *Pump) requests(in Frame) <-chan error { errs := make(chan error, 1) @@ -100,10 +100,11 @@ func (p *Pump) requests(in Frame) <-chan error { // responses relays the upstream's response header and then forwards the // upstream's response messages to the caller, one out per message, until either -// side fails, reporting that first failure on the returned channel. io.EOF means the upstream completed -// cleanly; any other error is the upstream's status or a failure sending to the -// caller. The header is relayed before the first message because gRPC flushes it -// on the first send, and a header sent late is a header the caller never sees. +// side fails, reporting that first failure on the returned channel. io.EOF +// means the upstream completed cleanly; any other error is the upstream's +// status or a failure sending to the caller. The header is relayed before the +// first message because gRPC flushes it on the first send, and a header sent +// late is a header the caller never sees. func (p *Pump) responses(out Frame) <-chan error { errs := make(chan error, 1) diff --git a/internal/server/loopback.go b/internal/server/loopback.go index 351fbc2..48b908a 100644 --- a/internal/server/loopback.go +++ b/internal/server/loopback.go @@ -37,10 +37,10 @@ type ( // // 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. + // 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 diff --git a/internal/server/reporter.go b/internal/server/reporter.go index 0850450..8f71449 100644 --- a/internal/server/reporter.go +++ b/internal/server/reporter.go @@ -80,8 +80,9 @@ func (r *Reporter) Observe(ctx context.Context, method string, code codes.Code, // and records the RPC's duration and final gRPC status code. It covers all // forwarded traffic, which grpc-go serves through the unknown-service handler as // streams. The local health service's unary Check is not metered, since unary -// calls do not pass through a stream interceptor; its streaming Watch, if a -// client uses it, would be metered under its own method name. +// calls do not pass through a stream interceptor; its streaming Watch is +// metered under its own method name, which includes the call the loopback +// health check makes every interval when [WithLoopbackHealthCheck] is set. func (r *Reporter) StreamInterceptor() grpc.StreamServerInterceptor { return func( srv any, diff --git a/internal/server/server.go b/internal/server/server.go index e214752..2768d5a 100644 --- a/internal/server/server.go +++ b/internal/server/server.go @@ -1,3 +1,5 @@ +// Package server provides the gRPC server the gateway listens on, with a +// built-in health service and a periodic health check. package server import ( @@ -94,8 +96,8 @@ type ( ) // New constructs a [Server]. When no options are supplied, it uses insecure -// credentials, a default health check that always reports SERVING, a CLI -// logger, and a five second drain budget. +// credentials, a default health check that always reports SERVING, +// [logger.Default], and a five second drain budget. func New(sopts ...Option) (*Server, error) { opts := &options{ creds: creds.NewListener(creds.Insecure()), @@ -256,9 +258,8 @@ func (s *Server) Start(ctx context.Context, lis net.Listener) error { // Stop shuts the server down, halting the health check loop and draining // in-flight RPCs. The drain is bounded by whichever expires first: the // [WithShutdownTimeout] budget or ctx. Past that, remaining calls are dropped. -// A forced shutdown is still a shutdown, so it is reported through a warning -// rather than an error; the only errors here would be a caller's to handle, and -// there are none. +// It always returns nil; a forced shutdown is logged as a warning rather than +// reported as an error. func (s *Server) Stop(ctx context.Context) error { s.mu.Lock() log := s.logger diff --git a/internal/services/services.go b/internal/services/services.go index 21b07d4..b0dae3e 100644 --- a/internal/services/services.go +++ b/internal/services/services.go @@ -33,7 +33,7 @@ const ( ) // aliases maps a service name to the additional names allowing it implies. -// Callers name the service they mean; Expand supplies the compatibility +// Callers name the service they mean; expand supplies the compatibility // spellings so configuration does not have to. var aliases = map[string][]string{ Reflection: {ReflectionV1Alpha}, @@ -46,9 +46,9 @@ func Default() []string { return []string{WorkflowService, OperatorService} } -// Known returns every service the proxy can forward, which is every service -// whose descriptors this package links in. It is the universe configuration may -// select from, and the set the namespace completeness guard audits. +// Known returns every service the proxy can forward, in its canonical spelling; +// compatibility aliases such as v1alpha reflection are left to All. It is the +// universe configuration may select from. func Known() []string { return []string{WorkflowService, OperatorService, Reflection} } diff --git a/internal/template/template.go b/internal/template/template.go index 73cae9f..4c2bcdb 100644 --- a/internal/template/template.go +++ b/internal/template/template.go @@ -6,6 +6,9 @@ import ( texttemplate "text/template" ) +// probeMeta is the sample metadata a template is rendered against at parse +// time, so a template that cannot execute fails when parsed rather than per +// request. var probeMeta = map[string]string{ "dc": "probe", "x-cluster": "probe", diff --git a/internal/transport/connect/doc.go b/internal/transport/connect/doc.go index c3d426b..6abe208 100644 --- a/internal/transport/connect/doc.go +++ b/internal/transport/connect/doc.go @@ -1,8 +1,8 @@ -// Package connect manages a pool of reusable gRPC client connections keyed by -// a caller-supplied logical key, distinct from the dial target. Use NewPool to +// Package connect manages a pool of reusable gRPC client connections keyed by a +// caller-supplied logical key, distinct from the dial target. Use NewPool to // create a Pool, Set to register a connection for a key, Conn to retrieve it, -// ConnOrCreate to retrieve or create one on first use, and Close to shut every connection -// down exactly once. +// ConnOrCreate to retrieve or create one on first use, and Close to shut every +// connection down exactly once. // // Conn wraps that pool as a [grpc.ClientConnInterface] whose dial target is // chosen per call by a Resolver: StaticResolver for a fixed upstream, or a diff --git a/internal/transport/connect/pool.go b/internal/transport/connect/pool.go index 1d96846..f88ebb2 100644 --- a/internal/transport/connect/pool.go +++ b/internal/transport/connect/pool.go @@ -54,13 +54,14 @@ func (p *Pool) Conn(key string) (*grpc.ClientConn, error) { return cn, nil } -// ConnOrCreate returns the connection registered for key, creating and registering -// one with grpc.NewClient(target, opts...) when none exists yet. key is the -// logical cache key and target is the dial address; callers that need distinct -// connections to the same target (e.g. identical host:port with different TLS -// server names) must pass distinct keys. If callers race to create the same -// key, each constructs a client but only one connection is kept; the losers are closed and -// every caller receives the same *grpc.ClientConn. +// ConnOrCreate returns the connection registered for key, creating and +// registering one with grpc.NewClient(target, opts...) when none exists yet. +// key is the logical cache key and target is the dial address; callers that +// need distinct connections to the same target (e.g. identical host:port with +// different TLS server names) must pass distinct keys. If callers race to +// create the same key, each constructs a client but only one connection is +// kept; the losers are closed and every caller receives the same +// *grpc.ClientConn. func (p *Pool) ConnOrCreate(key, target string, opts ...grpc.DialOption) (*grpc.ClientConn, error) { if conn, _ := p.Conn(key); conn != nil { return conn, nil diff --git a/internal/transport/creds/doc.go b/internal/transport/creds/doc.go index edcbdbd..c065312 100644 --- a/internal/transport/creds/doc.go +++ b/internal/transport/creds/doc.go @@ -12,6 +12,7 @@ // Security is the default: only [Insecure] yields a plaintext credential, and an // accidentally-empty client credential verifies the peer against the system root // pool rather than silently downgrading. Construction performs no file I/O; -// certificate material is read and parsed lazily (and once) when a credential is -// validated or used. +// certificate material is read and parsed lazily when a credential is validated +// or used. A Dialer parses it once and reuses it across [Dialer.DialOption] +// calls, while [Listener.TLSConfig] re-reads it on every call. package creds diff --git a/internal/transport/creds/options.go b/internal/transport/creds/options.go index 3d67ab4..22f299f 100644 --- a/internal/transport/creds/options.go +++ b/internal/transport/creds/options.go @@ -25,7 +25,7 @@ func (f optFunc) apply(o *options) { f(o) } // Insecure disables transport security. It is the only way to obtain a // plaintext credential: security is the default, so an accidentally-empty // credential fails toward TLS rather than silently downgrading. Use it -// deliberately, for example on the local loopback socket. +// deliberately, for example to reach a local development Temporal Service. func Insecure() Option { return optFunc(func(o *options) { o.insecure = true }) } diff --git a/internal/transport/meta/meta.go b/internal/transport/meta/meta.go index 058b7ac..b0e5f4e 100644 --- a/internal/transport/meta/meta.go +++ b/internal/transport/meta/meta.go @@ -1,7 +1,7 @@ // Package meta defines the internal contract for what the gateway learns about a // request once and every later stage reads: the [Target] on the context, and the -// namespace stamped on outgoing metadata for the per-upstream proxy. It depends -// on no other internal packages. +// namespace stamped on outgoing metadata for the forwarder's resolver and client +// interceptors. It depends on no other internal packages. package meta import ( @@ -12,7 +12,8 @@ import ( const ( // NamespaceHeader is the outgoing metadata key that carries the local (pre- - // translation) namespace from the router to the upstream proxy. + // translation) namespace from the router to the forwarder's resolver and + // client interceptors. NamespaceHeader = "x-temporal-proxy-namespace" // VersionHeader is the outgoing metadata key that carries the proxy's own diff --git a/pkg/api/ext/v1/key_material.go b/pkg/api/ext/v1/key_material.go index 56cbaa8..eb20d0a 100644 --- a/pkg/api/ext/v1/key_material.go +++ b/pkg/api/ext/v1/key_material.go @@ -1,3 +1,6 @@ +// Package ext holds [KeyMaterial], an optional framing an extension server can +// marshal as the ciphertext it returns from Encrypt, along with helpers to +// validate, marshal, and unmarshal it. package ext import ( @@ -37,9 +40,9 @@ func (km *KeyMaterial) Marshal() ([]byte, error) { return packed, nil } -// Validate reports whether km carries a wrapped DEK, treating a nil km as one -// that does not. Every other field is optional: an extension server may version -// no keys, key off no namespace, and carry nothing of its own. +// Validate returns an error unless km carries a wrapped DEK, treating a nil km +// as one that does not. Every other field is optional: an extension server may +// version no keys, key off no namespace, and carry nothing of its own. func (km *KeyMaterial) Validate() error { return validation.Validate( "", diff --git a/pkg/codec/chain.go b/pkg/codec/chain.go index 3d0130e..ffcccee 100644 --- a/pkg/codec/chain.go +++ b/pkg/codec/chain.go @@ -16,8 +16,8 @@ type ( } // Chain applies a set of codecs as one. Encode runs them in the order the - // chain holds them and Decode runs them in reverse, so a payload is - // compressed before it is sealed and unsealed before it is decompressed. + // chain holds them and Decode runs them in reverse, so the last codec to + // encode a payload is the first to decode it. // Note that this is the opposite of the SDK's own convention, where a // multi-codec list encodes last to first; a Chain is handed to the SDK whole, // as a single codec, so its order stays its own concern. diff --git a/pkg/codec/encryptor.go b/pkg/codec/encryptor.go index 03fc985..afacfbb 100644 --- a/pkg/codec/encryptor.go +++ b/pkg/codec/encryptor.go @@ -32,8 +32,8 @@ type ( cipher Cipher } - // Cipher encrypts and decrypts bytes. It is the subset of a key-management - // backend, typically a [crypto.Vault], that [Encryptor] depends on. + // Cipher encrypts and decrypts bytes. It is what [Encryptor] depends on, + // typically a wrapper that binds a [crypto.Vault] to a namespace. Cipher interface { Encrypt([]byte) (*crypto.Message, error) Decrypt(*crypto.Message) ([]byte, error) diff --git a/pkg/crypto/dek.go b/pkg/crypto/dek.go index 8b88850..bee3192 100644 --- a/pkg/crypto/dek.go +++ b/pkg/crypto/dek.go @@ -28,7 +28,7 @@ type ( // DEKMaterial defines the material needed in order to decrypt a payload. DEKMaterial struct { Version byte - KEKID string // The ID/URI of the KEK the encrypted the DEK. + KEKID string // The ID/URI of the KEK that encrypted the DEK. EncryptedDEK string // The base64-encoded encrypted DEK. } ) @@ -55,7 +55,9 @@ func (d *DEK) Encrypt(ctx context.Context, pt []byte) ([]byte, error) { } // Decrypt decrypts the ciphertext ct using AES-256-GCM. The ciphertext must be -// prefixed with the nonce, as produced by [DEK.Encrypt]. +// prefixed with the nonce, as produced by [DEK.Encrypt]. It returns an error +// matching [ErrMalformedCipherText] when ct is shorter than the nonce or fails +// authentication. func (d *DEK) Decrypt(ctx context.Context, ct []byte) ([]byte, error) { ns := d.gcm.NonceSize() if len(ct) < ns { diff --git a/pkg/crypto/kek.go b/pkg/crypto/kek.go index a2f6ce6..7976938 100644 --- a/pkg/crypto/kek.go +++ b/pkg/crypto/kek.go @@ -11,7 +11,7 @@ import ( ) type ( - // KEK defines an interface for a Key Encryption Keys. + // KEK defines an interface for a Key Encryption Key. // These keys are used to encrypt/decrypt DEKs and are customer-managed (e.g. via AWS/GCP KMS). KEK interface { io.Closer @@ -30,6 +30,8 @@ type ( // KEKRegistry holds the set of KEKs available for encrypting and decrypting DEKs. // It is keyed by namespace (for encryption) and by key ID (for decryption). // Close must be called when the registry is no longer needed to release KEK resources. + // It is not modified after construction, so it is safe for concurrent use as + // long as its KEKs are. KEKRegistry struct { defaultKey KEK // Fallback when no namespace key exists. keks map[string]KEK // map from NS -> KEK @@ -50,7 +52,9 @@ type ( // NewKEKRegistry constructs a KEKRegistry, applying opts in order. A default key is // required (see [WithDefaultKey]); construction fails if one is not provided. The -// key-ID index used by Decrypt is built after all options are applied. +// key-ID index used by Decrypt is built after all options are applied, and +// construction fails if any key ID appears more than once across the default, +// namespace, and decrypt-only keys. func NewKEKRegistry(opts ...KEKRegistryOption) (*KEKRegistry, error) { r := &KEKRegistry{ keks: map[string]KEK{}, @@ -109,7 +113,8 @@ func WithDefaultKey(k KEK) KEKRegistryOption { }) } -// WithKeyForNamespace registers k for ns, used when encrypting or decrypting DEKs for that namespace. +// WithKeyForNamespace registers k for ns, used when encrypting DEKs for that namespace. +// Decryption selects the key by KEK ID, not by namespace. func WithKeyForNamespace(ns string, k KEK) KEKRegistryOption { return kekRegOpt(func(r *KEKRegistry) error { if k == nil { @@ -144,7 +149,8 @@ func WithDecryptOnlyKey(k KEK) KEKRegistryOption { }) } -// Encrypt encrypts the given DEK using the KEK registered for the specified namespace. +// Encrypt encrypts the given DEK using the KEK registered for the specified namespace, +// falling back to the default key when none is registered. // It returns DEKMaterial containing the KEK ID and the base64-encoded ciphertext. func (r *KEKRegistry) Encrypt(ctx context.Context, ns string, dek *DEK) (*DEKMaterial, error) { if dek == nil { diff --git a/pkg/crypto/keys.go b/pkg/crypto/keys.go index cc3912e..26fff4a 100644 --- a/pkg/crypto/keys.go +++ b/pkg/crypto/keys.go @@ -8,6 +8,7 @@ import ( "strings" "gocloud.dev/secrets" + // Register the gocloud keeper URL openers NewCloudKey relies on. _ "gocloud.dev/secrets/awskms" _ "gocloud.dev/secrets/azurekeyvault" _ "gocloud.dev/secrets/gcpkms" @@ -183,7 +184,8 @@ func (f *KeyFactory) Create(ctx context.Context, uri string) (KEK, error) { return fn(ctx, uri) } -// ID returns a unique ID for this KEK, e.g. a KMS ARN. +// ID returns the key URI passed to [NewCloudKey], with a "testing://" scheme +// rewritten to "base64key://". func (k *CloudKey) ID() string { return k.id } diff --git a/pkg/crypto/vault.go b/pkg/crypto/vault.go index 7fe9e3b..409f1aa 100644 --- a/pkg/crypto/vault.go +++ b/pkg/crypto/vault.go @@ -36,7 +36,8 @@ type ( ns string } - // KeyConfig controls the lifetime of a namespace's DEK. + // KeyConfig controls the lifetime of a namespace's DEK. [NewVault] rejects + // a config unless Duration > 0 and 0 <= RenewBefore < Duration. KeyConfig struct { // Duration is how long a DEK is valid before it must be rotated. Duration time.Duration @@ -168,7 +169,8 @@ func WithDefaultKeyConfig(cfg KeyConfig) VaultOption { } // WithKeyConfig sets the KeyConfig for a specific namespace. Registering the -// same namespace more than once is an error surfaced by [NewVault]. +// same namespace more than once, or passing a cfg that breaks the [KeyConfig] +// constraints, is an error surfaced by [NewVault]. func WithKeyConfig(ns string, cfg KeyConfig) VaultOption { return func(e *vaultOptions) { if _, ok := e.config[ns]; ok { @@ -522,7 +524,8 @@ func (v *NamespacedVault) Seal(ctx context.Context, data []byte) (*Message, erro return v.inner.Seal(ctx, v.ns, data) } -// Open decrypts msg within the bound namespace. See [Vault.Open]. +// Open decrypts msg by delegating to the inner vault. The bound namespace is +// unused; the key is selected from msg's key material. See [Vault.Open]. func (v *NamespacedVault) Open(ctx context.Context, msg *Message) ([]byte, error) { return v.inner.Open(ctx, msg) } diff --git a/pkg/ext/binding.go b/pkg/ext/binding.go index 3259f0d..eeaf5fb 100644 --- a/pkg/ext/binding.go +++ b/pkg/ext/binding.go @@ -15,7 +15,9 @@ import ( // travel beside it in the clear, for a [KeySealer] to hand its key service as // an encryption context. Material framed by [NewSealWrapper] names no cipher, // so there is none to pass here: a sealer that wants to record which -// construction it used puts that in Opaque, which this binds. +// construction it used puts that in Opaque, which this binds. It returns an +// error when namespace or version is longer than 65535 bytes, or opaque is +// longer than [math.MaxUint32] bytes. func BindingContext(namespace, version string, opaque []byte) ([]byte, error) { return bindingContext(ext.KeyMaterial_CIPHER_UNSPECIFIED, namespace, version, opaque) } diff --git a/pkg/ext/keywrap.go b/pkg/ext/keywrap.go index 6128634..518893b 100644 --- a/pkg/ext/keywrap.go +++ b/pkg/ext/keywrap.go @@ -104,6 +104,7 @@ func NewKeyWrapper(lookup KeyLookup, opts ...KeyWrapperOption) (KMS, error) { // WithCipher seals new material with id instead of AES-256-GCM. It has no // bearing on opening, which uses whatever cipher the material names. +// [NewKeyWrapper] fails if no [CipherFunc] is registered for id. func WithCipher(id CipherID) KeyWrapperOption { return keyWrapperOpt(func(w *keyWrapper) error { w.cipher = id @@ -120,7 +121,7 @@ func WithCipher(id CipherID) KeyWrapperOption { // 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. +// 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/server.go b/pkg/ext/server.go index 475409a..2ac6280 100644 --- a/pkg/ext/server.go +++ b/pkg/ext/server.go @@ -48,7 +48,8 @@ type ( // Serve runs an extension server until ctx is cancelled or the process is // signalled, then shuts down and returns nil. A non-nil return means the server -// never started, never that a caller was turned away. +// failed to start or stopped serving with an error, never that a caller was +// turned away. // // Both generated services are registered whether or not [WithAuth] and [WithKMS] // were given, and one left unset answers Unimplemented. Defaults are :8900 on @@ -153,10 +154,10 @@ func WithKMS(kms KMS) Option { return func(o *options) { o.kms = kms } } -// WithLogger sets the logger for the server's lifecycle, defaulting to -// [logger.Default]. Handlers are not given it. A nil logger is ignored rather -// than installed, matching [logger.SetDefault]; the alternative is a panic on the -// first line the server logs. +// WithLogger sets the logger for the server's lifecycle, its plaintext warning, +// and the KMS service handler, defaulting to [logger.Default]. A nil logger is +// ignored rather than installed, matching [logger.SetDefault]; the alternative +// is a panic on the first line the server logs. func WithLogger(l logger.Logger) Option { return func(o *options) { if l != nil { diff --git a/pkg/logger/doc.go b/pkg/logger/doc.go index c03aa6b..8912a05 100644 --- a/pkg/logger/doc.go +++ b/pkg/logger/doc.go @@ -1,7 +1,8 @@ // Package logger provides a small, leveled, structured logging interface for -// the proxy along with a [zerolog]-backed implementation and a no-op -// implementation for tests. +// the proxy along with a [zerolog]-backed implementation, plus a no-op +// implementation and a recording [TestLogger] for tests. // // Package-level functions ([Debug], [Info], [Warn], [Error], [Fatal], [With]) -// delegate to a default [Logger] that writes to os.Stderr at [LevelInfo]. +// delegate to a default [Logger] that writes to os.Stderr at [LevelInfo]. The +// default is replaceable via [SetDefault]. package logger diff --git a/pkg/logger/logger.go b/pkg/logger/logger.go index 224e1c5..02aebc9 100644 --- a/pkg/logger/logger.go +++ b/pkg/logger/logger.go @@ -43,7 +43,8 @@ type ( } ) -// Default returns a default [Logger] which writes to os.Stderr at LevelInfo. +// Default returns the [Logger] last installed by [SetDefault], initially one +// that writes to os.Stderr at [LevelInfo]. func Default() Logger { return defaultLogger } diff --git a/pkg/logger/noop.go b/pkg/logger/noop.go index ffe1d3b..87fbc2b 100644 --- a/pkg/logger/noop.go +++ b/pkg/logger/noop.go @@ -4,8 +4,8 @@ import "github.com/temporalio/temporal-proxy/pkg/logger/tag" type ( // NoopLogger is a [Logger] that discards every entry. It is useful in tests - // and anywhere a non-nil Logger is required but output is not wanted. Unlike - // other implementations, its Fatal does not exit the process. + // and anywhere a non-nil Logger is required but output is not wanted. Its + // Fatal does not exit the process. NoopLogger struct{} ) @@ -14,21 +14,27 @@ func NewNoopLogger() *NoopLogger { return new(NoopLogger) } +// Debug implements [Logger] by discarding the entry. func (n *NoopLogger) Debug(string, ...tag.Tag) { } +// Error implements [Logger] by discarding the entry. func (n *NoopLogger) Error(string, ...tag.Tag) { } +// Fatal implements [Logger] by discarding the entry. It does not exit. func (n *NoopLogger) Fatal(string, ...tag.Tag) { } +// Info implements [Logger] by discarding the entry. func (n *NoopLogger) Info(string, ...tag.Tag) { } +// Warn implements [Logger] by discarding the entry. func (n *NoopLogger) Warn(string, ...tag.Tag) { } +// With implements [Logger] by returning the receiver; tags are discarded. func (n *NoopLogger) With(...tag.Tag) Logger { return n } diff --git a/pkg/logger/tag/tag.go b/pkg/logger/tag/tag.go index 3eca233..45da7df 100644 --- a/pkg/logger/tag/tag.go +++ b/pkg/logger/tag/tag.go @@ -15,7 +15,8 @@ func Component(c string) Tag { return String("component", c) } -// Error returns a Tag with key "error" carrying err's message. +// Error returns a Tag with key "error" carrying err's message. A nil err +// yields an empty string value. func Error(err error) Tag { msg := "" if err != nil { @@ -30,7 +31,8 @@ func String(k, v string) Tag { return Tag{Key: k, Value: v} } -// Stringer returns a Tag whose value is v.String(), evaluated immediately. +// Stringer returns a Tag whose value is v.String(), evaluated immediately. A +// nil v panics. func Stringer(k string, v fmt.Stringer) Tag { return String(k, v.String()) } diff --git a/pkg/logger/test.go b/pkg/logger/test.go index 268db39..f03f1ff 100644 --- a/pkg/logger/test.go +++ b/pkg/logger/test.go @@ -101,11 +101,21 @@ func (t *TestLogger) TagsOf(l Level, msg string) map[string]any { return nil } +// Debug implements [Logger] by recording the entry at [LevelDebug]. func (t *TestLogger) Debug(msg string, tags ...tag.Tag) { t.record(LevelDebug, msg, tags) } + +// Error implements [Logger] by recording the entry at [LevelError]. func (t *TestLogger) Error(msg string, tags ...tag.Tag) { t.record(LevelError, msg, tags) } + +// Fatal implements [Logger] by recording the entry at [LevelError]. It does +// not exit. func (t *TestLogger) Fatal(msg string, tags ...tag.Tag) { t.record(LevelError, msg, tags) } -func (t *TestLogger) Info(msg string, tags ...tag.Tag) { t.record(LevelInfo, msg, tags) } -func (t *TestLogger) Warn(msg string, tags ...tag.Tag) { t.record(LevelWarn, msg, tags) } + +// Info implements [Logger] by recording the entry at [LevelInfo]. +func (t *TestLogger) Info(msg string, tags ...tag.Tag) { t.record(LevelInfo, msg, tags) } + +// Warn implements [Logger] by recording the entry at [LevelWarn]. +func (t *TestLogger) Warn(msg string, tags ...tag.Tag) { t.record(LevelWarn, msg, tags) } // With returns a derived [Logger] that shares the receiver's recording store // and prepends the receiver's tags, followed by tags, to every entry it logs. diff --git a/pkg/logger/zero.go b/pkg/logger/zero.go index 4f3a28e..554cb13 100644 --- a/pkg/logger/zero.go +++ b/pkg/logger/zero.go @@ -24,26 +24,34 @@ func NewZeroLogger(w io.Writer, lvl Level) *ZeroLogger { } } +// Debug implements [Logger]. func (l *ZeroLogger) Debug(msg string, tags ...tag.Tag) { logEvent(l.log.Debug(), msg, tags) } +// Error implements [Logger]. func (l *ZeroLogger) Error(msg string, tags ...tag.Tag) { logEvent(l.log.Error(), msg, tags) } +// Fatal implements [Logger]. It writes the entry, then exits the process +// with status 1. func (l *ZeroLogger) Fatal(msg string, tags ...tag.Tag) { logEvent(l.log.Fatal(), msg, tags) } +// Info implements [Logger]. func (l *ZeroLogger) Info(msg string, tags ...tag.Tag) { logEvent(l.log.Info(), msg, tags) } +// Warn implements [Logger]. func (l *ZeroLogger) Warn(msg string, tags ...tag.Tag) { logEvent(l.log.Warn(), msg, tags) } +// With implements [Logger], returning a child ZeroLogger that includes tags +// on every entry. func (l *ZeroLogger) With(tags ...tag.Tag) Logger { ctx := l.log.With() for _, t := range tags { diff --git a/pkg/testutil/certs.go b/pkg/testutil/certs.go index 4952599..730d713 100644 --- a/pkg/testutil/certs.go +++ b/pkg/testutil/certs.go @@ -46,10 +46,10 @@ func ECDSACert(t *testing.T, tmpl *x509.Certificate) []byte { } // GenerateSelfSignedCert writes a self-signed ECDSA P-256 certificate and its -// matching PKCS#1-style EC private key to a fresh [testing.T.TempDir] and -// returns the paths. The certificate is valid for one hour, advertises CN -// "localhost" with DNSNames=["localhost"], and is suitable for loading via -// [crypto/tls.LoadX509KeyPair] in server-auth scenarios. +// matching SEC 1 ("EC PRIVATE KEY") private key to a fresh +// [testing.T.TempDir] and returns the paths. The certificate is valid for one +// hour, advertises CN "localhost" with DNSNames=["localhost"], and is suitable +// for loading via [crypto/tls.LoadX509KeyPair] in server-auth scenarios. func GenerateSelfSignedCert(t *testing.T) (certFile, keyFile string) { t.Helper() @@ -120,12 +120,10 @@ func GenerateRSACert(t *testing.T) (certFile, keyFile string) { // GenerateMTLSCerts writes a self-signed ECDSA P-256 CA certificate plus an // RSA-2048 leaf certificate signed by that CA (with its matching key) to a -// fresh [testing.T.TempDir] and returns the three paths. The leaf is RSA -// because credential validation checks the leaf's key type against an RSA-only -// cipher suite allowlist; an ECDSA leaf fails that check even though it -// verifies fine against the CA. The leaf advertises CN "localhost" with -// DNSNames=["localhost"]; both certificates are valid for one hour. Use this -// when a test needs the leaf to verify against the CA. +// fresh [testing.T.TempDir] and returns the three paths. The leaf key is +// written as PKCS#1 ("RSA PRIVATE KEY"). The leaf advertises CN "localhost" +// with DNSNames=["localhost"]; both certificates are valid for one hour. Use +// this when a test needs the leaf to verify against the CA. func GenerateMTLSCerts(t *testing.T) (caFile, certFile, keyFile string) { t.Helper() diff --git a/pkg/validation/certs/certs.go b/pkg/validation/certs/certs.go index 2f167d8..c32d008 100644 --- a/pkg/validation/certs/certs.go +++ b/pkg/validation/certs/certs.go @@ -107,6 +107,9 @@ func ValidatePEMKeyFile(path string) error { // ValidatePEM parses all CERTIFICATE blocks from pemData and runs each check // against every parsed certificate, collecting all failures into an // [validation.Errors] value. Returns nil immediately when no checks are provided. +// A check failure that is not a [validation.Error] is recorded with Field +// "unknown". PEM data with no certificates, or a certificate that fails to +// parse, returns a plain error rather than a [validation.Errors]. func ValidatePEM(pemData []byte, checks ...Check) error { // With no checks there is nothing to verify; skip PEM parsing entirely. if len(checks) == 0 { @@ -194,17 +197,17 @@ func IsCA() Check { // 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 +// 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 +// 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, via -// leafChecks, 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. +// 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. func SecureAlgorithm(allowedSuites ...uint16) Check { return func(cert *x509.Certificate) error { selfIssued := bytes.Equal(cert.RawIssuer, cert.RawSubject) diff --git a/pkg/validation/certs/doc.go b/pkg/validation/certs/doc.go index ee28f39..1da6810 100644 --- a/pkg/validation/certs/doc.go +++ b/pkg/validation/certs/doc.go @@ -1,6 +1,7 @@ // Package certs provides reusable [validation.Check] building blocks for // inspecting X.509 certificates and PEM material: expiry, CA basic constraint, -// signature-algorithm and key-type strength, and key size. ValidatePEM and its -// file variants parse PEM data and run a set of checks against every certificate -// they contain, aggregating failures into a [validation.Errors]. +// signature-algorithm and key-type strength, and key size. ValidatePEM and +// ValidatePEMFile parse PEM data and run a set of checks against every +// certificate they contain, aggregating failures into a [validation.Errors]. +// ValidatePEMKeyFile takes no checks; it only confirms a private key is present. package certs diff --git a/pkg/validation/checks.go b/pkg/validation/checks.go index fb77941..acf012d 100644 --- a/pkg/validation/checks.go +++ b/pkg/validation/checks.go @@ -18,7 +18,8 @@ type Check[V any] func(V) error // purely lexical: no DNS lookup, no /etc/services port resolution. Both // "host:port" and ":port" (listen-on-all-interfaces) forms are valid, and the // port must be a decimal integer in [0, 65535]. Port 0 is accepted because it -// is a valid listener form meaning "let the OS pick". +// is a valid listener form meaning "let the OS pick". A host containing "/" is +// rejected. func IsHostPort() Check[string] { return func(s string) error { host, port, err := net.SplitHostPort(s) @@ -41,8 +42,8 @@ func IsHostPort() Check[string] { // GT rejects any value not strictly greater than mark; a value equal to mark // fails. It works for any cmp.Ordered type (numbers, strings), comparing with -// the language's < operator, so for floats a NaN bound or input never -// satisfies the check. +// the language's comparison operators, so for floats a NaN bound or input +// always passes the check. func GT[V cmp.Ordered](mark V) Check[V] { return func(v V) error { if v <= mark { @@ -54,7 +55,8 @@ func GT[V cmp.Ordered](mark V) Check[V] { } // GTE rejects any value less than mark; a value equal to mark passes. It is the -// inclusive counterpart to GT and shares its cmp.Ordered and NaN semantics. +// inclusive counterpart to GT and shares its cmp.Ordered semantics; a NaN bound +// or input always passes. func GTE[V cmp.Ordered](mark V) Check[V] { return func(v V) error { if v < mark { @@ -66,7 +68,8 @@ func GTE[V cmp.Ordered](mark V) Check[V] { } // LT rejects any value not strictly less than mark; a value equal to mark -// fails. It is the mirror of GT and shares its cmp.Ordered and NaN semantics. +// fails. It is the mirror of GT and shares its cmp.Ordered semantics; a NaN +// bound or input always passes. func LT[V cmp.Ordered](mark V) Check[V] { return func(v V) error { if v >= mark { diff --git a/pkg/validation/doc.go b/pkg/validation/doc.go index 7a89160..29eb08b 100644 --- a/pkg/validation/doc.go +++ b/pkg/validation/doc.go @@ -5,8 +5,9 @@ // where Subject identifies the thing being validated (e.g. a hostname or // certificate CN), Field names the attribute that failed, and Message // describes the failure in human-readable form. [Errors] aggregates multiple -// [Error] values into a single error while remaining compatible with -// [errors.Is], [errors.As], and [errors.Join] via its Unwrap method. +// [Error] values into a single error; its Unwrap method lets [errors.Is] and +// [errors.As] see each entry, and it composes with [errors.Join] like any +// other error. // // Higher-level validation is expressed by passing a list of [Rule] values to // [Validate], typically constructed via [Field] and the built-in [Check] diff --git a/pkg/validation/validate.go b/pkg/validation/validate.go index 99038c5..f8541e4 100644 --- a/pkg/validation/validate.go +++ b/pkg/validation/validate.go @@ -2,9 +2,9 @@ package validation import "errors" -// Validate runs each rule, accumulating failures into a single error (which -// will typically be an [Errors] instance). Any [Error] whose Subject is empty is -// stamped with subject. +// Validate runs each rule, accumulating failures into a single error. It +// returns nil when every rule passes and an [Errors] otherwise. Any [Error] +// whose Subject is empty is stamped with subject. func Validate(subject string, rules ...Rule) error { var errs Errors for _, r := range rules {