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
5 changes: 3 additions & 2 deletions cmd/e2a/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,8 @@ func main() {
// migrations and exits without starting the server.
var spFlags sendingProtectionFlags
flag.BoolVar(&spFlags.inspect, "sending-protection-policy", false, "print the stored runtime policy state and the hash this config's policy would activate, then exit")
flag.BoolVar(&spFlags.activate, "activate-sending-protection-policy", false, "CAS-activate this config's sending-protection policy (requires -expected-generation, -expected-policy-sha256, -reason), then exit")
flag.BoolVar(&spFlags.activate, "activate-sending-protection-policy", false, "CAS-activate the selected sending-protection policy (file, or config if absent; requires -expected-generation, -expected-policy-sha256, -reason), then exit")
flag.StringVar(&spFlags.policyFile, "sending-protection-policy-file", "", "reviewed runtime-policy JSON file for inspection or activation (maximum 64 KiB; no config/env overrides)")
flag.Int64Var(&spFlags.expectedGeneration, "expected-generation", -1, "policy generation the operator reviewed (activation CAS)")
flag.StringVar(&spFlags.expectedPolicySHA, "expected-policy-sha256", "", "reviewed canonical policy hash (activation CAS)")
flag.BoolVar(&spFlags.grandfather, "grandfather-current-sending-domains", false, "one-shot: exempt currently sending-verified domains in the same transaction as the activation")
Expand Down Expand Up @@ -1000,7 +1001,7 @@ func main() {
// /v1 contract, so it lives on this mux and never enters api/openapi.yaml.
// The verdict is refreshed on a background ticker, so this handler never
// competes for a pooled connection — see readyz.go for why that matters.
readiness := newReadinessMonitor(pool, &draining)
readiness := newReadinessMonitor(pool, &draining, outboundSending.module)
defer readiness.Stop()
router.HandleFunc("/readyz", readiness.handler()).Methods(http.MethodGet)
// /selftest — deep dependency diagnostics (health+json), auth-gated by the
Expand Down
24 changes: 18 additions & 6 deletions cmd/e2a/readyz.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,14 @@ import (
"github.com/jackc/pgx/v5/pgxpool"

"github.com/tokencanopy/e2a/internal/identity"
"github.com/tokencanopy/e2a/internal/sendingpolicy"
"github.com/tokencanopy/e2a/migrations"
)

// Readiness is evaluated on a background ticker rather than inside the request,
// and a failing probe only flips readiness once it has been failing for
// readinessFailureGrace.
// readinessFailureGrace. An unreadable or invalid database-source policy or
// selected recipient registry entry fails immediately after its probe.
//
// The request-path version of this check was an availability hazard. It called
// pool.Ping with a 2s deadline, so when the pool was saturated the check queued
Expand All @@ -46,6 +48,8 @@ const (
poolProbeTimeout = 250 * time.Millisecond
)

const readinessPolicyUnavailable = "sending protection policy unavailable"

type readinessState struct {
ready bool
reason string
Expand All @@ -57,6 +61,7 @@ type readinessMonitor struct {
pool *pgxpool.Pool
draining *atomic.Bool
latest string
policy *sendingpolicy.Module

interval time.Duration
grace time.Duration
Expand Down Expand Up @@ -85,14 +90,15 @@ type readinessMonitor struct {
// Before that first success the instance is not ready, which preserves the
// guarantee this check exists for: a freshly deployed instance whose migrations
// did not apply never joins the rotation.
func newReadinessMonitor(pool *pgxpool.Pool, draining *atomic.Bool) *readinessMonitor {
return newReadinessMonitorWithConfig(pool, draining,
func newReadinessMonitor(pool *pgxpool.Pool, draining *atomic.Bool, policy *sendingpolicy.Module) *readinessMonitor {
return newReadinessMonitorWithConfig(pool, draining, policy,
readinessProbeInterval, readinessFailureGrace, readinessProbeTimeout)
}

func newReadinessMonitorWithConfig(pool *pgxpool.Pool, draining *atomic.Bool, interval, grace, timeout time.Duration) *readinessMonitor {
func newReadinessMonitorWithConfig(pool *pgxpool.Pool, draining *atomic.Bool, policy *sendingpolicy.Module, interval, grace, timeout time.Duration) *readinessMonitor {
m := &readinessMonitor{
pool: pool,
policy: policy,
draining: draining,
latest: latestMigration(),
interval: interval,
Expand Down Expand Up @@ -130,8 +136,9 @@ func (m *readinessMonitor) evaluate() {
// for transient pressure is what turned a slow database into a total
// outage. Before the first success lastOK is zero, so an instance that has
// never been ready flips immediately rather than being granted a grace
// period it has not earned.
if m.lastOK.IsZero() || now.Sub(m.lastOK) > m.grace {
// period it has not earned. Policy failures are not connectivity pressure:
// retaining a healthy verdict could admit a slot with the wrong controls.
if reason == readinessPolicyUnavailable || m.lastOK.IsZero() || now.Sub(m.lastOK) > m.grace {
m.state.Store(&readinessState{ready: false, reason: reason})
}
}
Expand Down Expand Up @@ -207,6 +214,11 @@ func (m *readinessMonitor) probe() (string, error) {
if !applied {
return "migrations not applied", errNotApplied
}
if m.policy != nil {
if err := m.policy.CheckReadiness(ctx, conn); err != nil {
return readinessPolicyUnavailable, err
}
}
return "", nil
}

Expand Down
14 changes: 7 additions & 7 deletions cmd/e2a/readyz_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,7 @@ func TestLatestMigrationAppliedRecognizesLegacyAlias(t *testing.T) {
func TestReadyzHandler_Ready(t *testing.T) {
pool := migratedTestDB(t)
rec := httptest.NewRecorder()
newReadinessMonitor(pool, nil).handler()(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
newReadinessMonitor(pool, nil, nil).handler()(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil))

if rec.Code != http.StatusOK {
t.Fatalf("status = %d, want 200; body=%s", rec.Code, rec.Body.String())
Expand All @@ -114,7 +114,7 @@ func TestReadyzHandler_DBUnreachable(t *testing.T) {
pool.Close()

rec := httptest.NewRecorder()
newReadinessMonitor(pool, nil).handler()(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
newReadinessMonitor(pool, nil, nil).handler()(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want 503", rec.Code)
}
Expand All @@ -140,7 +140,7 @@ func TestReadyzHandler_Draining(t *testing.T) {
var drain atomic.Bool
drain.Store(true)
rec := httptest.NewRecorder()
newReadinessMonitor(pool, &drain).handler()(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
newReadinessMonitor(pool, &drain, nil).handler()(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil))
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("status = %d, want 503", rec.Code)
}
Expand All @@ -160,7 +160,7 @@ func TestReadyzHandler_Draining(t *testing.T) {
// DOWN and answered every request — including ones that never touch the DB —
// with its own 503.
func TestReadinessToleratesTransientFailureWithinGrace(t *testing.T) {
m := newReadinessMonitorWithConfig(migratedTestDB(t), nil, time.Hour, time.Hour, 2*time.Second)
m := newReadinessMonitorWithConfig(migratedTestDB(t), nil, nil, time.Hour, time.Hour, 2*time.Second)
defer m.Stop()

rec := httptest.NewRecorder()
Expand All @@ -183,7 +183,7 @@ func TestReadinessToleratesTransientFailureWithinGrace(t *testing.T) {
// A failure that outlives the grace window must still flip readiness — the
// tolerance above must not become "never report a dead database".
func TestReadinessFlipsAfterGraceExpires(t *testing.T) {
m := newReadinessMonitorWithConfig(migratedTestDB(t), nil, time.Hour, 0, 2*time.Second)
m := newReadinessMonitorWithConfig(migratedTestDB(t), nil, nil, time.Hour, 0, 2*time.Second)
defer m.Stop()

m.probeFn = func() (string, error) { return "database unreachable", errNotApplied }
Expand All @@ -204,7 +204,7 @@ func TestReadinessFlipsAfterGraceExpires(t *testing.T) {
// held open and the probe must still succeed.
func TestReadinessProbeSurvivesPoolExhaustion(t *testing.T) {
pool := migratedTestDB(t)
m := newReadinessMonitorWithConfig(pool, nil, time.Hour, 0, 3*time.Second)
m := newReadinessMonitorWithConfig(pool, nil, nil, time.Hour, 0, 3*time.Second)
defer m.Stop()

ctx := context.Background()
Expand Down Expand Up @@ -246,7 +246,7 @@ func TestReadinessHandlerDoesNotTouchPool(t *testing.T) {
}
dead.Close()

m := newReadinessMonitorWithConfig(dead, nil, time.Hour, time.Hour, 2*time.Second)
m := newReadinessMonitorWithConfig(dead, nil, nil, time.Hour, time.Hour, 2*time.Second)
defer m.Stop()

start := time.Now()
Expand Down
41 changes: 34 additions & 7 deletions cmd/e2a/sending_policy.go
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@ type sendingProtectionFlags struct {
pauseClass string
evidenceRef string

policyFile string
expectedGeneration int64
expectedPolicySHA string
grandfather bool
Expand All @@ -73,6 +74,9 @@ func (f *sendingProtectionFlags) commandRequested() bool {
// modify, before anything starts: `-all` alone would otherwise be ignored
// and the server would boot as if nothing had been asked.
func (f *sendingProtectionFlags) validateStandalone() error {
if f.policyFile != "" && !f.inspect && !f.activate {
return errors.New("-sending-protection-policy-file requires -sending-protection-policy or -activate-sending-protection-policy")
}
if f.listAll && !f.listExternal {
return errors.New("-all is only valid with -list-external-sending-requests")
}
Expand Down Expand Up @@ -133,13 +137,37 @@ func runSendingProtectionCommand(ctx context.Context, cfg *config.Config, pool *
if err != nil {
return err
}
candidate := policy
if f.policyFile != "" {
candidate, err = readSendingPolicyFile(f.policyFile)
if err != nil {
return err
}
}
module := sendingpolicy.NewModule(pool, secrets)

switch {
case f.inspect:
return runPolicyInspect(ctx, module, source, policy, stdout)
if err := runPolicyInspect(ctx, module, source, policy, stdout); err != nil {
return err
}
if f.policyFile != "" {
hash, err := sendingpolicy.Hash(candidate)
if err != nil {
return err
}
canonical, err := sendingpolicy.CanonicalBytes(candidate)
if err != nil {
return err
}
fmt.Fprintf(stdout, "candidate_policy_sha256: %s\ncandidate_policy_canonical: %s\n", hash, canonical)
}
return nil
case f.activate:
return runPolicyActivate(ctx, module, policy, f, stdout)
if !candidate.AllControlsDisabled() && (secrets.Keyring == nil || secrets.Recipients == nil) {
return errors.New("activating enabled sending controls requires both sending-protection trust roots")
}
return runPolicyActivate(ctx, module, candidate, f, stdout)
case f.register:
return runOperatorRegister(ctx, module, secrets.Recipients, f, stdout)
case f.attest:
Expand Down Expand Up @@ -206,10 +234,9 @@ func runPolicyInspect(ctx context.Context, module *sendingpolicy.Module, source
return nil
}

// runPolicyActivate performs one reviewed CAS. There is no separately mounted
// policy file: the payload is built from the same validated config the server
// would use, and mutation requires the operator to present the exact hash they
// reviewed. Any mismatch is zero writes.
// runPolicyActivate performs one reviewed CAS using the selected file or config
// payload. The operator must present its exact canonical hash; a mismatch
// performs no policy or audit writes.
func runPolicyActivate(ctx context.Context, module *sendingpolicy.Module, policy sendingpolicy.RuntimePolicy, f *sendingProtectionFlags, stdout io.Writer) error {
if strings.TrimSpace(f.reason) == "" {
return errors.New("-activate-sending-protection-policy requires a nonblank -reason")
Expand All @@ -230,7 +257,7 @@ func runPolicyActivate(ctx context.Context, module *sendingpolicy.Module, policy
return err
}
if reviewedHash != configHash {
return fmt.Errorf("reviewed hash does not match this config's policy (%s); zero writes performed", configHash)
return fmt.Errorf("reviewed hash does not match the selected policy (%s); zero writes performed", configHash)
}

snapshot, err := module.ActivatePolicy(ctx, sendingpolicy.ActivationRequest{
Expand Down
37 changes: 37 additions & 0 deletions cmd/e2a/sending_policy_file.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package main

import (
"errors"
"io"
"os"

"github.com/tokencanopy/e2a/internal/sendingpolicy"
)

// readSendingPolicyFile reads one complete runtime-policy JSON document. This
// is a reviewed payload, not a server config: environment/config defaults must
// never alter it. Errors do not echo paths or possibly sensitive file contents.
func readSendingPolicyFile(path string) (sendingpolicy.RuntimePolicy, error) {
const maxBytes = 64 * 1024
info, err := os.Stat(path)
if err != nil || !info.Mode().IsRegular() {
return sendingpolicy.RuntimePolicy{}, errors.New("policy file must be a readable regular file")
}
file, err := os.Open(path)
if err != nil {
return sendingpolicy.RuntimePolicy{}, errors.New("cannot open policy file")
}
defer file.Close()
raw, err := io.ReadAll(io.LimitReader(file, maxBytes+1))
if err != nil {
return sendingpolicy.RuntimePolicy{}, errors.New("cannot read policy file")
}
if len(raw) > maxBytes {
return sendingpolicy.RuntimePolicy{}, errors.New("policy file exceeds 64 KiB")
}
policy, err := sendingpolicy.ParsePolicy(raw)
if err != nil {
return sendingpolicy.RuntimePolicy{}, errors.New("policy file does not match the runtime policy schema or validation rules")
}
return policy, nil
}
Loading
Loading