diff --git a/cmd/e2a/main.go b/cmd/e2a/main.go index 4e24ec487..8187fb09b 100644 --- a/cmd/e2a/main.go +++ b/cmd/e2a/main.go @@ -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") @@ -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 diff --git a/cmd/e2a/readyz.go b/cmd/e2a/readyz.go index 5b3dc2f70..05b5c1c50 100644 --- a/cmd/e2a/readyz.go +++ b/cmd/e2a/readyz.go @@ -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 @@ -46,6 +48,8 @@ const ( poolProbeTimeout = 250 * time.Millisecond ) +const readinessPolicyUnavailable = "sending protection policy unavailable" + type readinessState struct { ready bool reason string @@ -57,6 +61,7 @@ type readinessMonitor struct { pool *pgxpool.Pool draining *atomic.Bool latest string + policy *sendingpolicy.Module interval time.Duration grace time.Duration @@ -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, @@ -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}) } } @@ -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 } diff --git a/cmd/e2a/readyz_test.go b/cmd/e2a/readyz_test.go index f5c92c924..6832de7ba 100644 --- a/cmd/e2a/readyz_test.go +++ b/cmd/e2a/readyz_test.go @@ -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()) @@ -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) } @@ -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) } @@ -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() @@ -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 } @@ -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() @@ -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() diff --git a/cmd/e2a/sending_policy.go b/cmd/e2a/sending_policy.go index a97d75004..b34dcb8d5 100644 --- a/cmd/e2a/sending_policy.go +++ b/cmd/e2a/sending_policy.go @@ -49,6 +49,7 @@ type sendingProtectionFlags struct { pauseClass string evidenceRef string + policyFile string expectedGeneration int64 expectedPolicySHA string grandfather bool @@ -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") } @@ -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: @@ -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") @@ -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{ diff --git a/cmd/e2a/sending_policy_file.go b/cmd/e2a/sending_policy_file.go new file mode 100644 index 000000000..ef7bf2cfb --- /dev/null +++ b/cmd/e2a/sending_policy_file.go @@ -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 +} diff --git a/cmd/e2a/sending_policy_file_test.go b/cmd/e2a/sending_policy_file_test.go new file mode 100644 index 000000000..dfa83b8d7 --- /dev/null +++ b/cmd/e2a/sending_policy_file_test.go @@ -0,0 +1,170 @@ +package main + +import ( + "bytes" + "context" + "os" + "path/filepath" + "strings" + "testing" + + "github.com/tokencanopy/e2a/internal/sendingpolicy" + "github.com/tokencanopy/e2a/internal/testutil" +) + +func writePolicyFixture(t *testing.T, raw []byte) string { + t.Helper() + path := filepath.Join(t.TempDir(), "reviewed-policy.json") + if err := os.WriteFile(path, raw, 0600); err != nil { + t.Fatal(err) + } + return path +} + +func TestPolicyFileValidation(t *testing.T) { + p := sendingpolicy.DisabledPolicy() + raw, _ := sendingpolicy.CanonicalBytes(p) + good := writePolicyFixture(t, raw) + got, err := readSendingPolicyFile(good) + if err != nil { + t.Fatal(err) + } + wantHash, _ := sendingpolicy.Hash(p) + gotHash, _ := sendingpolicy.Hash(got) + if gotHash != wantHash { + t.Fatal("file changed the reviewed policy") + } + for name, input := range map[string][]byte{ + "unknown": []byte(`{"private-secret-field":"secret-value"}`), + "trailing": append(append([]byte{}, raw...), ']'), + "second object": append(append([]byte{}, raw...), raw...), + "oversized": []byte(strings.Repeat(" ", 65537)), + "null": []byte("null"), + } { + t.Run(name, func(t *testing.T) { + _, err := readSendingPolicyFile(writePolicyFixture(t, input)) + if err == nil { + t.Fatal("invalid policy file accepted") + } + if strings.Contains(err.Error(), "secret") { + t.Fatalf("file content leaked: %v", err) + } + }) + } + if _, err := readSendingPolicyFile(t.TempDir()); err == nil { + t.Fatal("directory accepted") + } + if err := (&sendingProtectionFlags{policyFile: good}).validateStandalone(); err == nil { + t.Fatal("orphan file flag accepted") + } + if err := (&sendingProtectionFlags{register: true, policyFile: good}).validateStandalone(); err == nil { + t.Fatal("file ignored by unrelated command") + } +} + +func TestPolicyFileCarriesApprovalGateIntoDatabase(t *testing.T) { + ctx := context.Background() + pool := testutil.TestDB(t) + cfg := spTestConfig() + recipients, err := sendingpolicy.LoadOperatorRecipients(spPolicyOperatorMap) + if err != nil { + t.Fatal(err) + } + secrets := sendingpolicy.Secrets{Recipients: recipients} + module := sendingpolicy.NewModule(pool, secrets) + if _, err := module.RegisterOperatorRecipients(ctx, "synthetic-operator", "test registration"); err != nil { + t.Fatal(err) + } + // The server's config remains unchanged; only the reviewed file carries the gate. + p, err := sendingpolicy.FromConfig(cfg) + if err != nil { + t.Fatal(err) + } + p.ExternalSendingAccess = &sendingpolicy.ExternalSendingAccessPolicy{Mode: sendingpolicy.ModeEnforce, AccountsCreatedAtOrAfter: "2000-01-01T00:00:00Z", Unlocks: []sendingpolicy.ExternalUnlock{sendingpolicy.UnlockOperatorApproval}} + raw, err := sendingpolicy.CanonicalBytes(p) + if err != nil { + t.Fatal(err) + } + path := writePolicyFixture(t, raw) + hash, _ := sendingpolicy.Hash(p) + var out bytes.Buffer + if err := runSendingProtectionCommand(ctx, cfg, pool, secrets, &sendingProtectionFlags{inspect: true, policyFile: path}, &out); err != nil { + t.Fatal(err) + } + if !strings.Contains(out.String(), "candidate_policy_sha256:") || !strings.Contains(out.String(), hash) { + t.Fatalf("missing candidate review: %s", out.String()) + } + before, _ := module.InspectPolicy(ctx) + if before.Generation != 0 { + t.Fatal("inspect wrote policy") + } + flags := &sendingProtectionFlags{activate: true, policyFile: path, expectedGeneration: 0, expectedPolicySHA: strings.Repeat("0", 64), reason: "carry admission gate"} + if err := runSendingProtectionCommand(ctx, cfg, pool, secrets, flags, &out); err == nil { + t.Fatal("wrong reviewed hash accepted") + } + still, _ := module.InspectPolicy(ctx) + if still.Generation != 0 { + t.Fatal("bad hash wrote policy") + } + flags.expectedPolicySHA = hash + if err := runSendingProtectionCommand(ctx, cfg, pool, secrets, flags, &out); err != nil { + t.Fatal(err) + } + stored, err := module.InspectPolicy(ctx) + if err != nil { + t.Fatal(err) + } + if stored.Generation != 1 || stored.Policy.ExternalSendingMode() != sendingpolicy.ModeEnforce || stored.Policy.BudgetMode != sendingpolicy.ModeDisabled || !stored.Policy.ExternalSendingAccess.Allows(sendingpolicy.UnlockOperatorApproval) || stored.Policy.ExternalSendingAccess.Allows(sendingpolicy.UnlockPaidEntitlement) { + t.Fatalf("carryover changed controls: %+v", stored) + } + if err := runSendingProtectionCommand(ctx, cfg, pool, secrets, flags, &out); err == nil { + t.Fatal("stale file replay accepted") + } +} + +func TestPolicyFileActivationRequiresTrustRootsForEnabledControls(t *testing.T) { + ctx := context.Background() + pool := testutil.TestDB(t) + recipients, err := sendingpolicy.LoadOperatorRecipients(spPolicyOperatorMap) + if err != nil { + t.Fatal(err) + } + secrets := sendingpolicy.Secrets{Recipients: recipients} + module := sendingpolicy.NewModule(pool, secrets) + if _, err = module.RegisterOperatorRecipients(ctx, "synthetic-operator", "test"); err != nil { + t.Fatal(err) + } + p := sendingpolicy.DisabledPolicy() + p.BudgetMode = sendingpolicy.ModeShadow + raw, _ := sendingpolicy.CanonicalBytes(p) + hash, _ := sendingpolicy.Hash(p) + var out bytes.Buffer + err = runSendingProtectionCommand(ctx, spTestConfig(), pool, secrets, &sendingProtectionFlags{activate: true, policyFile: writePolicyFixture(t, raw), expectedGeneration: 0, expectedPolicySHA: hash, reason: "test"}, &out) + if err == nil { + t.Fatal("file enabled budgets without the signing keyring") + } + stored, err := module.InspectPolicy(ctx) + if err != nil || stored.Generation != 0 { + t.Fatalf("missing trust root mutated policy: %+v %v", stored, err) + } +} + +func TestPolicyFileRejectsDuplicateKeys(t *testing.T) { + p := sendingpolicy.DisabledPolicy() + p.ExternalSendingAccess = &sendingpolicy.ExternalSendingAccessPolicy{Mode: sendingpolicy.ModeEnforce, AccountsCreatedAtOrAfter: "2000-01-01T00:00:00Z", Unlocks: []sendingpolicy.ExternalUnlock{sendingpolicy.UnlockOperatorApproval}} + raw, err := sendingpolicy.CanonicalBytes(p) + if err != nil { + t.Fatal(err) + } + for name, input := range map[string]string{ + "top level": strings.Replace(string(raw), `"budget_mode":"disabled"`, `"budget_mode":"shadow","budget_mode":"disabled"`, 1), + "escaped key": strings.Replace(string(raw), `"budget_mode":"disabled"`, `"budget_mode":"shadow","\u0062udget_mode":"disabled"`, 1), + "nested": strings.Replace(string(raw), `"mode":"enforce"`, `"mode":"disabled","mode":"enforce"`, 1), + } { + t.Run(name, func(t *testing.T) { + if _, err := readSendingPolicyFile(writePolicyFixture(t, []byte(input))); err == nil { + t.Fatal("ambiguous reviewed policy accepted") + } + }) + } +} diff --git a/cmd/e2a/sending_readiness_test.go b/cmd/e2a/sending_readiness_test.go new file mode 100644 index 000000000..c3699c15c --- /dev/null +++ b/cmd/e2a/sending_readiness_test.go @@ -0,0 +1,204 @@ +package main + +import ( + "context" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + "github.com/tokencanopy/e2a/internal/sendingpolicy" +) + +func readinessPolicySecrets(t *testing.T) sendingpolicy.Secrets { + t.Helper() + keys, err := sendingpolicy.LoadKeyring(`{"active":1,"keys":{"1":"AAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"}}`) + if err != nil { + t.Fatal(err) + } + recipients, err := sendingpolicy.LoadOperatorRecipients(spPolicyOperatorMap) + if err != nil { + t.Fatal(err) + } + return sendingpolicy.Secrets{Keyring: keys, Recipients: recipients} +} + +func assertReadiness(t *testing.T, m *readinessMonitor, status int) { + t.Helper() + rec := httptest.NewRecorder() + m.handler()(rec, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + if rec.Code != status { + t.Fatalf("readiness=%d %s, want %d", rec.Code, rec.Body.String(), status) + } + if status != http.StatusOK && rec.Body.String() != `{"status":"not_ready","reason":"sending protection policy unavailable"}` { + t.Fatalf("unbounded readiness error: %s", rec.Body.String()) + } +} + +func TestSendingReadinessRequiresStoredPolicyAndRegistry(t *testing.T) { + ctx := context.Background() + pool := migratedTestDB(t) + module := sendingpolicy.NewPolicyModule(pool, readinessPolicySecrets(t), sendingpolicy.PolicySourceDatabase, sendingpolicy.DisabledPolicy()) + m := newReadinessMonitorWithConfig(pool, nil, module, time.Hour, time.Hour, 2*time.Second) + defer m.Stop() + // Initial missing registry refuses readiness and does not auto-register. + assertReadiness(t, m, http.StatusServiceUnavailable) + var rows int + if err := pool.QueryRow(ctx, `SELECT count(*) FROM sending_operator_recipient_versions`).Scan(&rows); err != nil || rows != 0 { + t.Fatalf("readiness wrote registry: %d %v", rows, err) + } + if _, err := module.RegisterOperatorRecipients(ctx, "synthetic-operator", "bootstrap registry"); err != nil { + t.Fatal(err) + } + m.evaluate() + assertReadiness(t, m, http.StatusOK) + before, err := module.InspectPolicy(ctx) + if err != nil { + t.Fatal(err) + } + if _, err = pool.Exec(ctx, `UPDATE sending_protection_runtime_policy SET policy_sha256=repeat('0',64) WHERE singleton`); err != nil { + t.Fatal(err) + } + m.evaluate() + // A previously healthy slot must not retain the connectivity grace for a bad policy. + assertReadiness(t, m, http.StatusServiceUnavailable) + if _, err = pool.Exec(ctx, `UPDATE sending_protection_runtime_policy SET policy_sha256=$1 WHERE singleton`, before.PolicySHA256); err != nil { + t.Fatal(err) + } + m.evaluate() + assertReadiness(t, m, http.StatusOK) + p := before.Policy + p.OperatorNoticeRecipientVersion = 2 + raw, _ := sendingpolicy.CanonicalBytes(p) + hash, _ := sendingpolicy.Hash(p) + if _, err = pool.Exec(ctx, `UPDATE sending_protection_runtime_policy SET policy=$1,policy_sha256=$2 WHERE singleton`, raw, hash); err != nil { + t.Fatal(err) + } + m.evaluate() + assertReadiness(t, m, http.StatusServiceUnavailable) +} + +func TestSendingReadinessRegistryMismatchAndConfigFallback(t *testing.T) { + pool := migratedTestDB(t) + secrets := readinessPolicySecrets(t) + module := sendingpolicy.NewPolicyModule(pool, secrets, sendingpolicy.PolicySourceDatabase, sendingpolicy.DisabledPolicy()) + if _, err := module.RegisterOperatorRecipients(context.Background(), "synthetic-operator", "test"); err != nil { + t.Fatal(err) + } + wrong, err := sendingpolicy.LoadOperatorRecipients(strings.ReplaceAll(spPolicyOperatorMap, "policy-operator@example.test", "other-operator@example.test")) + if err != nil { + t.Fatal(err) + } + secrets.Recipients = wrong + for _, source := range []sendingpolicy.PolicySource{sendingpolicy.PolicySourceDatabase, sendingpolicy.PolicySourceConfig} { + t.Run(string(source), func(t *testing.T) { + m := newReadinessMonitorWithConfig(pool, nil, sendingpolicy.NewPolicyModule(pool, secrets, source, sendingpolicy.DisabledPolicy()), time.Hour, time.Hour, time.Second) + defer m.Stop() + status := http.StatusOK + if source == sendingpolicy.PolicySourceDatabase { + status = http.StatusServiceUnavailable + } + assertReadiness(t, m, status) + }) + } +} + +func TestSendingReadinessDoesNotBorrowSharedPool(t *testing.T) { + pool := migratedTestDB(t) + module := sendingpolicy.NewPolicyModule(pool, readinessPolicySecrets(t), sendingpolicy.PolicySourceDatabase, sendingpolicy.DisabledPolicy()) + if _, err := module.RegisterOperatorRecipients(context.Background(), "synthetic-operator", "test"); err != nil { + t.Fatal(err) + } + m := newReadinessMonitorWithConfig(pool, nil, module, time.Hour, 0, 2*time.Second) + defer m.Stop() + held := []*pgxpool.Conn{} + for i := int32(0); i < pool.Config().MaxConns; i++ { + conn, err := pool.Acquire(context.Background()) + if err != nil { + t.Fatal(err) + } + held = append(held, conn) + } + defer func() { + for _, conn := range held { + conn.Release() + } + }() + m.evaluate() + assertReadiness(t, m, http.StatusOK) +} + +func TestSendingReadinessMissingStoredPolicy(t *testing.T) { + pool := migratedTestDB(t) + module := sendingpolicy.NewPolicyModule(pool, readinessPolicySecrets(t), sendingpolicy.PolicySourceDatabase, sendingpolicy.DisabledPolicy()) + if _, err := pool.Exec(context.Background(), `DELETE FROM sending_protection_runtime_policy WHERE singleton`); err != nil { + t.Fatal(err) + } + m := newReadinessMonitorWithConfig(pool, nil, module, time.Hour, time.Hour, time.Second) + defer m.Stop() + assertReadiness(t, m, http.StatusServiceUnavailable) +} + +func TestSendingReadinessRequiresSelectedPermanentVersion(t *testing.T) { + ctx := context.Background() + pool := migratedTestDB(t) + secrets := readinessPolicySecrets(t) + module := sendingpolicy.NewPolicyModule(pool, secrets, sendingpolicy.PolicySourceDatabase, sendingpolicy.DisabledPolicy()) + if _, err := module.RegisterOperatorRecipients(ctx, "synthetic-operator", "test v1"); err != nil { + t.Fatal(err) + } + superset, err := sendingpolicy.LoadOperatorRecipients(`{"commitment_key":"` + spTestCommitmentKey + `","recipients":{"1":"policy-operator@example.test","2":"next-operator@example.test"}}`) + if err != nil { + t.Fatal(err) + } + secrets.Recipients = superset + module = sendingpolicy.NewPolicyModule(pool, secrets, sendingpolicy.PolicySourceDatabase, sendingpolicy.DisabledPolicy()) + p := sendingpolicy.DisabledPolicy() + p.OperatorNoticeRecipientVersion = 2 + raw, _ := sendingpolicy.CanonicalBytes(p) + hash, _ := sendingpolicy.Hash(p) + if _, err = pool.Exec(ctx, `UPDATE sending_protection_runtime_policy SET policy=$1,policy_sha256=$2 WHERE singleton`, raw, hash); err != nil { + t.Fatal(err) + } + m := newReadinessMonitorWithConfig(pool, nil, module, time.Hour, time.Hour, time.Second) + defer m.Stop() + assertReadiness(t, m, http.StatusServiceUnavailable) + if _, err = module.RegisterOperatorRecipients(ctx, "synthetic-operator", "test v2"); err != nil { + t.Fatal(err) + } + m.evaluate() + assertReadiness(t, m, http.StatusOK) +} + +func TestSendingReadinessPolicyTimeoutRecovers(t *testing.T) { + ctx := context.Background() + pool := migratedTestDB(t) + module := sendingpolicy.NewPolicyModule(pool, readinessPolicySecrets(t), sendingpolicy.PolicySourceDatabase, sendingpolicy.DisabledPolicy()) + if _, err := module.RegisterOperatorRecipients(ctx, "synthetic-operator", "test"); err != nil { + t.Fatal(err) + } + m := newReadinessMonitorWithConfig(pool, nil, module, time.Hour, time.Hour, 300*time.Millisecond) + defer m.Stop() + assertReadiness(t, m, http.StatusOK) + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + if _, err = tx.Exec(ctx, `LOCK TABLE sending_protection_runtime_policy IN ACCESS EXCLUSIVE MODE`); err != nil { + t.Fatal(err) + } + start := time.Now() + m.evaluate() + if elapsed := time.Since(start); elapsed > 2*time.Second { + t.Fatalf("policy probe escaped timeout: %v", elapsed) + } + assertReadiness(t, m, http.StatusServiceUnavailable) + if err = tx.Rollback(ctx); err != nil { + t.Fatal(err) + } + m.evaluate() + assertReadiness(t, m, http.StatusOK) +} diff --git a/docs/deployment.md b/docs/deployment.md index 467577f03..17f6cca6a 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -77,9 +77,10 @@ If you leave `shared_domain` empty, slug registration is disabled and every agen Wire your orchestrator to the two probe endpoints — they answer different questions and must not be swapped: `GET /api/health` is shallow liveness (restart policy; never checks the DB), `GET /readyz` is instance-local -readiness (DB reachable + migrations applied + not draining; use it for load +readiness (DB reachable + migrations applied + not draining; database-policy +mode also validates the runtime policy and selected recipient registry entry). Use it for load balancer routing so instances leave rotation during deploys and graceful -shutdown). Enable Prometheus metrics with the `metrics:` config block (or +shutdown. Enable Prometheus metrics with the `metrics:` config block (or `E2A_METRICS_ENABLED=true`) — exposition binds a separate loopback-default listener, never the public API handler. For continuous black-box monitoring of the full critical path (SMTP round-trip, outbound, WebSocket push, MCP), diff --git a/docs/runbooks/sending-policy-activation.md b/docs/runbooks/sending-policy-activation.md new file mode 100644 index 000000000..f852d8c10 --- /dev/null +++ b/docs/runbooks/sending-policy-activation.md @@ -0,0 +1,88 @@ +# Sending-policy activation and readiness + +These are operator-only server commands. They do not enable sending protection +by themselves. Runtime defaults remain unchanged, and hosted deployment, +rollback compatibility, and activation approvals belong to the deployment +runbook. Do not change the running server's config merely to prepare a future +policy generation. + +## Review a complete policy file + +`-sending-protection-policy-file` accepts a complete runtime-policy JSON object +(not a server YAML config), up to 64 KiB. Unknown fields, duplicate object keys, invalid values, +trailing content, and non-regular files are rejected. Config defaults and +environment overrides never alter the file's policy. Read/validation errors do +not echo file paths or contents. + +Inspect using the same binary and trust roots as the target deployment: + +```sh +e2a -config config.yaml -sending-protection-policy \ + -sending-protection-policy-file reviewed-policy.json +``` + +The existing `config_policy_*` fields describe the local config. With a file, +`candidate_policy_sha256` and `candidate_policy_canonical` describe the actual +activation candidate. Review that canonical payload and hash as well as the +stored generation. Inspection does not modify policy or audit rows. + +Activate the exact reviewed file: + +```sh +e2a -config config.yaml -activate-sending-protection-policy \ + -sending-protection-policy-file reviewed-policy.json \ + -expected-generation 0 \ + -expected-policy-sha256 "$REVIEWED_POLICY_SHA256" \ + -reason "Carry current admission policy into database source" +``` + +Use the generation actually inspected; `0` above illustrates the first +activation. A mismatched candidate hash, stale generation, absent/mismatched +selected recipient commitment, or missing trust roots for enabled controls +rejects the operation without policy/audit writes. Normal startup migrations +still run before operator commands. Successful activation writes the next +generation and audit event atomically. After a lost response, inspect before +retrying. Without the file flag, inspection and activation retain their +existing config-payload behavior. The file flag is invalid with other commands. + +## Preserve admission rules when switching sources + +The migration-seeded generation zero predates the external-sending admission +gate. A deployment using that gate must activate a reviewed generation carrying +its current effective admission mode, cohort cutoff, and unlock set **before** +selecting database source. In an operator-approval-only deployment, preserve +`mode: enforce` and `unlocks: [operator_approval]`; copying only the mode would +restore the older default unlock routes. Carry all other effective controls and +limits as well. This preparation need not enable budgets or trust progression. + +Register the selected recipient commitment explicitly with +`-register-sending-protection-operator-recipients` before activation. Readiness +never registers or repairs anything. Inspect the stored generation/hash and +compare effective behavior on every serving and rollback-capable slot. Keep the +same admission rules in the config fallback so changing the source back does +not remove the gate. A healthy readiness result proves validity, not that a +policy matches the operator's intended rollout: generation/hash and admission +carryover checks remain required deployment gates. + +## Database-source readiness + +In database-source mode the existing background readiness probe reads policy +and the selected permanent recipient commitment in one consistent read-only +transaction. It validates the supported schema, canonical hash, policy fields, +and exact agreement between the selected logical version, local commitment +key, and registry commitment. Both trust roots must be loaded. A missing, +unreadable, corrupt, or mismatched state returns HTTP 503 with the fixed reason +`sending protection policy unavailable`; raw errors and row contents are never +returned. It does not fall back to config or another recipient version. + +This check uses the monitor's dedicated database connection, not the shared +application pool, and performs no database work on the `/readyz` request path. +It shares the existing 2-second probe interval and 5-second probe timeout. +Policy failures bypass the connectivity failure grace period once observed; +recovery takes effect on the next successful probe. Ordinary connectivity +failures keep their existing 15-second grace, and draining still overrides +readiness immediately. `/api/health` remains shallow liveness. + +Config-source deployments do not depend on stored policy or registry rows for +readiness. This preserves the existing self-host defaults and the explicit +config-source rollback path. diff --git a/internal/sendingpolicy/policy.go b/internal/sendingpolicy/policy.go index 7d41bf933..86fd0a01a 100644 --- a/internal/sendingpolicy/policy.go +++ b/internal/sendingpolicy/policy.go @@ -551,7 +551,11 @@ func HashBytes(canonical []byte) string { // ParsePolicy decodes and validates a stored or configured policy. Unknown // fields are rejected: a payload written by a newer binary carrying a control // this one does not implement must fail closed, not be silently ignored. +// Duplicate object keys and trailing content are rejected before typed decoding. func ParsePolicy(raw []byte) (RuntimePolicy, error) { + if err := rejectDuplicateJSONKeys(string(raw), "policy"); err != nil { + return RuntimePolicy{}, err + } dec := json.NewDecoder(strings.NewReader(string(raw))) dec.DisallowUnknownFields() @@ -559,9 +563,6 @@ func ParsePolicy(raw []byte) (RuntimePolicy, error) { if err := dec.Decode(&p); err != nil { return RuntimePolicy{}, fmt.Errorf("sendingpolicy: decode policy: %w", err) } - if dec.More() { - return RuntimePolicy{}, fmt.Errorf("sendingpolicy: trailing content after policy object") - } if err := p.Validate(); err != nil { return RuntimePolicy{}, err } diff --git a/internal/sendingpolicy/readiness.go b/internal/sendingpolicy/readiness.go new file mode 100644 index 000000000..84ce54aab --- /dev/null +++ b/internal/sendingpolicy/readiness.go @@ -0,0 +1,43 @@ +package sendingpolicy + +import ( + "context" + "errors" + "time" + + "github.com/jackc/pgx/v5" +) + +// CheckReadiness validates database-source policy and its selected permanent +// recipient commitment against this process's trust roots in one read-only +// snapshot. It never repairs or registers rows. Config-source deployments have +// no dependency on those rows. The caller supplies its dedicated connection so +// readiness does not compete for the application's shared connection pool. +func (m *Module) CheckReadiness(ctx context.Context, conn *pgx.Conn) error { + if m.source == PolicySourceConfig { + return nil + } + if m.source != PolicySourceDatabase { + return errors.New("sendingpolicy: invalid readiness policy source") + } + if m.secrets.Keyring == nil || m.secrets.Recipients == nil { + return errors.New("sendingpolicy: database policy requires both trust roots") + } + tx, err := conn.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}) + if err != nil { + return err + } + defer func() { + cleanup, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + _ = tx.Rollback(cleanup) + }() + snapshot, err := scanPolicy(tx.QueryRow(ctx, policySelect)) + if err != nil { + return err + } + if err = m.checkSelectedOperatorRecipient(ctx, tx, snapshot.Policy.OperatorNoticeRecipientVersion, ""); err != nil { + return err + } + return tx.Commit(ctx) +} diff --git a/internal/sendingpolicy/store.go b/internal/sendingpolicy/store.go index bfa6b12d9..6eac78c8f 100644 --- a/internal/sendingpolicy/store.go +++ b/internal/sendingpolicy/store.go @@ -368,6 +368,12 @@ func (m *Module) ActivatePolicy(ctx context.Context, req ActivationRequest) (Pol // keeps the read/selection boundary explicit and composes with bootstrap or // repair tooling that might be running concurrently. func (m *Module) requireSelectedOperatorRecipient(ctx context.Context, tx pgx.Tx, version int) error { + return m.checkSelectedOperatorRecipient(ctx, tx, version, " FOR SHARE") +} + +// Readiness uses a consistent read-only snapshot; authorization and activation +// retain their share lock through the wrapper above. +func (m *Module) checkSelectedOperatorRecipient(ctx context.Context, tx pgx.Tx, version int, lockClause string) error { recipients := m.secrets.Recipients if recipients == nil { return fmt.Errorf("%w: local operator recipient map is not loaded", ErrOperatorRecipientUnavailable) @@ -382,7 +388,7 @@ func (m *Module) requireSelectedOperatorRecipient(ctx context.Context, tx pgx.Tx err := tx.QueryRow(ctx, `SELECT commitment_key_id, recipient_commitment FROM sending_operator_recipient_versions - WHERE logical_version = $1 FOR SHARE`, version, + WHERE logical_version = $1`+lockClause, version, ).Scan(®isteredKeyID, ®isteredCommitment) if errors.Is(err, pgx.ErrNoRows) { return fmt.Errorf("%w: logical version %d is not registered",