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 {