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
49 changes: 30 additions & 19 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,12 @@ reaches a different upstream with no change to the Worker.
responses, so the upstream only ever sees ciphertext while local Workers keep exchanging cleartext. DEKs are wrapped
by a KMS key (AWS KMS, Azure Key Vault, or GCP KMS), rotate automatically, and can be overridden per Namespace.
Set `encryption.failures` to seal failure messages and stack traces too, as the SDK's `EncodeCommonAttributes` does.

Payloads a Worker's own codec already encrypted, or that an earlier proxy sealed, can be forwarded as they are by
listing their encoding under `encryption.skipEncodings`. Listing an encoding trusts every Worker that sets it: the
proxy cannot tell ciphertext from plaintext labeled that way. Listing `binary/encrypted` also returns responses sealed
under a KEK this proxy doesn't hold, including its own if a key was removed from config, rather than failing the
call, so alert on `vault_ops_total{result="unknown_key"}`.
- **Pluggable key management.** For a backend the proxy has no built-in support for, such as an on-prem HSM or an
internal key service, point it at an extension server you run and it wraps DEKs through that instead. Only key
material is exchanged; payloads never reach it.
Expand Down Expand Up @@ -234,25 +240,30 @@ label reports blank there even though the same series populates it for a request
fixed label is added to every series in the table, and to nothing registered outside the proxy's own collectors, so
the runtime's `go_*` and `process_*` series stay as they are.

| Subsystem | Metric | Type | Labels |
| -------------- | ---------------------------- | --------- | ---------------------------------- |
| `server` | `requests_total` | counter | `method`, `code` |
| `server` | `request_duration_seconds` | histogram | `method` |
| `server` | `panics_total` | counter | `method` |
| `router` | `decisions_total` | counter | `upstream`, `outcome` |
| `router` | `forwarding_errors_total` | counter | `upstream`, `reason` |
| `encryption` | `vault_ops_total` | counter | `operation`, `result`, `namespace` |
| `encryption` | `vault_ops_duration_seconds` | histogram | `operation`, `namespace` |
| `encryption` | `kek_ops_total` | counter | `provider`, `operation`, `result` |
| `encryption` | `kek_ops_duration_seconds` | histogram | `provider`, `operation` |
| `encryption` | `dek_ops_total` | counter | `operation`, `result` |
| `encryption` | `dek_ops_duration_seconds` | histogram | `operation` |
| `encryption` | `dek_rotations_total` | counter | `reason` |
| `encryption` | `dek_cache_hits_total` | counter | none |
| `encryption` | `dek_cache_misses_total` | counter | none |
| `encryption` | `dek_cache_size` | gauge | none |
| `codec_server` | `requests_total` | counter | `route`, `code` |
| `codec_server` | `request_duration_seconds` | histogram | `route` |
`payloads_skipped_total` carries the metadata labels too (blank, as with `vault_ops`, for operations from the codec
server, since an HTTP request carries no gRPC metadata), and `vault_ops_total` reports `result="unknown_key"` for a
payload sealed under a KEK the proxy doesn't hold.

| Subsystem | Metric | Type | Labels |
| -------------- | ---------------------------- | --------- | ------------------------------------ |
| `server` | `requests_total` | counter | `method`, `code` |
| `server` | `request_duration_seconds` | histogram | `method` |
| `server` | `panics_total` | counter | `method` |
| `router` | `decisions_total` | counter | `upstream`, `outcome` |
| `router` | `forwarding_errors_total` | counter | `upstream`, `reason` |
| `encryption` | `vault_ops_total` | counter | `operation`, `result`, `namespace` |
| `encryption` | `vault_ops_duration_seconds` | histogram | `operation`, `namespace` |
| `encryption` | `payloads_skipped_total` | counter | `operation`, `encoding`, `namespace` |
| `encryption` | `kek_ops_total` | counter | `provider`, `operation`, `result` |
| `encryption` | `kek_ops_duration_seconds` | histogram | `provider`, `operation` |
| `encryption` | `dek_ops_total` | counter | `operation`, `result` |
| `encryption` | `dek_ops_duration_seconds` | histogram | `operation` |
| `encryption` | `dek_rotations_total` | counter | `reason` |
| `encryption` | `dek_cache_hits_total` | counter | none |
| `encryption` | `dek_cache_misses_total` | counter | none |
| `encryption` | `dek_cache_size` | gauge | none |
| `codec_server` | `requests_total` | counter | `route`, `code` |
| `codec_server` | `request_duration_seconds` | histogram | `route` |

The `encryption` subsystem only reports once encryption keys are configured, and `codec_server` only reports once the
codec server is enabled. Its `route` label is the matched pattern (`/decode`, and so on), never the request path,
Expand Down
4 changes: 4 additions & 0 deletions dev/config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,10 @@ routing:
encryption:
enabled: true
cacheSize: 200
# Encodings already encrypted before they reach the proxy, forwarded unsealed.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit:
Maybe note in the commend that this is a header?
It looks like a header, but might as well be very explicit.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

They're actually fields (Payload.Metadata["encoding"]) set by PayloadCodecs.

# List binary/encrypted on a proxy chained behind another one.
# skipEncodings:
# - binary/encrypted

default:
uri: testing://-ynIaZzFbAjp9VPgu0Ohk9YeQSLS9ta0m9mtnOnGZqo=
Expand Down
100 changes: 99 additions & 1 deletion e2e/encryption_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -90,15 +90,113 @@ func TestEndToEndPayloadEncryption(t *testing.T) {
require.NotEmpty(t, sealed.GetMetadata()[wireDEK], "sealed payload must carry the wrapped DEK")
}

// TestEndToEndSkipsWorkerEncryptedPayloads drives a payload a worker's own codec
// already encrypted through the full stack. Not copying SkipEncodings into
// CodecOptions in dataplane.New fails the byte-for-byte check.
func TestEndToEndSkipsWorkerEncryptedPayloads(t *testing.T) {
t.Parallel()

up := dataplanetest.NewUpstream(t)

cfg := dataplanetest.Config(up)
cfg.Encryption = config.Encryption{
Enabled: true,
SkipEncodings: []string{"acme/aes-gcm"},
Default: &config.KeyPolicy{URI: testingKeyURI(t), Duration: time.Hour},
}

f := dataplanetest.StartApp(t, cfg)

workerSealed := &common.Payload{
Metadata: map[string][]byte{wireEncoding: []byte("acme/aes-gcm")},
Data: []byte("ciphertext-from-a-worker"),
}

resp, err := f.Client().QueryWorkflow(f.Context(), queryWith(workerSealed), grpc.WaitForReady(true))
require.NoError(t, err)
require.True(t, proto.Equal(workerSealed, resp.GetQueryResult().GetPayloads()[0]))

reqs := up.Requests()
require.Len(t, reqs, 1)
sent := reqs[0].(*workflowservice.QueryWorkflowRequest).GetQuery().GetQueryArgs().GetPayloads()
require.Len(t, sent, 1)
require.True(t, proto.Equal(workerSealed, sent[0]), "upstream must receive the worker's payload unchanged")
}

// TestEndToEndChainedProxiesSealOnce runs client -> hop A -> hop B -> upstream,
// where only A holds A's key and B lists binary/encrypted. A round trip alone
// can't prove B skipped, since B would decode its own extra layer, so B's skip
// metric is the evidence. Removing the unknown-key pass-through in
// Encryptor.Decode fails the QueryWorkflow call.
func TestEndToEndChainedProxiesSealOnce(t *testing.T) {
t.Parallel()

up := dataplanetest.NewUpstream(t)

cfgB := dataplanetest.Config(up)
cfgB.Encryption = config.Encryption{
Enabled: true,
SkipEncodings: []string{wireEncryptedMarker},
Default: &config.KeyPolicy{URI: testingKeyURIOf(t, 0x0b), Duration: time.Hour},
}
hopB := dataplanetest.StartApp(t, cfgB)

cfgA := dataplanetest.Config(up)
cfgA.Upstreams[0].Listen = config.ListenConfig{HostPort: hopB.Addr(), Insecure: true}
cfgA.Encryption = config.Encryption{
Enabled: true,
Default: &config.KeyPolicy{URI: testingKeyURIOf(t, 0x0a), Duration: time.Hour},
}
hopA := dataplanetest.StartApp(t, cfgA)

secret := &common.Payload{
Metadata: map[string][]byte{wireEncoding: []byte("json/plain")},
Data: []byte(`"the-answer-is-42"`),
}

resp, err := hopA.Client().QueryWorkflow(hopA.Context(), queryWith(secret), grpc.WaitForReady(true))
require.NoError(t, err)
require.True(t, proto.Equal(secret, resp.GetQueryResult().GetPayloads()[0]))

reqs := up.Requests()
require.Len(t, reqs, 1)
payloads := reqs[0].(*workflowservice.QueryWorkflowRequest).GetQuery().GetQueryArgs().GetPayloads()
require.Len(t, payloads, 1)
require.Equal(t, wireEncryptedMarker, string(payloads[0].GetMetadata()[wireEncoding]))

requireLabel(t, hopB, "test_encryption_payloads_skipped_total", "operation", "encrypt")
requireLabel(t, hopB, "test_encryption_payloads_skipped_total", "operation", "decrypt")
}

// testingKeyURI builds a local testing:// key URI with a fixed 32-byte key. The
// kms module rewrites testing:// to gocloud's base64key:// local keeper, so no
// cloud KMS is needed.
func testingKeyURI(t *testing.T) url.URL {
t.Helper()

key := base64.StdEncoding.EncodeToString(bytes.Repeat([]byte{0x2a}, 32))
return testingKeyURIOf(t, 0x2a)
}

// testingKeyURIOf is testingKeyURI with the key filled with b, so two proxies in
// one test can hold different keys.
func testingKeyURIOf(t *testing.T, b byte) url.URL {
t.Helper()

key := base64.StdEncoding.EncodeToString(bytes.Repeat([]byte{b}, 32))
u, err := url.Parse("testing://" + key)
require.NoError(t, err)

return *u
}

// queryWith is a QueryWorkflow request carrying p as its only query argument.
func queryWith(p *common.Payload) *workflowservice.QueryWorkflowRequest {
return &workflowservice.QueryWorkflowRequest{
Namespace: "ns1",
Execution: &common.WorkflowExecution{WorkflowId: "wf-1"},
Query: &query.WorkflowQuery{
QueryType: "state",
QueryArgs: &common.Payloads{Payloads: []*common.Payload{p}},
},
}
}
31 changes: 25 additions & 6 deletions internal/config/encryption.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,13 +29,15 @@ type (
// (local) namespace names, matching the namespace the vault seals under at
// request time. Failures also seals the message and stack trace of outbound
// failures, the way the Temporal SDK's EncodeCommonAttributes does; it requires
// Enabled.
// Enabled. SkipEncodings lists payload encodings already encrypted before they
// reach the proxy, which are forwarded unsealed.
Encryption struct {
Enabled bool `yaml:"enabled"`
Failures bool `yaml:"failures"`
CacheSize *int `yaml:"cacheSize"`
Default *KeyPolicy `yaml:"default"`
Overrides map[string]KeyPolicy `yaml:"overrides"`
Enabled bool `yaml:"enabled"`
Failures bool `yaml:"failures"`
CacheSize *int `yaml:"cacheSize"`
Default *KeyPolicy `yaml:"default"`
Overrides map[string]KeyPolicy `yaml:"overrides"`
SkipEncodings []string `yaml:"skipEncodings"`
}

// KeyPolicy describes the KMS key backing a DEK and its rotation schedule.
Expand Down Expand Up @@ -95,6 +97,11 @@ func (e *Encryption) Validate() error {
)
}

rules = append(rules,
validation.Field("skipEncodings", e.SkipEncodings, validation.Unique[string]()),
validation.Children("skipEncodings", e.SkipEncodings, nonBlankEncoding()),
)

return validation.Validate("", rules...)
}

Expand Down Expand Up @@ -195,3 +202,15 @@ func validKeyURIRef() validation.Check[*url.URL] {
return nil
}
}

// nonBlankEncoding rejects a blank skipEncodings entry, which would match every
// payload that carries no encoding at all.
func nonBlankEncoding() validation.Check[*string] {
return func(s *string) error {
if *s == "" {
return errors.New("must not be blank")
}

return nil
}
}
29 changes: 29 additions & 0 deletions internal/config/encryption_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -170,6 +170,35 @@ func TestEncryptionValidate(t *testing.T) {
},
wantErr: "overrides",
},
{
name: "skip encodings",
cfg: config.Encryption{
Enabled: true,
Default: &valid,
SkipEncodings: []string{"binary/encrypted", "acme/aes-gcm"},
},
},
{
// A blank entry would match every payload with no encoding. Dropping the
// Children rule fails this row.
name: "blank skip encoding",
cfg: config.Encryption{
Enabled: true,
Default: &valid,
SkipEncodings: []string{"acme/aes-gcm", ""},
},
wantErr: "skipEncodings[1]",
},
{
// Dropping the Unique check fails this row.
name: "duplicate skip encoding",
cfg: config.Encryption{
Enabled: true,
Default: &valid,
SkipEncodings: []string{"acme/aes-gcm", "acme/aes-gcm"},
},
wantErr: "skipEncodings",
},
}

for _, tt := range tests {
Expand Down
11 changes: 10 additions & 1 deletion internal/dataplane/dataplane.go
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,11 @@ func New(ctx context.Context, cfg *config.Config, opts ...Option) (*Dataplane, e
// Every upstream applies the same chain: the vault and the encryption switch
// are global, so nothing here varies per upstream. Building it once is also
// what lets the codec server apply the identical chain.
codecOpts := proxy.CodecOptions{Encrypt: cfg.Encryption.Enabled, EncodeFailures: cfg.Encryption.Failures}
codecOpts := proxy.CodecOptions{
Encrypt: cfg.Encryption.Enabled,
EncodeFailures: cfg.Encryption.Failures,
SkipEncodings: cfg.Encryption.SkipEncodings,
}

// Only assign the vault once it is known to be there. o.vault is a concrete
// pointer and the field is an interface, so assigning unconditionally would
Expand Down Expand Up @@ -144,6 +148,11 @@ func New(ctx context.Context, cfg *config.Config, opts ...Option) (*Dataplane, e
)
}

// With no keys there is no encryption codec, so the list has nothing to skip.
if len(cfg.Encryption.SkipEncodings) > 0 && o.vault == nil {
o.logger.Warn("encryption.skipEncodings is set but no encryption keys are configured, so it has no effect")
}

Comment on lines +151 to +155

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

im fine with this but there is an argument for just failing invalid/non-effect configs

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. So far, I've been failing cases where behaviour would be invalid (e.g., enabling encryption without a default key policy) and warning about things that are ineffectual or don't break an invariant (e.g., apiTranslations defined with no Cloud upstream).

I absolutely don't feel super strongly about this, so if you'd prefer an error case here, say the word.

dp := &Dataplane{
ctx: ctx,
hostPort: cfg.Listen.HostPort,
Expand Down
38 changes: 38 additions & 0 deletions internal/dataplane/dataplane_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -222,6 +222,39 @@ func TestNewTwiceOverOneMetricsFactoryDoesNotPanic(t *testing.T) {
}
}

// TestNewWarnsWhenSkipEncodingsHasNoKeys covers a list with nothing to apply it
// to. Removing the warning fails the first case; dropping the length check fails
// the second.
func TestNewWarnsWhenSkipEncodingsHasNoKeys(t *testing.T) {
t.Parallel()

tests := []struct {
name string
skip []string
wantWarn bool
}{
{name: "list without keys", skip: []string{"acme/aes-gcm"}, wantWarn: true},
{name: "no list"},
}

for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
t.Parallel()

cfg := testConfig()
cfg.Encryption.SkipEncodings = tc.skip

log := logger.NewTestLogger()
deps := newTestDeps(t, cfg)
deps.logger = log

_, err := dataplane.New(deps.ctx, cfg, deps.opts()...)
require.NoError(t, err)
require.Equal(t, tc.wantWarn, log.Contains(skipEncodingsUnusedWarning))
})
}
}

// opts returns d as the options [dataplane.New] takes, less any named in omit,
// so a caller can prove New reports one as missing. Names are the ones New
// reports.
Expand Down Expand Up @@ -303,6 +336,11 @@ func testingKeyURL(t *testing.T) url.URL {
const cloudAPIUnusedWarning = "apiTranslations is configured but no upstream is Temporal Cloud, so no method " +
"will be translated"

// skipEncodingsUnusedWarning is the message New logs for a skip list with no
// encryption codec to act on.
const skipEncodingsUnusedWarning = "encryption.skipEncodings is set but no encryption keys are configured, so it " +
"has no effect"

// stagingCloudAPI is an override that says something, which is what an inert
// block has to be to be worth warning about: one that says nothing is
// indistinguishable from no block at all, and describes the defaults anyway.
Expand Down
21 changes: 21 additions & 0 deletions internal/proxy/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,12 @@ type (
// are sealed with every other payload. It requires Encrypt.
EncodeFailures bool

// SkipEncodings lists payload encodings already encrypted before they
// reach the proxy. Outbound payloads under one are forwarded unsealed, and
// listing codec.EncryptionEncoding also returns inbound payloads sealed
// under a KEK the vault doesn't hold, rather than failing the call.
SkipEncodings []string

// Reporter records the duration and result of each vault operation. It is
// required whenever a Vault is set.
Reporter *Reporter
Expand Down Expand Up @@ -83,6 +89,21 @@ func NewCodecs(opts CodecOptions) (*Codecs, error) {
if opts.Encrypt {
c.outbound = append(c.outbound, enc)
}

// The set is built once here; only the observer is bound per request.
if len(opts.SkipEncodings) > 0 {
encodings := codec.WithSkipEncodings(opts.SkipEncodings...)
skip := func(ctx context.Context, ns string) codec.Option {
return codec.WithEncryptorOptions(encodings, codec.WithSkipObserver(func(op, encoding string) {
opts.Reporter.PayloadSkipped(ctx, op, encoding, ns)
}))
}

c.inbound = append(c.inbound, skip)
if opts.Encrypt {
c.outbound = append(c.outbound, skip)
}
}
}

return c, nil
Expand Down
Loading
Loading