Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions cmd/proxy/main.go
Original file line number Diff line number Diff line change
@@ -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 (
Expand Down
2 changes: 1 addition & 1 deletion e2e/doc.go
Original file line number Diff line number Diff line change
@@ -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
13 changes: 8 additions & 5 deletions examples/authz/server/authorizer.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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},
}
Expand Down
5 changes: 3 additions & 2 deletions examples/authz/server/claims.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
4 changes: 2 additions & 2 deletions examples/authz/token.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)

Expand Down
6 changes: 3 additions & 3 deletions internal/api/fx.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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)
Expand Down
2 changes: 1 addition & 1 deletion internal/api/kms.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand Down
14 changes: 7 additions & 7 deletions internal/auth/jwks.go
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
5 changes: 3 additions & 2 deletions internal/auth/outbound/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 4 additions & 4 deletions internal/auth/outbound/static_provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 "<scheme> <apiKey>" 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
Expand Down
9 changes: 5 additions & 4 deletions internal/cloud/doc.go
Original file line number Diff line number Diff line change
@@ -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
// "<name>.<account-id>", 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 "<name>.<account-id>",
// 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
3 changes: 2 additions & 1 deletion internal/cloud/namespace.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 2 additions & 1 deletion internal/cloud/translation/translation.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion internal/codecserver/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 4 additions & 4 deletions internal/codecserver/fx.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 3 additions & 5 deletions internal/codecserver/reporter.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
2 changes: 1 addition & 1 deletion internal/config/auth.go
Original file line number Diff line number Diff line change
Expand Up @@ -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(
"",
Expand Down
5 changes: 2 additions & 3 deletions internal/config/cloudapi.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
//
Expand Down Expand Up @@ -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()
}
Expand Down
3 changes: 2 additions & 1 deletion internal/config/codecserver.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
2 changes: 1 addition & 1 deletion internal/config/encryption.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
11 changes: 6 additions & 5 deletions internal/config/extensions.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions internal/config/fx.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 3 additions & 3 deletions internal/config/routing.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Expand Down
14 changes: 11 additions & 3 deletions internal/config/upstream.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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.
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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

Expand Down
4 changes: 2 additions & 2 deletions internal/dataplane/dataplane.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down
7 changes: 3 additions & 4 deletions internal/dataplane/dataplanetest/dataplanetest.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading
Loading