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
2 changes: 0 additions & 2 deletions cmd/proxy/serve.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,6 @@ const shutdownTimeout = 30 * time.Second
// fxLogger swallows fx's event stream except for Started/Stopped failures,
// which it forwards to the app logger. Configuration errors are caught
// separately via app.Err() so they are visible at startup.
//
// The net effect is fx.NopLogger plus visible lifecycle hook failures.
type fxLogger struct {
log logger.Logger
}
Expand Down
45 changes: 12 additions & 33 deletions internal/api/auth.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,16 +18,10 @@ import (
)

// Auth authenticates an inbound stream by delegating the decision to an
// extension server implementing api.auth.v1.AuthService. It is the escape hatch
// for identity systems the built-in authenticators do not cover: the proxy
// asks, the operator's server decides.
//
// The headers name the metadata carrying the caller's credentials. They are
// declared rather than discovered because a verdict reports only admit-or-deny
// and says nothing about which headers mattered, and the proxy needs to know two
// things: which values to lift into the request, and which to report as
// [Auth.SecureHeaders] so they are stripped from the stream before it reaches an
// upstream, where a caller credential would collide with the proxy's own.
// extension server implementing api.auth.v1.AuthService. Its headers name the
// metadata carrying the caller's credentials: their values are lifted into the
// request, and they are reported as [Auth.SecureHeaders] so they are stripped
// from the stream before it reaches an upstream.
type Auth struct {
client auth.AuthServiceClient
headers []string
Expand All @@ -50,31 +44,16 @@ func NewAuth(cc grpc.ClientConnInterface, secureHeaders []string) *Auth {
}

// Authenticate asks the extension server whether the caller may proceed. Only an
// explicit DECISION_ALLOW admits the stream: an error, a denial, and an answer
// carrying no verdict all deny it, so a server that is down, misconfigured, or
// newer than this build fails the request closed rather than opening the gateway
// to everyone for as long as it is that way.
//
// A denial reaches the caller as an [rpc.Reject], whose message is generic while
// the provider's reason becomes the server-side detail: a provider writes that
// reason for whoever operates it, not for the caller it just turned away. An
// error keeps the provider's status code, since it tells a worker whether to fix
// its credential or retry; a decision carries no code, so the proxy supplies one.
// explicit DECISION_ALLOW admits the stream; an error, a denial, or an answer
// carrying no verdict denies it, so the request fails closed. Every rejection is
// an [rpc.Reject] with a generic message and the provider's reason as the
// server-side detail. An error keeps the provider's status code; a denial is
// PermissionDenied and a missing verdict is Internal.
//
// The declared credential headers are lifted into the request and withheld from
// the forwarded metadata, so each credential reaches the server in exactly one
// place. That separation is what lets the proxy hold a credential of its own to
// this server: metadata carries the proxy's, the request carries the caller's,
// and neither has to be told apart from the other on a shared header. It also
// puts the caller's credential out of reach of the interceptor [outbound.DialOptions]
// installs, which deletes the proxy's credential header from forwarded metadata
// and cannot tell that on this one call that header is the subject of the request
// rather than incidental cargo.
//
// The caller's remaining metadata is forwarded so the server can weigh context
// such as the method being invoked. gRPC drops reserved keys (":authority",
// "user-agent", "content-type", "grpc-*") when writing the request, so a caller
// cannot reach the extension server's transport this way.
// the forwarded metadata. The caller's remaining metadata is forwarded, except
// the reserved keys gRPC drops (":authority", "user-agent", "content-type",
// "grpc-*"), so a caller cannot reach the extension server's transport this way.
func (a *Auth) Authenticate(ctx context.Context, target meta.Target, md metadata.MD) error {
req := &auth.AuthRequest{
Target: &auth.Target{FullName: target.FullName, Namespace: target.Namespace},
Expand Down
29 changes: 9 additions & 20 deletions internal/api/fx.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,14 +13,10 @@ import (
)

// extensionKeyPrefix namespaces extension-server entries in the shared
// connection pool.
//
// A static resolver uses its dial address as the pool key, and Pool.ConnOrCreate
// returns the existing connection for a key while ignoring the options passed
// with it. Upstream hostPorts are unique among themselves, but nothing stops an
// extension server from sitting on the same host:port as an upstream, so without
// a distinct key the two would collapse onto whichever was dialed first and
// silently inherit its TLS settings and credentials.
// connection pool. A static resolver keys the pool by dial address and
// Pool.ConnOrCreate ignores the options for an existing key, so without the
// prefix an extension server on an upstream's host:port would share that
// connection and silently inherit its TLS settings and credentials.
const extensionKeyPrefix = "extension:"

// Module provides the pooled connection for every configured extension server,
Expand Down Expand Up @@ -86,24 +82,17 @@ type (

// Connections maps an extension server name to a connection to that server.
// It carries no lifecycle: closing a connection is the owner's
// responsibility, not the caller's.
//
// Callers get connections rather than finished clients because the two do not
// correspond one-to-one: several keys may live on one extension server, so a
// caller builds one [KMS] per key over the shared connection.
// responsibility, not the caller's. Several keys may live on one extension
// server, so a caller builds one [KMS] per key over the shared connection.
Connections map[string]grpc.ClientConnInterface
)

// extensionConn builds the pooled connection for a single extension server.
// Config rejects a templated hostPort, so the target is always static: the
// resolver is fixed and [connect.NewConn] creates the connection here, which the
// module then opens on start.
//
// This is deliberately narrower than the equivalent upstream path in
// internal/proxy. There is no namespace translation, because an extension
// server is not a Temporal service and has no namespaces to rewrite, and no
// payload encryption interceptor, because an extension server is the thing that
// wraps DEKs; sealing its traffic with the vault it backs would be circular.
// module then opens on start. Unlike an upstream connection it installs no
// namespace translation or payload encryption interceptor; an extension server
// wraps DEKs, so sealing its traffic with the vault it backs would be circular.
func extensionConn(pool *connect.Pool, s *config.ExtensionServer) (*connect.Conn, error) {
var opts []grpc.DialOption

Expand Down
5 changes: 2 additions & 3 deletions internal/auth/authenticator.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,8 @@ type (
// status error to reject it. SecureHeaders reports the metadata headers the
// authenticator consumes, so the proxy can strip the caller's credentials
// before forwarding upstream; it returns nil when the authenticator consumes
// no header. An authenticator may name more than one because it need not own
// the header it reads: an external one is told which headers its server
// consumes.
// no header and may name more than one (an external authenticator reports
// the headers its server consumes).
//
// The target is what the gateway resolved for this stream. Its Namespace is
// empty for a request that named none, so an implementation weighing it must
Expand Down
28 changes: 6 additions & 22 deletions internal/auth/jwks.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,20 +57,9 @@ type JWKSAuthenticator struct {
}

// NewJWKSAuthenticator builds a JWKSAuthenticator that resolves signing keys
// from the JWKS at rawURL.
//
// keyfunc.NewDefaultOverrideCtx performs its first key fetch synchronously
// (see github.com/MicahParks/jwkset's NewStorageFromHTTP), so calling it
// directly here would block construction for up to its HTTP timeout if the
// IdP is unreachable. To keep construction non-blocking, the fetch is kicked
// off in a background goroutine (via deferredKeyfunc); until it completes,
// Authenticate reports codes.Unavailable (fail closed, retryable) instead of
// blocking startup.
//
// rawURL is validated synchronously before the goroutine is started, so a
// malformed configuration (bad scheme/host, not a transient IdP outage) fails
// construction immediately instead of surfacing as a permanent
// codes.Unavailable for every future request.
// from the JWKS at rawURL. It returns an error if rawURL is malformed but does
// not block on the IdP: the first key fetch runs in the background, and until
// it completes Authenticate fails closed with codes.Unavailable.
func NewJWKSAuthenticator(rawURL string, audiences []string, issuer, header, scheme string) (*JWKSAuthenticator, error) {
if err := validateJWKSURL(rawURL); err != nil {
return nil, err
Expand Down Expand Up @@ -233,14 +222,9 @@ func deferredKeyfunc(load func() (jwt.Keyfunc, error)) (jwt.Keyfunc, *atomic.Poi
// availability problem (errKeysUnavailable -> codes.Unavailable), not a bad
// token. Any other resolver error is treated as an availability problem.
//
// keysPresent reports whether the keyset currently holds any keys. resolve is
// the underlying key resolver (in production, keyfunc.Keyfunc's method value).
//
// When keysPresent cannot positively confirm keys exist (e.g. its underlying
// read errors), it should return false; the failure then maps to
// codes.Unavailable (retryable) rather than masking as a bad token - the
// conservative fail-closed-toward-retryable choice for an infra-level read
// error.
// keysPresent reports whether the keyset currently holds any keys; it must
// return false when it cannot confirm keys exist (e.g. its read errors), so
// the failure maps to codes.Unavailable rather than a bad token.
func wrapKeyfunc(resolve jwt.Keyfunc, keysPresent func() bool) jwt.Keyfunc {
return func(t *jwt.Token) (any, error) {
key, err := resolve(t)
Expand Down
12 changes: 5 additions & 7 deletions internal/cloud/endpoint.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,13 +13,11 @@ const endpointSuffix = ".tmprl.cloud"
// same address the Cloud SDK dials by default.
const APIHostPort = "saas-api" + endpointSuffix + ":443"

// IsEndpoint reports whether hostPort addresses Temporal Cloud. The port is
// optional, and a template action in place of the host is tolerated, since a
// templated address is rendered per request but keeps its domain.
//
// This recognizes per-namespace and regional endpoints. Private-link endpoints
// use per-VPC hostnames that carry no Cloud domain, so those have to be declared
// rather than detected.
// IsEndpoint reports whether hostPort addresses Temporal Cloud through a
// per-namespace or regional endpoint. The port is optional, and a template
// action in place of the host is tolerated. Private-link endpoints use per-VPC
// hostnames with no Cloud domain, so those have to be declared rather than
// detected.
func IsEndpoint(hostPort string) bool {
host := hostPort
if h, _, err := net.SplitHostPort(host); err == nil {
Expand Down
2 changes: 1 addition & 1 deletion internal/cloud/namespace.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ func ValidateAccountID(id string) error {

// ValidateNamespace checks that ns is a well-formed Temporal Cloud namespace
// identifier, meaning "<name>.<account-id>". Every broken rule is reported, not
// just the first, so a caller can show the whole story at once.
// just the first.
//
// See: https://docs.temporal.io/cloud/namespaces for details
func ValidateNamespace(ns string) error {
Expand Down
18 changes: 3 additions & 15 deletions internal/cloud/translation/interceptor.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,9 @@ type options struct {
}

// Via sends a translated call over cc instead of continuing down the chain to
// the connection the interceptor is installed on.
//
// The upstream method belongs to a different service, which the connection that
// received the call does not serve: the caller asked a Temporal Service for
// ListNamespaces, and only Temporal Cloud's control plane can answer it. Where
// that service lives is as fixed as the conversions themselves, so the
// translation carries the connection rather than the request being routed to it,
// and a request reaching any upstream is answered the same way.
// the connection the interceptor is installed on, for an upstream method that
// connection does not serve, such as one only Temporal Cloud's control plane
// answers.
func Via(cc grpc.ClientConnInterface) Option {
return func(o *options) { o.via = cc }
}
Expand All @@ -32,13 +27,6 @@ func Via(cc grpc.ClientConnInterface) Option {
// connection, last, so translation is the innermost interceptor: every other
// interceptor on the chain then sees the method and message types the caller
// asked for rather than the substitute sent upstream.
//
// The result is a slice though it holds a single option today. A [Translation]
// substitutes a unary method, so a unary interceptor is all there is to install;
// translating a streaming method would add a stream interceptor beside it, the
// way the namespace and Cloud-namespace helpers in internal/proxy already pair
// the two. Keeping the slice means that arrives without changing this signature
// or the call sites, which already spread the result.
func DialOptions(r *Registry, opts ...Option) []grpc.DialOption {
return []grpc.DialOption{grpc.WithChainUnaryInterceptor(unaryClientInterceptor(r, opts...))}
}
Expand Down
34 changes: 11 additions & 23 deletions internal/cloud/translation/namespaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,13 +40,10 @@ const (
)

// listNamespaces translates WorkflowService.ListNamespaces onto
// CloudService.GetNamespaces.
//
// The Cloud API version header is required - without it GetNamespaces fails with
// InvalidArgument - and is pinned to the version the SDK this package compiles
// against defaults to, since that is the same module the message types and their
// versioned fields come from. Bumping go.temporal.io/cloud-sdk moves the
// conversions and the version they were written against together.
// CloudService.GetNamespaces. GetNamespaces fails with InvalidArgument without
// the Cloud API version header, so it is pinned to the default of the
// go.temporal.io/cloud-sdk the message types come from; bumping that module
// moves the conversions and the version together.
func listNamespaces() *Translation {
return Adapt(listNamespacesMethod, getNamespacesMethod, listNamespacesRequest, listNamespacesResponse).
WithHeader(cloudclient.TemporalCloudAPIVersionHeader(), cloudclient.DefaultAPIVersion())
Expand Down Expand Up @@ -97,23 +94,14 @@ func listNamespacesResponse(
// response ListNamespaces returns per namespace, with state already mapped by
// the caller so the deleted filter and this conversion agree on it.
//
// Only fields Cloud actually reports carry a value. Cloud has no namespace UUID,
// description, or owner email to give, and it describes replication as regional
// replicas rather than as the clusters ReplicationConfig names, so
// IsGlobalNamespace is derived from how many replicas there are while
// ReplicationConfig is left empty rather than filled with region ids a client
// would read as cluster names. FailoverVersion and FailoverHistory have no
// Cloud equivalent at all.
// Only fields Cloud reports carry a value. Cloud has no namespace UUID,
// description, owner email, or failover data, and its regional replicas are not
// the clusters ReplicationConfig names, so IsGlobalNamespace is derived from the
// replica count and ReplicationConfig is left empty.
//
// Empty is not the same as absent, though, and the difference is load-bearing:
// every sub-message a Temporal Service would populate is allocated here even
// when there is nothing to put in it. A frontend builds NamespaceInfo, Config
// and ReplicationConfig unconditionally on every path (the server funnels them
// all through namespaceHandler.createResponse), so clients are written against a
// reply where they are always present - the temporal CLI reads
// resp.ReplicationConfig.ActiveClusterName with no nil check and dies on a nil
// one. An empty message says "Cloud did not report this" just as well as an
// absent one, without breaking a client that has never had to handle absence.
// Every sub-message a Temporal Service populates is allocated even when empty:
// clients assume they are present, and the temporal CLI reads
// resp.ReplicationConfig.ActiveClusterName with no nil check.
func describeNamespace(
ns *cloudnamespace.Namespace,
state enumspb.NamespaceState,
Expand Down
15 changes: 5 additions & 10 deletions internal/cloud/translation/translation.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,16 +114,11 @@ func Answer[Req, Resp proto.Message](from string, answer func(Req, Resp) error)
}
}

// WithHeader stamps key: value on the substituted call and returns t, so a
// mapping can declare the dialect the upstream method needs alongside the
// conversions themselves. It replaces any value the caller sent rather than
// adding to it: the caller did not ask for this upstream method and cannot know
// what its API expects, so its own header is not intent worth preserving. A
// caller invoking that API directly is forwarded untranslated and keeps its
// header.
//
// Headers travel only on a call this translation substituted; a method the
// registry does not translate is untouched.
// WithHeader stamps key: value on the substituted call and returns t. It
// replaces any value the caller sent rather than adding to it. Headers travel
// only on a call this translation substituted: a method the registry does not
// translate is untouched, and a caller invoking the upstream API directly is
// forwarded untranslated and keeps its header.
func (t *Translation) WithHeader(key, value string) *Translation {
if t.headers == nil {
t.headers = make(map[string]string, 1)
Expand Down
4 changes: 2 additions & 2 deletions internal/codecserver/doc.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
// # Why It Exists
//
// A worker or client connecting through the gateway never sees ciphertext: the
// per-upstream proxy opens inbound payloads whether or not sealing is enabled.
// upstream's forwarder opens inbound payloads whether or not sealing is enabled.
// Two callers do not have that path. The Temporal Cloud UI reaches Cloud
// directly and calls a codec endpoint from the operator's browser, and the
// Temporal CLI run with --codec-endpoint reaches a Temporal Service directly.
Expand Down Expand Up @@ -33,7 +33,7 @@
//
// The handler applies a [Codecs], in practice
// [github.com/temporalio/temporal-proxy/internal/proxy.Codecs], which is the
// same value the per-upstream proxies install as a gRPC client interceptor. A
// same value each upstream's forwarder installs as a gRPC client interceptor. A
// payload therefore transforms identically whichever path it travelled,
// because the chain has one construction site rather than two free to drift
// apart.
Expand Down
Loading
Loading