diff --git a/api/openapi.yaml b/api/openapi.yaml index fb837e04f..d09ce65ce 100644 --- a/api/openapi.yaml +++ b/api/openapi.yaml @@ -88,6 +88,35 @@ components: - scope - created_at type: object + AccountDailyLimit: + additionalProperties: true + properties: + clean_active_days: + format: int64 + type: integer + limit: + format: int64 + type: integer + resets_at: + format: date-time + type: string + shared_limit: + format: int64 + type: integer + shared_used: + format: int64 + type: integer + used: + format: int64 + type: integer + required: + - limit + - used + - shared_limit + - shared_used + - clean_active_days + - resets_at + type: object AccountMetricsView: additionalProperties: true properties: @@ -174,6 +203,9 @@ components: properties: agent_email: type: string + daily_limit: + $ref: "#/components/schemas/AccountDailyLimit" + description: External-recipient allowance for this UTC day. Used includes pending or uncertain provider submissions. Shared-identity usage is included in total usage and also bounded by shared_limit. Internal recipients (own live agents and verified owner mailbox) do not count. Omitted when this deployment does not enable the account trust ladder. deleted_at: description: When the account was moved to the trash. Absent for a live account. format: date-time @@ -2658,6 +2690,9 @@ components: description: The account's usage at the time the cap was hit (matches usage.). format: int64 type: integer + daily_limit: + $ref: "#/components/schemas/AccountDailyLimit" + description: Current external-recipient allowance, usage and reset time when the account trust ladder refuses an immediate send. limit: description: The cap that was hit (matches limits.max_). format: int64 @@ -2666,7 +2701,7 @@ components: description: The account's plan label. type: string resource: - description: "The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (per-UTC-day send cap; no AccountView field — resets at midnight UTC)." + description: "The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (daily send cap; daily_limit reports external-recipient usage on deployments with the account trust ladder — resets at midnight UTC)." type: string upgrade_url: description: An upgrade affordance URL, when the operator has configured one. @@ -4527,7 +4562,7 @@ components: format: int64 type: integer daily_recipient_limit: - description: Current UTC-day recipient allowance. Zero means no ramp cap applies. + description: Current UTC-day recipient allowance. Zero means no per-domain ramp cap applies; account daily_limit can still apply. format: int64 type: integer estimated_completion_at: @@ -4545,7 +4580,7 @@ components: format: date-time type: string status: - description: "Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt." + description: "Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt, account_managed (the account daily_limit replaces the per-domain ramp)." type: string required: - status diff --git a/cli/src/__tests__/whoami.test.ts b/cli/src/__tests__/whoami.test.ts index 2401a37d8..5f8de7388 100644 --- a/cli/src/__tests__/whoami.test.ts +++ b/cli/src/__tests__/whoami.test.ts @@ -40,6 +40,15 @@ describe("whoami command", () => { vi.clearAllMocks(); }); + it("shows the external daily allowance and shared subset", async () => { + mockAccountGet.mockResolvedValue(makeAccount({dailyLimit:{limit:88,used:17,sharedLimit:50,sharedUsed:6,cleanActiveDays:1,resetsAt:new Date("2026-01-02T00:00:00Z")}})); + const {whoami}=await import("../commands/whoami.js");await whoami({}); + const output=mockStdout.mock.calls.map((c:unknown[])=>c[0]).join(""); + expect(output).toContain("daily: 17/88 external recipients"); + expect(output).toContain("shared identity: 6/50"); + expect(output).toContain("2026-01-02T00:00:00.000Z"); + }); + it("prints identity, scope, plan, and usage for an account key", async () => { mockAccountGet.mockResolvedValue(makeAccount()); const { whoami } = await import("../commands/whoami.js"); diff --git a/cli/src/commands/whoami.ts b/cli/src/commands/whoami.ts index a1c11d478..f454b14b5 100644 --- a/cli/src/commands/whoami.ts +++ b/cli/src/commands/whoami.ts @@ -37,6 +37,10 @@ export async function whoami(opts: WhoamiOptions): Promise { `usage: ${account.usage.agents}/${account.limits.maxAgents} agents, ` + `${account.usage.messagesMonth}/${account.limits.maxMessagesMonth} messages this month\n`, ); + if (account.dailyLimit) { + const d=account.dailyLimit; + process.stdout.write(`daily: ${d.used}/${d.limit} external recipients; shared identity: ${d.sharedUsed}/${d.sharedLimit}; resets ${d.resetsAt.toISOString()}\n`); + } // Additive, optional: only ever present right after a dashboard restore // from the trash, so most accounts print nothing new here. if (account.restoredAt) { diff --git a/cmd/e2a/main.go b/cmd/e2a/main.go index c529ea1e3..2a7050932 100644 --- a/cmd/e2a/main.go +++ b/cmd/e2a/main.go @@ -894,6 +894,7 @@ func main() { }, time.Duration(cfg.Limits.CacheTTLSeconds)*time.Second, ) + enforcer.SetAccountDailyControl(outboundSending.module.AccountTrustEnabled) api.SetEnforcer(enforcer) // Master switch for the outbound footer; the enforcer above carries the // per-account entitlement + row-less default the decision reads. diff --git a/config.example.yaml b/config.example.yaml index 085054a23..24d2827eb 100644 --- a/config.example.yaml +++ b/config.example.yaml @@ -308,6 +308,13 @@ sender_identity: # established self-host-compatible default of 50, so the hashes intentionally # differ even though both policies are disabled. sending_ramp: + # Opt-in replacement: one account allowance (20 to 2,000), with shared + # identity capped at 50 and only external recipients counted. Independent + # of the legacy enabled/schedule fields below. Review existing-account + # seeding before activation; see docs/design/account-trust-ladder.md. + account_trust_enabled: false + # Retire probation/account/shared budget caps; global and notice pools stay. + disable_legacy_daily_budgets: false enabled: false start_daily: 50 target_daily: 2000 diff --git a/docs/api.md b/docs/api.md index bacaa60c7..a86ff470d 100644 --- a/docs/api.md +++ b/docs/api.md @@ -708,13 +708,15 @@ usually a truncated TXT.) Every domain response also carries **`sending_ramp`** — the platform-managed recipient-volume ramp state for newly verified custom sender domains: -`status` (open set; known values `inactive | ramping | complete | exempt`), -`daily_recipient_limit` (zero means no cap applies), `recipients_used_today`, +`status` (open set; known values `inactive | ramping | complete | exempt | account_managed`), +`daily_recipient_limit` (zero means no per-domain cap applies), `recipients_used_today`, `active_days` / `ramp_days`, and `resets_at` / `estimated_completion_at`. You can read this state but cannot change the schedule, exempt yourself, or reset progression through the API — see [`docs/runbooks/sending-ramp.md`](runbooks/sending-ramp.md) for the -operator-side mechanics. +operator-side mechanics. With `account_managed`, the optional `daily_limit` +object on `GET /v1/account` supplies the account-wide external allowance, +reserved usage, shared-identity subset, and UTC reset time. ### Agents (`/v1/agents`) diff --git a/docs/design/account-trust-ladder.md b/docs/design/account-trust-ladder.md new file mode 100644 index 000000000..79cfa9beb --- /dev/null +++ b/docs/design/account-trust-ladder.md @@ -0,0 +1,94 @@ +# Account sending trust ladder + +This opt-in control replaces the custom-domain ramp with one account-level +external-recipient allowance. It does not change monthly recipient-delivery +quotas or grant permission to email external recipients. Operator approval, +pauses, suppression, verified sending identity, and provider authorization +continue to apply independently. + +## Allowance and accounting + +The account starts at 20 external recipients per UTC day. Each completed clean +active day advances a linear schedule: `20 + floor(1980 * min(days, 29) / 29)`. +The thirtieth active day's allowance is 2,000. A day qualifies after at least +one provider-accepted external recipient; idle days and failed attempts earn +nothing. A detector breach disqualifies the day. The detector that records +breaches and reduces trust is a later implementation phase. + +Shared-identity recipients count toward the same account allowance and also +have a ceiling of 50 per day. The shared ceiling is a subset, not a separate +allowance. Adding or changing a domain never resets or multiplies trust. +The effective account allowance is the lower of trust and `max_messages_day`: +a null plan cap means unlimited, zero means no external sending, and accounts +without a provisioned limits row conservatively use 20. Thus bare Free remains +at 20; a paid plan or add-on removes only the plan cap, never the trust gate. + +External recipients are the admission gate's recipient classes: exclude the +account's live agents and its currently verified owner mailbox. Normalize and +deduplicate the complete To/Cc/Bcc envelope. Internal recipients still count +against monthly quota. This classification is repeated at authorization and +immediately before a provider call, so an old grant cannot acquire new external +recipients after an ownership or verification change. Any count change, including becoming internal, requires a fresh authorization. Legacy/all-recipient attempts never earn trust credit, including across runtime-policy toggles. + +Reservations serialize on an account row after the existing account-control, +operation, and budget locks. All identities compete for that same allowance. +Retries reuse the reservation; authoritative rejection/cancellation releases it; +uncertain outcomes keep it. Authoritative late acceptance restores a released +reservation. A refused midnight rollover retains the old reservation so delayed +provider evidence is not lost. Clean-day credit uses the settled attempt’s day and external units, separately from conservative reserved capacity; internal-only acceptance earns no credit. Late activity for a pruned historical bucket is not recreated, preventing a second award for archived days. Final authorization remains the enforcement +boundary; the immediate-send preflight is guidance, not a reservation. + +## Configuration and compatibility + +Both new settings default to false: + +```yaml +sending_ramp: + account_trust_enabled: true + disable_legacy_daily_budgets: true +``` + +`account_trust_enabled` selects the fixed account schedule instead of the legacy +per-domain ramp, even when the budget mode is disabled. Existing domain history +and legacy self-host behavior remain intact when it is false. +`disable_legacy_daily_budgets` removes the probation, account daily, and account +shared daily caps from budget admission while retaining their counter keys for +safe settlement across policy changes. Platform and operational notice pools +remain independently governed by their existing policy. External-only platform counting starts with `account_trust_enabled`; disabling legacy caps alone retains the existing all-recipient platform accounting. A preceding shadow phase must not interpret those counters as external-only evidence. + +The corresponding optional runtime-policy keys are omitted when false, preserving +canonical hashes of policies written before these keys existed. Database-sourced +policy remains authoritative at runtime, including daily quota delegation. + +## Visibility + +`GET /v1/account` has an optional `daily_limit` object with `limit`, `used`, +`shared_limit`, `shared_used`, `clean_active_days`, and `resets_at`. Usage includes +pending or uncertain reservations. The endpoint returns `limits_unavailable` +when the enabled control cannot be read instead of implying that it is absent. +Internal/system accounts are exempt and omit the object. + +An immediate request exceeding the allowance returns 402 `limit_exceeded`, +resource `messages_day`, and the same snapshot in `error.details.daily_limit`. +Scheduled sends are checked when they fire. Queued sends encountering the cap +are held until UTC midnight, subject to the existing finite retry horizon; +queued-message and terminal hold diagnostics include usage, allowance, and reset time. +The generated TypeScript and Python models, CLI `whoami`, MCP `whoami`, and +usage dashboard carry the same contract. No platform-wide capacity is exposed. + +## Persistence and rollout + +Migrations 130–131 add account trust, daily buckets, message reservations, and immutable per-attempt external-unit provenance without +activating the control or exempting existing accounts. Foreign keys erase these +rows with their owning account. Daily maintenance prunes settled history older +than 90 days, folding clean-day credit into the account row first. It retains +unresolved reservations and their buckets. Each run bounds both accounts and +rows processed and skips active account locks. + +This PR supplies implementation, not hosted activation. The activation change +must seed already-approved accounts from reviewed observed external volume before +enabling the ladder; `grandfather_daily` supports a bounded floor (0–2,000) without +bypassing the shared ceiling or plan cap. No send-path code sets that floor. +The hosted policy-source switch, reviewed grandfathering command/payload, +seven-day observation gate, detector automation, and production enablement remain +separate rollout work. User-facing rollout guidance belongs with enablement. diff --git a/internal/agent/external_access.go b/internal/agent/external_access.go index 5622815d9..2bf97bb73 100644 --- a/internal/agent/external_access.go +++ b/internal/agent/external_access.go @@ -12,6 +12,7 @@ import ( "github.com/tokencanopy/e2a/internal/outbound" "github.com/tokencanopy/e2a/internal/sendingpolicy" + "github.com/tokencanopy/e2a/internal/sendramp" ) // ExternalSendingNotEnabledCode is the stable 403 code for a send the account @@ -65,6 +66,23 @@ func (a *API) preflightExternalAccess(ctx context.Context, userID, agentID strin if verdict.Denied() { return a.externalSendingNotEnabledError(ctx, userID) } + if req.ScheduledAt == nil { + if daily, ok := a.externalAccess.(interface { + DailyLimitPreflight(context.Context, string, string, []string) (*sendramp.AccountDailyLimit, error) + }); ok { + d, err := daily.DailyLimitPreflight(ctx, userID, agentID, recipients) + if err != nil { + return &OutboundError{Status: http.StatusServiceUnavailable, Code: "limits_unavailable", Msg: "could not verify daily sending limit; retry shortly"} + } + if d != nil && !d.Allowed { + limit, used := d.Limit, d.Used + if d.SharedBinding { + limit, used = d.SharedLimit, d.SharedUsed + } + return &OutboundError{Status: http.StatusPaymentRequired, Code: "limit_exceeded", Msg: "daily external-recipient allowance reached; wait until the UTC reset", Details: map[string]any{"resource": "messages_day", "limit": limit, "current": used, "daily_limit": d}} + } + } + } return nil } diff --git a/internal/agent/external_access_test.go b/internal/agent/external_access_test.go index 004dd945e..cea4151a2 100644 --- a/internal/agent/external_access_test.go +++ b/internal/agent/external_access_test.go @@ -194,3 +194,25 @@ func TestDeliverOutboundExternalAccessMessageFollowsUnlocks(t *testing.T) { }) } } + +func TestDailyLimitSharedRefusalReportsBindingCap(t *testing.T) { + api, store, _, _, pool := setupAsyncAPIWithPool(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + api.SetExternalAccess(sendingpolicy.NewPolicyModule(pool, sendingpolicy.Secrets{}, sendingpolicy.PolicySourceConfig, p)) + ctx := context.Background() + user, ag := selfAgent(t, store, "sharedcap") + if _, err := pool.Exec(ctx, `INSERT INTO account_limits(user_id,max_agents,max_domains,max_messages_month,max_storage_bytes,max_messages_day) VALUES($1,100,10,10000,1073741824,NULL) ON CONFLICT(user_id) DO UPDATE SET max_messages_day=NULL`, user.ID); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `INSERT INTO account_sending_trust(user_id,grandfather_daily) VALUES($1,2000)`, user.ID); err != nil { + t.Fatal(err) + } + if _, err := pool.Exec(ctx, `INSERT INTO account_send_days(user_id,day,reserved_count,shared_count) VALUES($1,(clock_timestamp() AT TIME ZONE 'UTC')::date,50,50)`, user.ID); err != nil { + t.Fatal(err) + } + _, e := api.DeliverOutbound(ctx, user, ag, outbound.SendRequest{To: []string{"external@example.test"}, Subject: "synthetic", Body: "test"}, "send", "", nil, nil) + if e == nil || e.Status != 402 || e.Details["limit"] != 50 || e.Details["current"] != 50 { + t.Fatalf("shared refusal: %+v", e) + } +} diff --git a/internal/agent/outbound_async.go b/internal/agent/outbound_async.go index 5a8f99e10..77067fb06 100644 --- a/internal/agent/outbound_async.go +++ b/internal/agent/outbound_async.go @@ -212,6 +212,15 @@ func (a *outboundSendStore) RecordHold(ctx context.Context, messageID string, cl return a.store.RecordOutboundHold(ctx, messageID, string(class), anchor) } +// RecordHoldDetail persists only the current claim's safe quota diagnostic. +// The terminal sent/failed transitions clear or replace this field. +func (a *outboundSendStore) RecordHoldDetail(ctx context.Context, messageID string, jobID int64, detail string) error { + return a.store.WithTx(ctx, func(tx pgx.Tx) error { + _, err := tx.Exec(ctx, `UPDATE messages SET delivery_detail=$3 WHERE id=$1 AND send_job_id=$2 AND delivery_status IN ('accepted','sending')`, messageID, jobID, messagelifecycle.SafeDiagnostic(detail)) + return err + }) +} + // SuppressedRecipients backs the SendWorker's pre-provider suppression guard: // the effective account-wide + exact-agent subset (the store normalizes both // sides). diff --git a/internal/apiserver/apiserver.go b/internal/apiserver/apiserver.go index 1f13f62a4..eeef3937f 100644 --- a/internal/apiserver/apiserver.go +++ b/internal/apiserver/apiserver.go @@ -132,7 +132,19 @@ func BuildDeps(p Params) httpapi.Deps { } var rampSnapshot func(context.Context, string, string, time.Time) (sendramp.Snapshot, error) if p.Pool != nil { - rampSnapshot = sendramp.NewStore(p.Pool).Snapshot + legacyRamp := sendramp.NewStore(p.Pool) + rampSnapshot = func(ctx context.Context, user, domain string, now time.Time) (sendramp.Snapshot, error) { + if p.SendingAccess != nil { + enabled, err := p.SendingAccess.AccountTrustEnabled(ctx) + if err != nil { + return sendramp.Snapshot{}, err + } + if enabled { + return sendramp.Snapshot{Status: "account_managed"}, nil + } + } + return legacyRamp.Snapshot(ctx, user, domain, now) + } } var listMessageLifecycle httpapi.MessageLifecycleLister var countAgentMetrics httpapi.AgentMetricsCounter @@ -342,6 +354,7 @@ func BuildDeps(p Params) httpapi.Deps { Metrics: p.Metrics, } if p.SendingAccess != nil { + deps.AccountDailyLimit = p.SendingAccess.AccountDailyLimit deps.SendingAccessStatus = p.SendingAccess.ExternalAccessStatus deps.SubmitSendingAccessRequest = p.SendingAccess.SubmitAccessRequest deps.LatestSendingAccessRequest = p.SendingAccess.LatestAccessRequest diff --git a/internal/config/config.go b/internal/config/config.go index d3a6da3b3..0164ab123 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -520,10 +520,12 @@ type SenderIdentityConfig struct { // public API. Values are snapshotted when a domain first sends, so later config // changes do not reshape an in-flight ramp. type SendingRampConfig struct { - Enabled bool `yaml:"enabled"` - StartDaily int `yaml:"start_daily"` - TargetDaily int `yaml:"target_daily"` - RampDays int `yaml:"ramp_days"` + DisableLegacyDailyBudgets bool `yaml:"disable_legacy_daily_budgets"` + AccountTrustEnabled bool `yaml:"account_trust_enabled"` + Enabled bool `yaml:"enabled"` + StartDaily int `yaml:"start_daily"` + TargetDaily int `yaml:"target_daily"` + RampDays int `yaml:"ramp_days"` } // SendingProtectionConfig carries the non-schedule half of the sending diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 0bbc28e63..519a7b4df 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -530,7 +530,7 @@ func TestSendingRampDefaultsOverridesAndValidation(t *testing.T) { if err != nil { t.Fatalf("Load defaults: %v", err) } - if cfg.SendingRamp.Enabled || cfg.SendingRamp.StartDaily != 50 || cfg.SendingRamp.TargetDaily != 2000 || cfg.SendingRamp.RampDays != 30 { + if cfg.SendingRamp.AccountTrustEnabled || cfg.SendingRamp.DisableLegacyDailyBudgets || cfg.SendingRamp.Enabled || cfg.SendingRamp.StartDaily != 50 || cfg.SendingRamp.TargetDaily != 2000 || cfg.SendingRamp.RampDays != 30 { t.Fatalf("sending ramp defaults = %+v, want disabled 50/2000/30", cfg.SendingRamp) } diff --git a/internal/httpapi/account.go b/internal/httpapi/account.go index 46e4e789e..ba557bf00 100644 --- a/internal/httpapi/account.go +++ b/internal/httpapi/account.go @@ -9,6 +9,7 @@ import ( "github.com/danielgtaylor/huma/v2" "github.com/tokencanopy/e2a/internal/identity" + "github.com/tokencanopy/e2a/internal/sendramp" ) // AccountUserView is the authenticated principal's identity (A-1). Returned by @@ -25,7 +26,8 @@ type AccountUserView struct { // credentials (where the credential *is* a single agent) — omitted for // account scope, which spans many agents. type AccountView struct { - User AccountUserView `json:"user"` + DailyLimit *sendramp.AccountDailyLimit `json:"daily_limit,omitempty" doc:"External-recipient allowance for this UTC day. Used includes pending or uncertain provider submissions. Shared-identity usage is included in total usage and also bounded by shared_limit. Internal recipients (own live agents and verified owner mailbox) do not count. Omitted when this deployment does not enable the account trust ladder."` + User AccountUserView `json:"user"` // Scope is an OPEN set on this response view (evolving vocabulary), not a // closed enum — see docs/api.md "Versioning & stability". Scope string `json:"scope" doc:"Credential scope. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: account, agent."` @@ -345,6 +347,13 @@ func (s *Server) handleGetMyLimits(ctx context.Context, _ *struct{}) (*accountOu if err != nil { return nil, NewError(http.StatusInternalServerError, "internal_error", "limits lookup failed") } + var daily *sendramp.AccountDailyLimit + if s.deps.AccountDailyLimit != nil { + daily, err = s.deps.AccountDailyLimit(ctx, user.ID) + if err != nil { + return nil, NewError(http.StatusServiceUnavailable, "limits_unavailable", "daily sending limit is temporarily unavailable").WithDetails(RetryAfterDetails{RetryAfterSeconds: limitsUnavailableRetrySeconds}).WithRetryAfter(limitsUnavailableRetrySeconds) + } + } var usage LimitsUsageView if s.deps.GetUsage != nil { usage = s.deps.GetUsage(ctx, user.ID) @@ -356,6 +365,7 @@ func (s *Server) handleGetMyLimits(ctx context.Context, _ *struct{}) (*accountOu agentAddress = p.AgentID } return &accountOutput{Body: AccountView{ + DailyLimit: daily, User: AccountUserView{ID: user.ID, Email: user.Email}, Scope: p.Scope, AgentAddress: agentAddress, diff --git a/internal/httpapi/account_daily_limit_test.go b/internal/httpapi/account_daily_limit_test.go new file mode 100644 index 000000000..1a8a0821d --- /dev/null +++ b/internal/httpapi/account_daily_limit_test.go @@ -0,0 +1,44 @@ +package httpapi + +import ( + "context" + "errors" + "github.com/tokencanopy/e2a/internal/sendramp" + "testing" + "time" +) + +func TestAccountDailyLimit(t *testing.T) { + reset := time.Date(2026, 1, 2, 0, 0, 0, 0, time.UTC) + srv := testServer(t, func(d *Deps) { + d.AccountDailyLimit = func(_ context.Context, user string) (*sendramp.AccountDailyLimit, error) { + if user != "u_1" { + t.Errorf("wrong tenant: %s", user) + } + return &sendramp.AccountDailyLimit{Limit: 88, Used: 17, SharedLimit: 50, SharedUsed: 6, CleanActiveDays: 1, ResetsAt: reset}, nil + } + }) + code, body := getJSON(t, srv.URL+"/v1/account", "good") + if code != 200 { + t.Fatalf("%d %v", code, body) + } + d, ok := body["daily_limit"].(map[string]any) + if !ok || d["limit"] != float64(88) || d["used"] != float64(17) || d["resets_at"] != "2026-01-02T00:00:00Z" { + t.Fatalf("daily_limit: %v", d) + } + if _, leaked := d["allowed"]; leaked { + t.Fatal("internal decision leaked") + } +} + +func TestAccountDailyLimitUnavailable(t *testing.T) { + srv := testServer(t, func(d *Deps) { + d.AccountDailyLimit = func(context.Context, string) (*sendramp.AccountDailyLimit, error) { + return nil, errors.New("unavailable") + } + }) + code, body := getJSON(t, srv.URL+"/v1/account", "good") + if code != 503 || errCode(body) != "limits_unavailable" { + t.Fatalf("%d %v", code, body) + } +} diff --git a/internal/httpapi/domains.go b/internal/httpapi/domains.go index d130a52f5..a485b9a52 100644 --- a/internal/httpapi/domains.go +++ b/internal/httpapi/domains.go @@ -84,8 +84,8 @@ type DomainView struct { } type SendingRampView struct { - Status string `json:"status" doc:"Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt."` - DailyRecipientLimit int `json:"daily_recipient_limit" doc:"Current UTC-day recipient allowance. Zero means no ramp cap applies."` + Status string `json:"status" doc:"Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt, account_managed (the account daily_limit replaces the per-domain ramp)."` + DailyRecipientLimit int `json:"daily_recipient_limit" doc:"Current UTC-day recipient allowance. Zero means no per-domain ramp cap applies; account daily_limit can still apply."` RecipientsUsedToday int `json:"recipients_used_today" doc:"Recipient capacity reserved for the current UTC day, including submissions whose provider outcome is still pending."` ResetsAt *time.Time `json:"resets_at,omitempty"` ActiveDays int `json:"active_days" doc:"UTC days that reached the provider-accepted volume threshold."` diff --git a/internal/httpapi/errors.go b/internal/httpapi/errors.go index cce886ab9..c1042bd83 100644 --- a/internal/httpapi/errors.go +++ b/internal/httpapi/errors.go @@ -16,6 +16,7 @@ package httpapi import ( "encoding/json" + "github.com/tokencanopy/e2a/internal/sendramp" "net/http" "reflect" "strconv" @@ -159,10 +160,11 @@ type PayloadTooLargeDetails struct { // time. `plan_code`/`upgrade_url` are the account's plan label and any upgrade // affordance the operator configured. type LimitExceededDetails struct { + DailyLimit *sendramp.AccountDailyLimit `json:"daily_limit,omitempty" doc:"Current external-recipient allowance, usage and reset time when the account trust ladder refuses an immediate send."` // Resource is an OPEN set (evolving response-side vocabulary): a new // capped resource means a new value here, and that must not break // spec-generated clients. - Resource string `json:"resource" doc:"The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (per-UTC-day send cap; no AccountView field — resets at midnight UTC)."` + Resource string `json:"resource" doc:"The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (daily send cap; daily_limit reports external-recipient usage on deployments with the account trust ladder — resets at midnight UTC)."` Limit int64 `json:"limit" doc:"The cap that was hit (matches limits.max_)."` Current int64 `json:"current" doc:"The account's usage at the time the cap was hit (matches usage.)."` PlanCode string `json:"plan_code,omitempty" doc:"The account's plan label."` diff --git a/internal/httpapi/httpapi.go b/internal/httpapi/httpapi.go index 3d2362424..9c3c43faa 100644 --- a/internal/httpapi/httpapi.go +++ b/internal/httpapi/httpapi.go @@ -147,6 +147,7 @@ type MessageRestoreOp func(ctx context.Context, messageID, agentID string) (*ide // Deps are the collaborators the v1 layer needs. Everything is injected so // the package has no hidden globals and is straightforward to test. type Deps struct { + AccountDailyLimit func(context.Context, string) (*sendramp.AccountDailyLimit, error) // External sending access (all optional; nil = surface absent/501). // SendingAccessStatus backs the additive sending_access object on GET // /v1/account; the request pair backs /v1/account/sending-access/request; diff --git a/internal/limits/enforcer.go b/internal/limits/enforcer.go index 2ddfedff0..dba2aae8e 100644 --- a/internal/limits/enforcer.go +++ b/internal/limits/enforcer.go @@ -40,10 +40,11 @@ type limitsReader interface { // check; the win from caching limits is avoiding the join into // account_limits, which is the costlier read. type DBEnforcer struct { - store limitsReader - counter Counter - defaults Defaults - cacheTTL time.Duration + accountDailyControl func(context.Context) (bool, error) + store limitsReader + counter Counter + defaults Defaults + cacheTTL time.Duration mu sync.Mutex cache map[string]cachedLimits @@ -222,6 +223,12 @@ func (e *DBEnforcer) CheckDomainCreate(ctx context.Context, userID string) error return nil } +// SetAccountDailyControl delegates daily enforcement to the external-recipient +// authorization ledger. Install before serving; runtime policy reads fail closed. +func (e *DBEnforcer) SetAccountDailyControl(f func(context.Context) (bool, error)) { + e.accountDailyControl = f +} + // CheckMessageSend enforces the month-flow cap and the storage stock // cap for an outbound send of `units` recipient-deliveries. Either // being exceeded blocks the operation. The flow cap is checked first @@ -255,7 +262,14 @@ func (e *DBEnforcer) CheckMessageSend(ctx context.Context, userID string, units // the self-host default and every paid shape — so the extra count read // only happens for accounts that actually carry the cap. Resets at UTC // midnight with the usage_summaries bucket_date. - if lim.MaxMessagesDay != nil { + accountDaily := false + if e.accountDailyControl != nil { + accountDaily, err = e.accountDailyControl(ctx) + if err != nil { + return err + } + } + if lim.MaxMessagesDay != nil && !accountDaily { dayCount, err := e.counter.MessagesToday(ctx, userID) if err != nil { return err diff --git a/internal/limits/units_test.go b/internal/limits/units_test.go index dc8456794..4fe255b9b 100644 --- a/internal/limits/units_test.go +++ b/internal/limits/units_test.go @@ -143,3 +143,20 @@ func TestCheckMessageSend_DailyCap(t *testing.T) { } }) } + +func TestAccountDailyControlPreservesMonthlyQuota(t *testing.T) { + store := &fakeStore{found: true, row: Limits{MaxMessagesMonth: 100, MaxMessagesDay: intPtr(20), MaxStorageBytes: 1 << 40}} + counter := &fakeCounter{messagesMonth: 99, messagesToday: 1000} + e := newEnforcerWithReader(store, counter, defaultsForTest(), 0) + e.SetAccountDailyControl(func(context.Context) (bool, error) { return true, nil }) + if err := e.CheckMessageSend(context.Background(), "u", 1); err != nil { + t.Fatalf("legacy daily gate not delegated: %v", err) + } + if le, ok := IsLimitExceeded(e.CheckMessageSend(context.Background(), "u", 2)); !ok || le.Resource != "messages_month" { + t.Fatalf("monthly gate changed: %v", le) + } + e.SetAccountDailyControl(func(context.Context) (bool, error) { return false, errors.New("policy unavailable") }) + if err := e.CheckMessageSend(context.Background(), "u", 1); err == nil { + t.Fatal("unreadable policy failed open") + } +} diff --git a/internal/outboundsend/gate_worker_test.go b/internal/outboundsend/gate_worker_test.go index fd08ff137..e3528a5e6 100644 --- a/internal/outboundsend/gate_worker_test.go +++ b/internal/outboundsend/gate_worker_test.go @@ -3,6 +3,8 @@ package outboundsend_test import ( "context" "errors" + "github.com/tokencanopy/e2a/internal/sendramp" + "strings" "testing" "time" @@ -555,3 +557,16 @@ func TestGatedWorker_ExternalSendingNotEnabledFailsWithItsReason(t *testing.T) { t.Fatalf("legacy refusal: err=%v failed=%+v", err, st.failed) } } + +func TestGatedWorkerPersistsDailyLimitHoldDetail(t *testing.T) { + st := &fakeStore{job: acceptedJob("msg_daily_hold")} + reset := time.Now().UTC().Add(time.Hour) + g := &fakeGate{reserve: sendingpolicy.Decision{Reason: sendingpolicy.ReasonRampCapacity, RetryAt: reset, DailyLimit: &sendramp.AccountDailyLimit{Limit: 20, Used: 20, SharedLimit: 20, SharedUsed: 20, ResetsAt: reset}}} + err := outboundsend.NewSendWorker(st, &fakeDeliverer{}).WithGate(g).Work(context.Background(), gatedJob("msg_daily_hold", 1)) + if !isSnooze(err) { + t.Fatalf("hold: %v", err) + } + if len(st.holdDetails) != 1 || !strings.Contains(st.holdDetails[0], "limit=20 used=20") || !strings.Contains(st.holdDetails[0], reset.Format(time.RFC3339)) { + t.Fatalf("missing persisted allowance/reset: %v", st.holdDetails) + } +} diff --git a/internal/outboundsend/reconcile_test.go b/internal/outboundsend/reconcile_test.go index c048843bc..eb656509e 100644 --- a/internal/outboundsend/reconcile_test.go +++ b/internal/outboundsend/reconcile_test.go @@ -1093,3 +1093,7 @@ func TestRegisterJobs_RegistersTerminalReconcilePeriodic(t *testing.T) { t.Fatalf("RegisterJobs periodics = %d, want 1", len(periodics)) } } + +func (s failingTerminalStore) RecordHoldDetail(context.Context, string, int64, string) error { + return nil +} diff --git a/internal/outboundsend/worker.go b/internal/outboundsend/worker.go index 1ffc2ac6f..30c3723a3 100644 --- a/internal/outboundsend/worker.go +++ b/internal/outboundsend/worker.go @@ -316,6 +316,8 @@ type Store interface { // RecordHold persists the message's finite-hold class and anchor. Terminal // writes clear the pair. RecordHold(ctx context.Context, messageID string, class HoldClass, anchor time.Time) error + // RecordHoldDetail exposes a current hold diagnostic on the queued message. + RecordHoldDetail(ctx context.Context, messageID string, jobID int64, detail string) error // MarkSent records the provider outcome monotonically from a pre-terminal // state, including when trash won after ClaimSend. MarkSent(ctx context.Context, messageID string, jobID int64, attempt int, occurredAt time.Time, providerMessageID, sentAs string) error @@ -748,7 +750,14 @@ func (w *SendWorker) hold(ctx context.Context, job *river.Job[OutboundSendArgs], } return river.JobSnooze(delay) } - return w.holdFinite(ctx, job, j, attempt, class, "sending_policy_hold: "+d.Reason, delay, observedAt) + detail := "sending_policy_hold: " + d.Reason + if d.DailyLimit != nil { + detail += fmt.Sprintf(" limit=%d used=%d shared_limit=%d shared_used=%d resets_at=%s", d.DailyLimit.Limit, d.DailyLimit.Used, d.DailyLimit.SharedLimit, d.DailyLimit.SharedUsed, d.DailyLimit.ResetsAt.Format(time.RFC3339)) + if err := w.store.RecordHoldDetail(ctx, j.MessageID, job.ID, detail); err != nil { + return fmt.Errorf("record daily limit hold: %w", err) + } + } + return w.holdFinite(ctx, job, j, attempt, class, detail, delay, observedAt) } // holdFinite persists the hold state, expires the message when its derived diff --git a/internal/outboundsend/worker_test.go b/internal/outboundsend/worker_test.go index 1bcd55426..5de2b6d49 100644 --- a/internal/outboundsend/worker_test.go +++ b/internal/outboundsend/worker_test.go @@ -36,12 +36,13 @@ type fakeStore struct { suppressed []string suppressedErr error - sent []sentCall - holds []holdCall - failed []failedCall - deferred []failedCall - temporary []failedCall - released []string + holdDetails []string + sent []sentCall + holds []holdCall + failed []failedCall + deferred []failedCall + temporary []failedCall + released []string // suppressionUserID records the tenant the guard was scoped to. suppressionUserID string suppressionAgentID string @@ -443,3 +444,8 @@ func gatedJob(id string, attempt int) *river.Job[outboundsend.OutboundSendArgs] j.Args.OperationRef = &ref return j } + +func (f *fakeStore) RecordHoldDetail(_ context.Context, _ string, _ int64, detail string) error { + f.holdDetails = append(f.holdDetails, detail) + return nil +} diff --git a/internal/sendingpolicy/account_trust.go b/internal/sendingpolicy/account_trust.go new file mode 100644 index 000000000..a87ebc4bf --- /dev/null +++ b/internal/sendingpolicy/account_trust.go @@ -0,0 +1,227 @@ +package sendingpolicy + +import ( + "context" + "errors" + + "github.com/jackc/pgx/v5" + "github.com/tokencanopy/e2a/internal/sendramp" +) + +type accountCapacityError struct{ daily sendramp.AccountDailyLimit } + +func (e *accountCapacityError) Error() string { return errRampCapacity.Error() } +func (e *accountCapacityError) Unwrap() error { return errRampCapacity } + +// externalRecipientCount shares the admission gate's live-agent and verified +// owner definitions, independently of approval status or available unlocks. +func externalRecipientCount(ctx context.Context, q dbQuerier, user string, envelope []string) (int, error) { + facts, err := loadAccountAccessFacts(ctx, q, user) + if err != nil { + return 0, err + } + own, err := ownAgentRecipients(ctx, q, user, envelope) + if err != nil { + return 0, err + } + seen := make(map[string]bool) + n := 0 + for _, raw := range envelope { + addr := NormalizeOwnerMailbox(raw) + if seen[addr] { + continue + } + seen[addr] = true + if _, ok := own[addr]; ok { + continue + } + if facts.ownerRecipientVerified() && addr == facts.proofAddress { + continue + } + n++ + } + return n, nil +} + +func accountPlanCap(ctx context.Context, q dbQuerier, user string) (*int, error) { + var cap *int + err := q.QueryRow(ctx, `SELECT max_messages_day FROM account_limits WHERE user_id=$1`, user).Scan(&cap) + if errors.Is(err, pgx.ErrNoRows) { + v := 20 + return &v, nil + } + return cap, err +} + +func (m *Module) accountTrustSubject(ctx context.Context, tx pgx.Tx, op operationRow) (rampSubject, error) { + if op.Purpose != PurposeCustomerMessage { + return rampSubject{}, nil + } + facts, err := loadAccountAccessFacts(ctx, tx, op.accountRef()) + if err != nil { + return rampSubject{}, err + } + if accountClassExempt(facts.class) { + return rampSubject{}, nil + } + envelope, err := messageEnvelopeQ(ctx, tx, op.OperationID) + if err != nil { + return rampSubject{}, err + } + n, err := externalRecipientCount(ctx, tx, op.accountRef(), envelope) + if err != nil { + return rampSubject{}, err + } + if n == 0 { + return rampSubject{}, nil + } + if !op.Shared { + agent, _, err := messageSender(ctx, tx, op.OperationID) + if err != nil { + return rampSubject{}, err + } + verified, err := customIdentityVerified(ctx, tx, op.accountRef(), agent) + if err != nil { + return rampSubject{}, err + } + if !verified { + return rampSubject{}, errRampIdentityUnverified + } + } + cap, err := accountPlanCap(ctx, tx, op.accountRef()) + if err != nil { + return rampSubject{}, err + } + return rampSubject{messageID: op.OperationID, userID: op.accountRef(), units: n, applies: true, account: true, shared: op.Shared, planCap: cap}, nil +} + +// AccountTrustEnabled is queried at runtime, including database-policy changes. +func (m *Module) AccountTrustEnabled(ctx context.Context) (bool, error) { + p, err := m.policyForRead(ctx, m.pool) + return p.AccountTrustEnabled, err +} + +// AccountDailyLimit returns nil when the opt-in control is absent. The read +// transaction gives the account response a coherent plan/usage snapshot. +func (m *Module) AccountDailyLimit(ctx context.Context, user string) (*sendramp.AccountDailyLimit, error) { + tx, err := m.pool.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.RepeatableRead, AccessMode: pgx.ReadOnly}) + if err != nil { + return nil, err + } + defer tx.Rollback(ctx) + p, err := m.policyForRead(ctx, tx) + if err != nil { + return nil, err + } + if !p.AccountTrustEnabled { + return nil, nil + } + facts, err := loadAccountAccessFacts(ctx, tx, user) + if err != nil { + return nil, err + } + if accountClassExempt(facts.class) { + return nil, nil + } + cap, err := accountPlanCap(ctx, tx, user) + if err != nil { + return nil, err + } + day, err := ledgerDay(ctx, tx) + if err != nil { + return nil, err + } + d, err := sendramp.AccountSnapshotTx(ctx, tx, user, day, cap) + if err != nil { + return nil, err + } + return &d, nil +} + +// RecheckAccountTrustGrant is a refuse-only read immediately before provider +// I/O. It never takes account locks after the operation lock, avoiding a lock +// inversion with authorization. A stale classification must obtain a new grant. +func (m *Module) recheckAccountTrustGrant(ctx context.Context, tx pgx.Tx, op operationRow, stored reservationRow) (bool, error) { + facts, err := loadAccountAccessFacts(ctx, tx, op.accountRef()) + if err != nil { + return false, err + } + if accountClassExempt(facts.class) { + return true, nil + } + envelope, err := messageEnvelopeQ(ctx, tx, op.OperationID) + if err != nil { + return false, err + } + n, err := externalRecipientCount(ctx, tx, op.accountRef(), envelope) + if err != nil { + return false, err + } + if stored.AccountTrustUnits == nil || n != *stored.AccountTrustUnits { + return false, nil + } + if n == 0 { + return true, nil + } + if !op.Shared { + agent, _, err := messageSender(ctx, tx, op.OperationID) + if err != nil { + return false, err + } + ok, err := customIdentityVerified(ctx, tx, op.accountRef(), agent) + if err != nil || !ok { + return false, err + } + } + cap, err := accountPlanCap(ctx, tx, op.accountRef()) + if err != nil { + return false, err + } + now, err := ledgerDay(ctx, tx) + if err != nil { + return false, err + } + d, err := sendramp.AccountSnapshotTx(ctx, tx, op.accountRef(), now, cap) + if err != nil { + return false, err + } + var valid bool + err = tx.QueryRow(ctx, `SELECT EXISTS(SELECT 1 FROM account_send_reservations WHERE message_id=$1 AND user_id=$2 AND day=$3 AND units >= $4 AND shared=$5 AND state='reserved')`, op.OperationID, op.accountRef(), now, n, op.Shared).Scan(&valid) + return valid && d.Used <= d.Limit && (!op.Shared || d.SharedUsed <= d.SharedLimit), err +} + +// DailyLimitPreflight gives immediate sends a useful refusal before persistence. +// It does not reserve capacity: final authorization repeats the decision under +// account locks, so concurrent preflights cannot overrun the allowance. +func (m *Module) DailyLimitPreflight(ctx context.Context, user, agent string, recipients []string) (*sendramp.AccountDailyLimit, error) { + enabled, err := m.AccountTrustEnabled(ctx) + if err != nil || !enabled { + return nil, err + } + envelope, err := normalizeEnvelope(recipients) + if err != nil { + return nil, err + } + n, err := externalRecipientCount(ctx, m.pool, user, envelope) + if err != nil || n == 0 { + return nil, err + } + facts, err := loadAccountAccessFacts(ctx, m.pool, user) + if err != nil { + return nil, err + } + if accountClassExempt(facts.class) { + return nil, nil + } + d, err := m.AccountDailyLimit(ctx, user) + if err != nil || d == nil { + return d, err + } + own, err := customIdentityVerified(ctx, m.pool, user, agent) + if err != nil { + return nil, err + } + d.SharedBinding = !own && d.SharedLimit-d.SharedUsed < d.Limit-d.Used + d.Allowed = n <= d.Limit-d.Used && (own || n <= d.SharedLimit-d.SharedUsed) + return d, nil +} diff --git a/internal/sendingpolicy/account_trust_integration_test.go b/internal/sendingpolicy/account_trust_integration_test.go new file mode 100644 index 000000000..1a64d0c5f --- /dev/null +++ b/internal/sendingpolicy/account_trust_integration_test.go @@ -0,0 +1,252 @@ +package sendingpolicy_test + +import ( + "github.com/tokencanopy/e2a/internal/sendingpolicy" + "strings" + "sync" + "sync/atomic" + "testing" +) + +func TestAccountTrustAcrossSharedAndOwnIdentities(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + g := f.gate(p) + user := f.user("standard") + f.plan(user, "scale") + shared := f.agent(user) + own, _ := f.customDomainAgent(user) + if d := f.send(g, f.message(shared, "relay", 12)); !d.Allow { + t.Fatalf("shared: %+v", d) + } + if d := f.send(g, f.message(own, "own_address", 8)); !d.Allow { + t.Fatalf("own: %+v", d) + } + if d := f.send(g, f.message(own, "own_address", 1)); d.Allow || d.Reason != sendingpolicy.ReasonRampCapacity { + t.Fatalf("account cap: %+v", d) + } +} + +func TestAccountTrustInternalRecipientsDoNotUseDailyOrPlatformCapacity(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + p.DisableLegacyDailyBudgets = true + p.BudgetMode = sendingpolicy.ModeEnforce + p.AllCustomerGlobalDailyRecipients = 1 + g := f.gate(p) + user := f.user("standard") + f.plan(user, "free") + sender := f.agent(user) + recipient := "inside@agents.localhost" + if _, err := f.pool.Exec(f.ctx, `INSERT INTO domains(domain,user_id) VALUES('agents.localhost',$1) ON CONFLICT DO NOTHING`, user); err != nil { + t.Fatal(err) + } + if _, err := f.pool.Exec(f.ctx, `INSERT INTO agent_identities(id,user_id,registered_domain,name) VALUES($1,$2,'agents.localhost','inside')`, recipient, user); err != nil { + t.Fatal(err) + } + message := f.messageTo(sender, "relay", []string{recipient}) + if d := f.send(g, message); !d.Allow { + t.Fatalf("internal send: %+v", d) + } + if d := f.send(g, f.message(sender, "relay", 1)); !d.Allow { + t.Fatalf("internal consumed platform capacity: %+v", d) + } +} + +func TestAccountTrustConcurrentReservations(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + g := f.gate(p) + user := f.user("standard") + f.plan(user, "scale") + agent := f.agent(user) + var refs []sendingpolicy.OperationRef + for i := 0; i < 25; i++ { + _, ref := f.prepareMessage(g, f.message(agent, "relay", 1)) + refs = append(refs, ref) + } + var allowed atomic.Int32 + var wg sync.WaitGroup + for _, ref := range refs { + wg.Add(1) + go func(ref sendingpolicy.OperationRef) { + defer wg.Done() + d := f.authorize(g, ref) + if d.Allow { + allowed.Add(1) + } else if d.DailyLimit == nil || d.DailyLimit.Limit != 20 { + t.Errorf("missing limit on hold: %+v", d) + } + }(ref) + } + wg.Wait() + if allowed.Load() != 20 { + t.Fatalf("concurrent allowed=%d want 20", allowed.Load()) + } +} + +func TestAccountTrustBareFreeAndSharedCeiling(t *testing.T) { + for _, tc := range []struct { + name string + cap *int + shared bool + units int + }{{"bare free", intPtr(20), false, 21}, {"shared", nil, true, 51}, {"own", nil, false, 2001}} { + t.Run(tc.name, func(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + g := f.gate(p) + user := f.user("standard") + f.plan(user, "scale") + if _, err := f.pool.Exec(f.ctx, `UPDATE account_limits SET max_messages_day=$2 WHERE user_id=$1`, user, tc.cap); err != nil { + t.Fatal(err) + } + if _, err := f.pool.Exec(f.ctx, `INSERT INTO account_sending_trust(user_id,grandfather_daily) VALUES($1,2000)`, user); err != nil { + t.Fatal(err) + } + agent, _ := f.customDomainAgent(user) + sentAs := "own_address" + if tc.shared { + agent = f.agent(user) + sentAs = "relay" + } + if d := f.send(g, f.message(agent, sentAs, tc.units)); d.Allow || d.Reason != sendingpolicy.ReasonRampCapacity { + t.Fatalf("ceiling: %+v", d) + } + }) + } +} +func intPtr(n int) *int { return &n } + +func TestAccountTrustRedemptionRechecksPlanCap(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + g := f.gate(p) + user := f.user("standard") + f.plan(user, "scale") + agent := f.agent(user) + _, ref := f.prepareMessage(g, f.message(agent, "relay", 1)) + _, attempt, err := g.Reserve(f.ctx, ref) + if err != nil { + t.Fatal(err) + } + d, auth, err := g.ConsumeAttempt(f.ctx, attempt) + if err != nil || !d.Allow || auth == nil { + t.Fatalf("authorize: %+v %v", d, err) + } + if _, err = f.pool.Exec(f.ctx, `UPDATE account_limits SET max_messages_day=0 WHERE user_id=$1`, user); err != nil { + t.Fatal(err) + } + if err = g.RedeemProviderCall(f.ctx, *auth); err == nil { + t.Fatal("old grant bypassed reduced plan cap") + } +} + +func TestAccountTrustMixedEnvelopeAndOwnerProofRevocation(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + p.DisableLegacyDailyBudgets = true + p.BudgetMode = sendingpolicy.ModeEnforce + p.AllCustomerGlobalDailyRecipients = 1 + g := f.gate(p) + user := f.user("standard") + f.plan(user, "scale") + sender := f.agent(user) + f.proveOwner(user) + owner := f.ownerEmail(user) + f.exec(`INSERT INTO domains(domain,user_id) VALUES('agents.localhost',$1) ON CONFLICT DO NOTHING`, user) + f.exec(`INSERT INTO agent_identities(id,user_id,registered_domain,name) VALUES('inside@agents.localhost',$1,'agents.localhost','inside')`, user) + message := f.esaMessage(sender, "relay", []string{owner, "unique@example.test"}, []string{strings.ToUpper(owner), "UNIQUE@EXAMPLE.TEST"}, []string{"inside@agents.localhost", "unique@example.test"}) + _, ref := f.prepareMessage(g, message) + _, attempt, err := g.Reserve(f.ctx, ref) + if err != nil { + t.Fatal(err) + } + d, auth, err := g.ConsumeAttempt(f.ctx, attempt) + if err != nil || !d.Allow || auth == nil { + t.Fatalf("mixed envelope: %+v %v", d, err) + } + snapshot, err := g.(*sendingpolicy.Module).AccountDailyLimit(f.ctx, user) + if err != nil || snapshot.Used != 1 || snapshot.SharedUsed != 1 { + t.Fatalf("usage: %+v %v", snapshot, err) + } + // Losing the verified-owner exemption makes the old one-unit token stale. + f.exec(`UPDATE users SET owner_email_verified_at=NULL,owner_email_verified_address=NULL,owner_email_verified_source=NULL WHERE id=$1`, user) + if err = g.RedeemProviderCall(f.ctx, *auth); err == nil { + t.Fatal("revoked owner proof expanded old grant") + } +} + +func TestAccountTrustGrantCannotKeepExternalCreditAfterOwnerVerification(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + g := f.gate(p) + user := f.user("standard") + f.plan(user, "scale") + sender := f.agent(user) + _, ref := f.prepareMessage(g, f.messageTo(sender, "relay", []string{f.ownerEmail(user)})) + _, attempt, err := g.Reserve(f.ctx, ref) + if err != nil { + t.Fatal(err) + } + d, auth, err := g.ConsumeAttempt(f.ctx, attempt) + if err != nil || !d.Allow || auth == nil { + t.Fatalf("grant: %+v %v", d, err) + } + f.proveOwner(user) + if err = g.RedeemProviderCall(f.ctx, *auth); err == nil { + t.Fatal("newly internal recipient kept external-credit token") + } +} + +func TestAccountTrustDisabledRetryDoesNotEarnActivity(t *testing.T) { + f := newFixture(t) + p := sendingpolicy.DisabledPolicy() + p.AccountTrustEnabled = true + g := f.gate(p) + user := f.user("standard") + f.plan(user, "scale") + sender := f.agent(user) + _, ref := f.prepareMessage(g, f.messageTo(sender, "relay", []string{f.ownerEmail(user)})) + _, attempt, err := g.Reserve(f.ctx, ref) + if err != nil { + t.Fatal(err) + } + d, auth, err := g.ConsumeAttempt(f.ctx, attempt) + if err != nil || !d.Allow || auth == nil { + t.Fatalf("grant: %+v %v", d, err) + } + f.proveOwner(user) + disabled := f.gate(sendingpolicy.DisabledPolicy()) + if err = disabled.RedeemProviderCall(f.ctx, *auth); err == nil { + t.Fatal("changed accounting mode kept old grant") + } + _, retry, err := disabled.Reserve(f.ctx, ref) + if err != nil { + t.Fatal(err) + } + d, auth, err = disabled.ConsumeAttempt(f.ctx, retry) + if err != nil || !d.Allow || auth == nil { + t.Fatalf("retry: %+v %v", d, err) + } + if err = disabled.RedeemProviderCall(f.ctx, *auth); err != nil { + t.Fatal(err) + } + if err = disabled.SettleProvider(f.ctx, sendingpolicy.ProviderSettlement{Attempt: retry, Outcome: sendingpolicy.SettlementProviderAccepted}); err != nil { + t.Fatal(err) + } + var activity int + if err = f.pool.QueryRow(f.ctx, `SELECT COALESCE(sum(confirmed_count),0) FROM account_send_days WHERE user_id=$1`, user).Scan(&activity); err != nil { + t.Fatal(err) + } + if activity != 0 { + t.Fatalf("disabled internal retry earned %d units of clean activity", activity) + } +} diff --git a/internal/sendingpolicy/budget.go b/internal/sendingpolicy/budget.go index 32f414d6d..814d84797 100644 --- a/internal/sendingpolicy/budget.go +++ b/internal/sendingpolicy/budget.go @@ -6,6 +6,7 @@ import ( "encoding/hex" "errors" "fmt" + "math" "sort" "time" @@ -130,6 +131,11 @@ func scopeKeys(purpose Purpose, accountID string, shared, probation bool) ([]cou // limitFor resolves today's cap for one counter under the current policy and // the account's authoritative plan. func limitFor(key counterKey, policy RuntimePolicy, planCode string) int { + // Keep ledger keys for safe settlement across policy changes, but remove + // these legacy caps from admission. The account trust ledger owns them. + if policy.DisableLegacyDailyBudgets && (key.Scope == ScopeGlobalProbation || key.Scope == ScopeAccountDaily || key.Scope == ScopeAccountSharedDaily) { + return math.MaxInt32 + } switch key.Scope { case ScopeGlobalAll: return policy.AllCustomerGlobalDailyRecipients @@ -353,6 +359,9 @@ func (p *ledgerPlan) release(ref ledgerRef, units int) { // acquire records taking `units`, reporting whether there was room. func (p *ledgerPlan) acquire(ref ledgerRef, units int) bool { + if units == 0 { + return true + } row, ok := p.rows[ref] if !ok { return false diff --git a/internal/sendingpolicy/fromconfig.go b/internal/sendingpolicy/fromconfig.go index 673cedf6c..f588f01fe 100644 --- a/internal/sendingpolicy/fromconfig.go +++ b/internal/sendingpolicy/fromconfig.go @@ -19,6 +19,8 @@ func FromConfig(cfg *config.Config) (RuntimePolicy, error) { sp := cfg.SendingProtect policy := RuntimePolicy{ + DisableLegacyDailyBudgets: cfg.SendingRamp.DisableLegacyDailyBudgets, + AccountTrustEnabled: cfg.SendingRamp.AccountTrustEnabled, AllCustomerGlobalDailyRecipients: sp.AllCustomerGlobalDailyRecipients, BounceMinOutcomes: sp.BounceMinOutcomes, BouncePauseBasisPoints: sp.BouncePauseBasisPoints, diff --git a/internal/sendingpolicy/gate.go b/internal/sendingpolicy/gate.go index c4706fe51..01a57eb08 100644 --- a/internal/sendingpolicy/gate.go +++ b/internal/sendingpolicy/gate.go @@ -89,20 +89,21 @@ func (m *Module) effectivePolicy(ctx context.Context, tx pgx.Tx) (RuntimePolicy, // reservationRow is one row of sending_budget_reservations. type reservationRow struct { - OperationID string - Attempt int - SourceAccountRef *string - PolicySubjectRef string - Purpose Purpose - Day time.Time - Units int - Probation bool - State string - CallState string - Nonce *string - NoticeVersion *int - NoticeCommitment []byte - Exists bool + OperationID string + Attempt int + SourceAccountRef *string + PolicySubjectRef string + Purpose Purpose + Day time.Time + Units int + AccountTrustUnits *int + Probation bool + State string + CallState string + Nonce *string + NoticeVersion *int + NoticeCommitment []byte + Exists bool } // scopeKeys returns the counters this stored reservation charged, using only @@ -127,13 +128,13 @@ func lockReservation(ctx context.Context, tx pgx.Tx, operationID string, attempt err := tx.QueryRow(ctx, ` SELECT operation_id, submission_attempt, source_account_ref, policy_subject_ref, purpose, day, units, probation, state, call_state, authorization_nonce, - notice_recipient_version, notice_recipient_commitment + notice_recipient_version, notice_recipient_commitment, account_trust_units FROM sending_budget_reservations WHERE operation_id = $1 AND submission_attempt = $2 FOR UPDATE`, operationID, attempt, ).Scan(&r.OperationID, &r.Attempt, &r.SourceAccountRef, &r.PolicySubjectRef, &r.Purpose, &r.Day, &r.Units, &r.Probation, &r.State, &r.CallState, &r.Nonce, - &r.NoticeVersion, &r.NoticeCommitment) + &r.NoticeVersion, &r.NoticeCommitment, &r.AccountTrustUnits) if errors.Is(err, pgx.ErrNoRows) { return reservationRow{}, nil } @@ -299,6 +300,16 @@ func (m *Module) Reserve(ctx context.Context, ref OperationRef) (Decision, Attem if err != nil { return Decision{}, AttemptRef{}, err } + if policy.AccountTrustEnabled && op.Purpose == PurposeCustomerMessage { + envelope, e := messageEnvelopeQ(ctx, tx, op.OperationID) + if e != nil { + return Decision{}, AttemptRef{}, e + } + units, err = externalRecipientCount(ctx, tx, op.accountRef(), envelope) + if err != nil { + return Decision{}, AttemptRef{}, err + } + } day, err := ledgerDay(ctx, tx) if err != nil { @@ -707,6 +718,9 @@ func (m *Module) readAuthState(ctx context.Context, tx pgx.Tx, ref AttemptRef) ( if errors.Is(err, ErrSourceUnavailable) { return st, terminalHold(ReasonSourceUnavailable), nil } + if d, ok := rampHoldFor(err, st.day); ok { + return st, d, nil + } return st, Decision{}, err } @@ -798,6 +812,12 @@ func (m *Module) readAuthState(ctx context.Context, tx pgx.Tx, ref AttemptRef) ( } st.envelope = envelope st.units = len(st.envelope) + if policy.AccountTrustEnabled && op.Purpose == PurposeCustomerMessage { + st.units, err = externalRecipientCount(ctx, tx, op.accountRef(), st.envelope) + if err != nil { + return st, Decision{}, err + } + } // External sending access, decided under the locks just taken: the users // row (FOR SHARE), the account control row and the plan row are all held @@ -826,6 +846,9 @@ func (m *Module) readAuthState(ctx context.Context, tx pgx.Tx, ref AttemptRef) ( if errors.Is(err, ErrSourceUnavailable) { return st, terminalHold(ReasonSourceUnavailable), nil } + if d, ok := rampHoldFor(err, st.day); ok { + return st, d, nil + } return st, Decision{}, err } return st, allowDecision(), nil @@ -1208,6 +1231,16 @@ func (m *Module) releaseReacquiredUnits(ctx context.Context, tx pgx.Tx, st authS return nil } +// accountTrustUnits records the meaning of attempt units at authorization. +// Null marks legacy/all-recipient accounting, which cannot earn trust credit. +func accountTrustUnits(st authState) *int { + if st.policy.AccountTrustEnabled && st.op.Purpose == PurposeCustomerMessage { + n := st.units + return &n + } + return nil +} + // authorize confirms the capacity, mints the single-use nonce, records the // provenance correlation, and returns the token. func (m *Module) authorize(ctx context.Context, tx pgx.Tx, st authState) (*ProviderAuthorization, error) { @@ -1243,8 +1276,8 @@ func (m *Module) authorize(ctx context.Context, tx pgx.Tx, st authState) (*Provi INSERT INTO sending_budget_reservations (operation_id, submission_attempt, source_account_ref, policy_subject_ref, purpose, day, units, probation, state, call_state, authorization_nonce, - notice_recipient_version, notice_recipient_commitment) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, 'confirmed', 'authorized', $9, $10, $11) + notice_recipient_version, notice_recipient_commitment, account_trust_units) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, 'confirmed', 'authorized', $9, $10, $11, $12) ON CONFLICT (operation_id, submission_attempt) DO UPDATE SET day = EXCLUDED.day, units = EXCLUDED.units, @@ -1254,11 +1287,12 @@ func (m *Module) authorize(ctx context.Context, tx pgx.Tx, st authState) (*Provi authorization_nonce = EXCLUDED.authorization_nonce, notice_recipient_version = EXCLUDED.notice_recipient_version, notice_recipient_commitment = EXCLUDED.notice_recipient_commitment, + account_trust_units = EXCLUDED.account_trust_units, provider_call_started_at = NULL, updated_at = now()`, st.op.OperationID, attempt, st.op.SourceAccountRef, st.op.PolicySubjectRef, st.op.Purpose, st.day, st.units, st.probation, nonce, - noticeVersion, noticeCommitment, + noticeVersion, noticeCommitment, accountTrustUnits(st), ); err != nil { return nil, fmt.Errorf("sendingpolicy: confirm reservation: %w", err) } @@ -1633,6 +1667,19 @@ func (m *Module) RedeemProviderCall(ctx context.Context, auth ProviderAuthorizat return m.invalidate(ctx, tx, auth.attempt) } + if op.Purpose == PurposeCustomerMessage && policy.AccountTrustEnabled != (stored.AccountTrustUnits != nil) { + return m.invalidate(ctx, tx, auth.attempt) + } + if policy.AccountTrustEnabled && op.Purpose == PurposeCustomerMessage { + valid, err := m.recheckAccountTrustGrant(ctx, tx, op, stored) + if err != nil { + return err + } + if !valid { + return m.invalidate(ctx, tx, auth.attempt) + } + } + if auth.notice != nil { ok, err := m.noticeSelectorStillCurrent(ctx, tx, auth, stored, ownerEmail) if err != nil { @@ -2006,7 +2053,11 @@ func (m *Module) settle(ctx context.Context, operationID string, attempt int, se // Ramp keys come last in the normative order, after the correlation row, // which is keyed by this operation and already held under its lock. if op.Purpose == PurposeCustomerMessage { - if err := m.rampSettle(ctx, tx, op.OperationID, settlement.Outcome); err != nil { + acceptedUnits := 0 + if stored.AccountTrustUnits != nil { + acceptedUnits = *stored.AccountTrustUnits + } + if err := m.rampSettle(ctx, tx, op.OperationID, settlement.Outcome, stored.Day, acceptedUnits); err != nil { return err } } diff --git a/internal/sendingpolicy/policy.go b/internal/sendingpolicy/policy.go index e1dd981b6..7d41bf933 100644 --- a/internal/sendingpolicy/policy.go +++ b/internal/sendingpolicy/policy.go @@ -115,6 +115,9 @@ const maxBasisPoints = 9999 // struct is hashed, reviewed by a human, and then required by hash at // activation, so renaming a key is a policy-breaking change. type RuntimePolicy struct { + DisableLegacyDailyBudgets bool `json:"disable_legacy_daily_budgets,omitempty"` + // Omitted when disabled to preserve existing canonical policy hashes. + AccountTrustEnabled bool `json:"account_trust_enabled,omitempty"` AllCustomerGlobalDailyRecipients int `json:"all_customer_global_daily_recipients"` BounceMinOutcomes int `json:"bounce_min_outcomes"` BouncePauseBasisPoints int `json:"bounce_pause_basis_points"` diff --git a/internal/sendingpolicy/ramp.go b/internal/sendingpolicy/ramp.go index 9c5ca9546..0133a3e71 100644 --- a/internal/sendingpolicy/ramp.go +++ b/internal/sendingpolicy/ramp.go @@ -55,6 +55,9 @@ const ReasonRampCapacity = "sending_ramp_capacity_exhausted" // rampSubject is everything the ramp needs about one customer message, read // from the locked source rows. type rampSubject struct { + account bool + shared bool + planCap *int messageID string userID string domain string @@ -72,6 +75,9 @@ type rampSubject struct { // shared with that worker during the migration. Two different notions of // "eligible" would let one path reserve capacity the other never released. func (m *Module) rampSubjectFor(ctx context.Context, tx pgx.Tx, policy RuntimePolicy, op operationRow, units int) (rampSubject, error) { + if policy.AccountTrustEnabled { + return m.accountTrustSubject(ctx, tx, op) + } if !policy.RampEnabled || op.Purpose != PurposeCustomerMessage || op.Shared { return rampSubject{}, nil } @@ -115,6 +121,9 @@ func (m *Module) rampSubjectFor(ctx context.Context, tx pgx.Tx, policy RuntimePo // probation concept for custom domains at all, and the account and platform // pools still bound the traffic. func (m *Module) rampProbation(ctx context.Context, tx pgx.Tx, policy RuntimePolicy, op operationRow) (bool, error) { + if policy.AccountTrustEnabled { + return false, nil + } if op.Shared { return true, nil } @@ -145,6 +154,20 @@ func (m *Module) rampAuthorize(ctx context.Context, tx pgx.Tx, policy RuntimePol if !subject.applies { return nil } + if subject.account { + d, err := sendramp.ReserveAccountTx(ctx, tx, sendramp.AccountReserveRequest{UserID: subject.userID, MessageID: subject.messageID, Units: subject.units, Shared: subject.shared, Day: day, PlanCap: subject.planCap}) + if err != nil { + var permanent *sendramp.PermanentError + if errors.As(err, &permanent) { + return fmt.Errorf("%w: %v", errRampUnavailable, err) + } + return err + } + if !d.Allowed { + return &accountCapacityError{daily: d} + } + return nil + } decision, err := sendramp.ReserveTx(ctx, tx, sendramp.ReserveRequest{ MessageID: subject.messageID, UserID: subject.userID, @@ -188,6 +211,12 @@ func (m *Module) rampAuthorize(ctx context.Context, tx pgx.Tx, policy RuntimePol // writing one: the two call sites have released their units at different // points, and only they know which. func rampHoldFor(err error, day time.Time) (Decision, bool) { + var capacity *accountCapacityError + if errors.As(err, &capacity) { + d := holdDecision(ReasonRampCapacity, capacity.daily.ResetsAt) + d.DailyLimit = &capacity.daily + return d, true + } switch { case errors.Is(err, errRampCapacity): return holdDecision(ReasonRampCapacity, nextUTCMidnight(day)), true @@ -207,6 +236,9 @@ func rampHoldFor(err error, day time.Time) (Decision, bool) { // merely slowed down, and releasing its ramp claim would let the same message // re-qualify a stage it has already qualified. func (m *Module) rampRelease(ctx context.Context, tx pgx.Tx, messageID string) error { + if err := sendramp.SettleAccountTx(ctx, tx, messageID, false, time.Time{}, 0); err != nil { + return err + } if err := sendramp.ReleaseTx(ctx, tx, messageID); err != nil { return fmt.Errorf("sendingpolicy: release ramp capacity: %w", err) } @@ -221,13 +253,19 @@ func (m *Module) rampRelease(ctx context.Context, tx pgx.Tx, messageID string) e // gives the units back. Retryable and ambiguous results are deliberately absent // from the closed outcome set and leave the reservation standing: a message // that might have been delivered must not release ramp capacity. -func (m *Module) rampSettle(ctx context.Context, tx pgx.Tx, messageID string, outcome SettlementOutcome) error { +func (m *Module) rampSettle(ctx context.Context, tx pgx.Tx, messageID string, outcome SettlementOutcome, day time.Time, units int) error { switch outcome { case SettlementProviderAccepted: + if err := sendramp.SettleAccountTx(ctx, tx, messageID, true, day, units); err != nil { + return err + } if err := sendramp.ConfirmTx(ctx, tx, messageID); err != nil { return fmt.Errorf("sendingpolicy: confirm ramp capacity: %w", err) } case SettlementProviderPermanentlyRejected: + if err := sendramp.SettleAccountTx(ctx, tx, messageID, false, time.Time{}, 0); err != nil { + return err + } if err := sendramp.ReleaseTx(ctx, tx, messageID); err != nil { return fmt.Errorf("sendingpolicy: release ramp capacity: %w", err) } diff --git a/internal/sendingpolicy/types.go b/internal/sendingpolicy/types.go index 80c852120..51c3d44e1 100644 --- a/internal/sendingpolicy/types.go +++ b/internal/sendingpolicy/types.go @@ -4,6 +4,7 @@ import ( "encoding/json" "errors" "fmt" + "github.com/tokencanopy/e2a/internal/sendramp" "sort" "strings" "time" @@ -217,10 +218,11 @@ const ( // forever instead of failing the message once. A terminal hold carries no // RetryAt because there is no time at which the answer changes. type Decision struct { - Allow bool - Reason string - RetryAt time.Time - Terminal bool + DailyLimit *sendramp.AccountDailyLimit + Allow bool + Reason string + RetryAt time.Time + Terminal bool } func allowDecision() Decision { return Decision{Allow: true} } diff --git a/internal/sendramp/account.go b/internal/sendramp/account.go new file mode 100644 index 000000000..083bc7922 --- /dev/null +++ b/internal/sendramp/account.go @@ -0,0 +1,221 @@ +package sendramp + +import ( + "context" + "errors" + "time" + + "github.com/jackc/pgx/v5" +) + +// AccountDailyLimit reports external-recipient capacity. Shared identity usage +// is a subset of account usage, with a ceiling of 50. Reserved uncertain sends +// remain charged until an authoritative outcome resolves them. +type AccountDailyLimit struct { + Limit int `json:"limit"` + Used int `json:"used"` + SharedLimit int `json:"shared_limit"` + SharedUsed int `json:"shared_used"` + CleanActiveDays int `json:"clean_active_days"` + ResetsAt time.Time `json:"resets_at"` + Allowed bool `json:"-"` + SharedBinding bool `json:"-"` +} + +// AccountTrustLimit reaches 2,000 on the thirtieth clean active day. Only +// completed UTC days earn capacity; neither elapsed time nor retries do. +func AccountTrustLimit(days int) int { + if days < 0 { + days = 0 + } + if days > 29 { + days = 29 + } + return 20 + 1980*days/29 +} + +type AccountReserveRequest struct { + UserID, MessageID string + Units int + Shared bool + Day time.Time + PlanCap *int +} + +func lockAccount(ctx context.Context, tx pgx.Tx, user string) error { + if _, err := tx.Exec(ctx, `INSERT INTO account_sending_trust(user_id) VALUES($1) ON CONFLICT DO NOTHING`, user); err != nil { + return err + } + var id string + return tx.QueryRow(ctx, `SELECT user_id FROM account_sending_trust WHERE user_id=$1 FOR UPDATE`, user).Scan(&id) +} + +// AccountSnapshotTx is read-only. Callers reserving capacity first lock the +// account trust row; UI reads need no locks. A plan cap of nil is unlimited. +func AccountSnapshotTx(ctx context.Context, tx pgx.Tx, user string, now time.Time, planCap *int) (AccountDailyLimit, error) { + day := utcDay(now) + d := AccountDailyLimit{ResetsAt: day.AddDate(0, 0, 1)} + var floor int + err := tx.QueryRow(ctx, `SELECT COALESCE((SELECT grandfather_daily FROM account_sending_trust WHERE user_id=$1),0), + COALESCE((SELECT archived_clean_days FROM account_sending_trust WHERE user_id=$1),0) + (SELECT count(*) FROM account_send_days WHERE user_id=$1 AND day<$2 AND confirmed_count>0 AND NOT breached), + COALESCE((SELECT reserved_count FROM account_send_days WHERE user_id=$1 AND day=$2),0), + COALESCE((SELECT shared_count FROM account_send_days WHERE user_id=$1 AND day=$2),0)`, user, day).Scan(&floor, &d.CleanActiveDays, &d.Used, &d.SharedUsed) + if err != nil { + return d, err + } + d.Limit = max(floor, AccountTrustLimit(d.CleanActiveDays)) + if planCap != nil { + d.Limit = min(d.Limit, max(0, *planCap)) + } + d.SharedLimit = min(50, d.Limit) + return d, nil +} + +// ReserveAccountTx serializes all identities of one account. Call it after +// account-control and platform-budget locks, before issuing a provider grant. +func ReserveAccountTx(ctx context.Context, tx pgx.Tx, r AccountReserveRequest) (AccountDailyLimit, error) { + if r.UserID == "" || r.MessageID == "" || r.Units < 1 { + return AccountDailyLimit{}, permanentf("sendramp: invalid account reservation") + } + if err := lockAccount(ctx, tx, r.UserID); err != nil { + return AccountDailyLimit{}, err + } + day := utcDay(r.Day) + d, err := AccountSnapshotTx(ctx, tx, r.UserID, day, r.PlanCap) + if err != nil { + return d, err + } + var owner, state string + var oldDay time.Time + var units int + var shared bool + err = tx.QueryRow(ctx, `SELECT user_id,day,units,shared,state FROM account_send_reservations WHERE message_id=$1 FOR UPDATE`, r.MessageID).Scan(&owner, &oldDay, &units, &shared, &state) + switch { + case err == nil: + if owner != r.UserID || shared != r.Shared { + return d, permanentf("sendramp: account reservation changed") + } + if state == "released" { + return d, permanentf("sendramp: account reservation released") + } + if state == "confirmed" { + d.Allowed = true + return d, nil + } + // Classification may change when the owner mailbox or live-agent set + // changes. Acquire newly external units, but retain uncertain exposure + // already reserved by a prior attempt rather than refunding it here. + r.Units = max(r.Units, units) + if oldDay.Equal(day) { + extra := r.Units - units + if extra > d.Limit-d.Used || (r.Shared && extra > d.SharedLimit-d.SharedUsed) { + return d, nil + } + if extra > 0 { + if _, err = tx.Exec(ctx, `UPDATE account_send_days SET reserved_count=reserved_count+$3,shared_count=shared_count+$4 WHERE user_id=$1 AND day=$2`, r.UserID, day, extra, sharedUnits(extra, r.Shared)); err != nil { + return d, err + } + if _, err = tx.Exec(ctx, `UPDATE account_send_reservations SET units=$2 WHERE message_id=$1`, r.MessageID, r.Units); err != nil { + return d, err + } + } + d.Used += extra + d.SharedUsed += sharedUnits(extra, r.Shared) + d.Allowed = true + return d, nil + } + if r.Units > d.Limit-d.Used || (r.Shared && r.Units > d.SharedLimit-d.SharedUsed) { + return d, nil + } + // Retries on a later day move the unresolved reservation. Both day buckets + // remain serialized by the account lock, including late settlement. + if _, err = tx.Exec(ctx, `UPDATE account_send_days SET reserved_count=reserved_count-$3, shared_count=shared_count-$4 WHERE user_id=$1 AND day=$2`, r.UserID, oldDay, units, sharedUnits(units, shared)); err != nil { + return d, err + } + if _, err = tx.Exec(ctx, `DELETE FROM account_send_reservations WHERE message_id=$1`, r.MessageID); err != nil { + return d, err + } + case !errors.Is(err, pgx.ErrNoRows): + return d, err + } + if r.Units > d.Limit-d.Used || (r.Shared && r.Units > d.SharedLimit-d.SharedUsed) { + return d, nil + } + _, err = tx.Exec(ctx, `INSERT INTO account_send_days(user_id,day,reserved_count,shared_count) VALUES($1,$2,$3,$4) + ON CONFLICT(user_id,day) DO UPDATE SET reserved_count=account_send_days.reserved_count+EXCLUDED.reserved_count,shared_count=account_send_days.shared_count+EXCLUDED.shared_count`, r.UserID, day, r.Units, sharedUnits(r.Units, r.Shared)) + if err != nil { + return d, err + } + _, err = tx.Exec(ctx, `INSERT INTO account_send_reservations(message_id,user_id,day,units,shared,state) VALUES($1,$2,$3,$4,$5,'reserved')`, r.MessageID, r.UserID, day, r.Units, r.Shared) + if err != nil { + return d, err + } + d.Used += r.Units + d.SharedUsed += sharedUnits(r.Units, r.Shared) + d.Allowed = true + return d, nil +} + +func sharedUnits(units int, shared bool) int { + if shared { + return units + } + return 0 +} + +// SettleAccountTx preserves idempotency and authoritative late acceptance. +// It probes immutable ownership before taking locks in account->message order. +func SettleAccountTx(ctx context.Context, tx pgx.Tx, message string, accepted bool, acceptedDay time.Time, acceptedUnits int) error { + var user string + err := tx.QueryRow(ctx, `SELECT user_id FROM account_send_reservations WHERE message_id=$1`, message).Scan(&user) + if errors.Is(err, pgx.ErrNoRows) { + return nil + } + if err != nil { + return err + } + if err = lockAccount(ctx, tx, user); err != nil { + return err + } + var state string + var day time.Time + var units int + var shared bool + err = tx.QueryRow(ctx, `SELECT day,units,shared,state FROM account_send_reservations WHERE message_id=$1 FOR UPDATE`, message).Scan(&day, &units, &shared, &state) + if errors.Is(err, pgx.ErrNoRows) { + return nil + } + if err != nil { + return err + } + if state == "confirmed" || (!accepted && state == "released") { + return nil + } + next := "released" + reserved := -units + if accepted { + next = "confirmed" + reserved = 0 + if state == "released" { + reserved = units + } + } + _, err = tx.Exec(ctx, `UPDATE account_send_days SET reserved_count=reserved_count+$3,shared_count=shared_count+$4 WHERE user_id=$1 AND day=$2`, user, day, reserved, sharedUnits(reserved, shared)) + if err != nil { + return err + } + // Capacity follows the conservative message reservation, but activity follows + // the immutable accepted attempt. Internal-only retries earn no trust. Update + // only retained day buckets: recreating pruned history could double-credit a + // day already folded into archived_clean_days. + if accepted && acceptedUnits > 0 { + if acceptedDay.IsZero() { + return errors.New("sendramp: accepted attempt day is required") + } + if _, err = tx.Exec(ctx, `UPDATE account_send_days SET confirmed_count=confirmed_count+$3 WHERE user_id=$1 AND day=$2`, user, utcDay(acceptedDay), acceptedUnits); err != nil { + return err + } + } + _, err = tx.Exec(ctx, `UPDATE account_send_reservations SET state=$2 WHERE message_id=$1`, message, next) + return err +} diff --git a/internal/sendramp/account_maintenance.go b/internal/sendramp/account_maintenance.go new file mode 100644 index 000000000..5bef9db7e --- /dev/null +++ b/internal/sendramp/account_maintenance.go @@ -0,0 +1,54 @@ +package sendramp + +import ( + "context" + "time" +) + +// sweepAccounts bounds historical daily buckets without forgetting earned +// trust. Unresolved reservations are never pruned. Ninety days exceeds the +// detector lookback; older clean-day evidence is folded into the account row. +// Each run is bounded to 100 accounts and 1,000 reservations/day buckets per account and skips accounts currently sending. +func (s *Store) sweepAccounts(ctx context.Context, now time.Time) error { + tx, err := s.pool.Begin(ctx) + if err != nil { + return err + } + defer tx.Rollback(ctx) + cutoff := utcDay(now).AddDate(0, 0, -90) + rows, err := tx.Query(ctx, `SELECT a.user_id FROM account_sending_trust a WHERE EXISTS ( + SELECT 1 FROM account_send_days d WHERE d.user_id=a.user_id AND d.day<$1 + AND NOT EXISTS(SELECT 1 FROM account_send_reservations r WHERE r.user_id=d.user_id AND r.day=d.day AND r.state='reserved')) + ORDER BY a.user_id LIMIT 100 FOR UPDATE OF a SKIP LOCKED`, cutoff) + if err != nil { + return err + } + var users []string + for rows.Next() { + var user string + if err = rows.Scan(&user); err != nil { + rows.Close() + return err + } + users = append(users, user) + } + err = rows.Err() + rows.Close() + if err != nil { + return err + } + for _, user := range users { + if _, err = tx.Exec(ctx, `DELETE FROM account_send_reservations WHERE message_id IN (SELECT message_id FROM account_send_reservations WHERE user_id=$1 AND day<$2 AND state IN ('confirmed','released') ORDER BY day,message_id LIMIT 1000)`, user, cutoff); err != nil { + return err + } + _, err = tx.Exec(ctx, `WITH pruned AS ( + DELETE FROM account_send_days WHERE user_id=$1 AND day IN (SELECT d.day FROM account_send_days d WHERE d.user_id=$1 AND d.day<$2 + AND NOT EXISTS(SELECT 1 FROM account_send_reservations r WHERE r.user_id=d.user_id AND r.day=d.day) ORDER BY d.day LIMIT 1000) + RETURNING confirmed_count,breached) + UPDATE account_sending_trust SET archived_clean_days=archived_clean_days+(SELECT count(*) FROM pruned WHERE confirmed_count>0 AND NOT breached) WHERE user_id=$1`, user, cutoff) + if err != nil { + return err + } + } + return tx.Commit(ctx) +} diff --git a/internal/sendramp/account_test.go b/internal/sendramp/account_test.go new file mode 100644 index 000000000..df0c2c1c3 --- /dev/null +++ b/internal/sendramp/account_test.go @@ -0,0 +1,275 @@ +package sendramp_test + +import ( + "context" + "fmt" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/tokencanopy/e2a/internal/sendramp" +) + +func TestAccountTrustSchedule(t *testing.T) { + for _, tc := range []struct{ days, want int }{{0, 20}, {1, 88}, {14, 975}, {29, 2000}, {90, 2000}} { + if got := sendramp.AccountTrustLimit(tc.days); got != tc.want { + t.Errorf("days %d: got %d want %d", tc.days, got, tc.want) + } + } +} + +func TestAccountTrustReservationAndCleanDays(t *testing.T) { + _, pool, user, _, message := seedRampMessage(t, "account-trust") + ctx := context.Background() + day := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + reserve := func(units int, at time.Time) sendramp.AccountDailyLimit { + t.Helper() + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + d, err := sendramp.ReserveAccountTx(ctx, tx, sendramp.AccountReserveRequest{UserID: user, MessageID: message, Units: units, Day: at}) + if err != nil { + t.Fatal(err) + } + if err = tx.Commit(ctx); err != nil { + t.Fatal(err) + } + return d + } + d := reserve(21, day) + if d.Allowed || d.Limit != 20 || d.Used != 0 { + t.Fatalf("initial over-limit: %+v", d) + } + d = reserve(20, day) + if !d.Allowed || d.Used != 20 { + t.Fatalf("reserve: %+v", d) + } + d = reserve(20, day) + if !d.Allowed || d.Used != 20 { + t.Fatalf("retry double charged: %+v", d) + } + for i := 0; i < 2; i++ { + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + if err = sendramp.SettleAccountTx(ctx, tx, message, true, day, 20); err != nil { + t.Fatal(err) + } + if err = tx.Commit(ctx); err != nil { + t.Fatal(err) + } + } + check := func(at time.Time, limit, days int) { + t.Helper() + tx, err := pool.BeginTx(ctx, pgx.TxOptions{}) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + d, err := sendramp.AccountSnapshotTx(ctx, tx, user, at, nil) + if err != nil { + t.Fatal(err) + } + if d.Limit != limit || d.CleanActiveDays != days { + t.Fatalf("snapshot: %+v", d) + } + } + check(day, 20, 0) // Today cannot earn additional capacity during the same day. + check(day.AddDate(0, 0, 1), 88, 1) + check(day.AddDate(0, 0, 30), 88, 1) // Idle time earns no trust. +} + +func TestAccountTrustRetentionKeepsEarnedDaysAndUnresolvedReservations(t *testing.T) { + store, pool, user, _, message := seedRampMessage(t, "account-retention") + ctx := context.Background() + now := time.Date(2026, 9, 30, 12, 0, 0, 0, time.UTC) + old := now.AddDate(0, 0, -100) + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + if _, err = sendramp.ReserveAccountTx(ctx, tx, sendramp.AccountReserveRequest{UserID: user, MessageID: message, Units: 1, Day: old}); err != nil { + t.Fatal(err) + } + if err = sendramp.SettleAccountTx(ctx, tx, message, true, old, 1); err != nil { + t.Fatal(err) + } + if err = tx.Commit(ctx); err != nil { + t.Fatal(err) + } + if err = store.Sweep(ctx, now); err != nil { + t.Fatal(err) + } + var rows int + if err = pool.QueryRow(ctx, `SELECT count(*) FROM account_send_days WHERE user_id=$1`, user).Scan(&rows); err != nil { + t.Fatal(err) + } + if rows != 0 { + t.Fatalf("old account buckets retained: %d", rows) + } + tx, err = pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + d, err := sendramp.AccountSnapshotTx(ctx, tx, user, now, nil) + if err != nil { + t.Fatal(err) + } + if d.Limit != 88 || d.CleanActiveDays != 1 { + t.Fatalf("pruning lost trust: %+v", d) + } +} + +func TestAccountTrustRetentionPreservesUnresolvedAndBreachedDays(t *testing.T) { + store, pool, user, domain, message := seedRampMessage(t, "account-retain-open") + ctx := context.Background() + now := time.Date(2026, 9, 30, 12, 0, 0, 0, time.UTC) + old := now.AddDate(0, 0, -100) + second := createMessageForAgent(t, pool, "agent@"+domain, "account-retain-breach") + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + for i, msg := range []string{message, second} { + if _, err = sendramp.ReserveAccountTx(ctx, tx, sendramp.AccountReserveRequest{UserID: user, MessageID: msg, Units: 1, Day: old.AddDate(0, 0, i)}); err != nil { + t.Fatal(err) + } + } + if err = sendramp.SettleAccountTx(ctx, tx, second, true, old.AddDate(0, 0, 1), 1); err != nil { + t.Fatal(err) + } + if _, err = tx.Exec(ctx, `UPDATE account_send_days SET breached=true WHERE user_id=$1`, user); err != nil { + t.Fatal(err) + } + if err = tx.Commit(ctx); err != nil { + t.Fatal(err) + } + for i := 0; i < 2; i++ { + if err = store.Sweep(ctx, now); err != nil { + t.Fatal(err) + } + } + var reserved int + if err = pool.QueryRow(ctx, `SELECT count(*) FROM account_send_reservations WHERE user_id=$1 AND state='reserved'`, user).Scan(&reserved); err != nil { + t.Fatal(err) + } + if reserved != 1 { + t.Fatal("lost unresolved reservation") + } + tx, err = pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + d, err := sendramp.AccountSnapshotTx(ctx, tx, user, now, nil) + if err != nil { + t.Fatal(err) + } + if d.Limit != 20 || d.CleanActiveDays != 0 { + t.Fatalf("breached day earned trust: %+v", d) + } +} + +func TestAccountTrustMidnightRefusalRetainsLateEvidence(t *testing.T) { + _, pool, user, _, message := seedRampMessage(t, "account-midnight") + ctx := context.Background() + day := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + req := sendramp.AccountReserveRequest{UserID: user, MessageID: message, Units: 1, Day: day, Shared: true} + if d, err := sendramp.ReserveAccountTx(ctx, tx, req); err != nil || !d.Allowed { + t.Fatalf("reserve: %+v %v", d, err) + } + zero := 0 + req.PlanCap = &zero + req.Day = day.AddDate(0, 0, 1) + if d, err := sendramp.ReserveAccountTx(ctx, tx, req); err != nil || d.Allowed { + t.Fatalf("midnight cap: %+v %v", d, err) + } + // A permanent rejection can be corrected by authoritative acceptance; + // neither a midnight hold nor duplicate outcomes may lose/double credit. + for _, accepted := range []bool{false, false, true, true, false} { + if err = sendramp.SettleAccountTx(ctx, tx, message, accepted, day, 1); err != nil { + t.Fatal(err) + } + } + var reserved, confirmed int + if err = tx.QueryRow(ctx, `SELECT reserved_count,confirmed_count FROM account_send_days WHERE user_id=$1 AND day=$2`, user, day).Scan(&reserved, &confirmed); err != nil { + t.Fatal(err) + } + if reserved != 1 || confirmed != 1 { + t.Fatalf("late evidence: reserved=%d confirmed=%d", reserved, confirmed) + } + d, err := sendramp.AccountSnapshotTx(ctx, tx, user, req.Day, nil) + if err != nil { + t.Fatal(err) + } + if d.CleanActiveDays != 1 || d.Used != 0 { + t.Fatalf("day attribution: %+v", d) + } +} + +func TestAccountTrustReclassifiedRecipientsAcquireAdditionalCapacity(t *testing.T) { + _, pool, user, _, message := seedRampMessage(t, "account-reclassified") + ctx := context.Background() + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + r := sendramp.AccountReserveRequest{UserID: user, MessageID: message, Units: 1, Day: time.Now().UTC()} + if d, err := sendramp.ReserveAccountTx(ctx, tx, r); err != nil || !d.Allowed { + t.Fatalf("initial: %+v %v", d, err) + } + r.Units = 2 + if d, err := sendramp.ReserveAccountTx(ctx, tx, r); err != nil || !d.Allowed || d.Used != 2 { + t.Fatalf("reclassification must acquire difference: %+v %v", d, err) + } + r.Units = 21 + if d, err := sendramp.ReserveAccountTx(ctx, tx, r); err != nil || d.Allowed || d.Used != 2 { + t.Fatalf("refused resize must preserve reservation: %+v %v", d, err) + } +} + +func TestAccountTrustAcceptedAttemptActivity(t *testing.T) { + for _, units := range []int{0, 1} { + t.Run(fmt.Sprintf("accepted-units-%d", units), func(t *testing.T) { + _, pool, user, _, message := seedRampMessage(t, "account-attempt") + ctx := context.Background() + day := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer tx.Rollback(ctx) + for _, at := range []time.Time{day, day.AddDate(0, 0, 1)} { + d, err := sendramp.ReserveAccountTx(ctx, tx, sendramp.AccountReserveRequest{UserID: user, MessageID: message, Units: 1, Day: at}) + if err != nil || !d.Allowed { + t.Fatalf("reserve: %+v %v", d, err) + } + } + if err = sendramp.SettleAccountTx(ctx, tx, message, true, day, units); err != nil { + t.Fatal(err) + } + var original, retry int + if err = tx.QueryRow(ctx, `SELECT confirmed_count FROM account_send_days WHERE user_id=$1 AND day=$2`, user, day).Scan(&original); err != nil { + t.Fatal(err) + } + if err = tx.QueryRow(ctx, `SELECT confirmed_count FROM account_send_days WHERE user_id=$1 AND day=$2`, user, day.AddDate(0, 0, 1)).Scan(&retry); err != nil { + t.Fatal(err) + } + if original != units || retry != 0 { + t.Fatalf("activity original=%d retry=%d; want %d,0", original, retry, units) + } + }) + } +} diff --git a/internal/sendramp/maintenance.go b/internal/sendramp/maintenance.go index 1da0ff0eb..038b608a8 100644 --- a/internal/sendramp/maintenance.go +++ b/internal/sendramp/maintenance.go @@ -28,7 +28,10 @@ func (s *Store) Sweep(ctx context.Context, now time.Time) error { )`, utcDay(now).AddDate(0, 0, -35)); err != nil { return err } - return tx.Commit(ctx) + if err := tx.Commit(ctx); err != nil { + return err + } + return s.sweepAccounts(ctx, now) } type MaintenanceArgs struct{} diff --git a/internal/testutil/contract_server.go b/internal/testutil/contract_server.go index 1a84f5cf2..ba1acf6fa 100644 --- a/internal/testutil/contract_server.go +++ b/internal/testutil/contract_server.go @@ -191,6 +191,7 @@ func StartContractServer(ctx context.Context, dbURL string) (*ContractServer, er // so every other scenario is unaffected while the restriction's contract // is exercised over the wire. sendingPolicy := sendingpolicy.DisabledPolicy() + sendingPolicy.AccountTrustEnabled = true sendingPolicy.ExternalSendingAccess = &sendingpolicy.ExternalSendingAccessPolicy{ Mode: sendingpolicy.ModeEnforce, AccountsCreatedAtOrAfter: ContractExternalAccessCutoff, @@ -216,6 +217,7 @@ func StartContractServer(ctx context.Context, dbURL string) (*ContractServer, er api.SetProviderSubmitter(providerSubmitter, sendingGate) api.SetExternalAccess(sendingModule) api.SetIdempotencyStore(idempotencyStore) + enforcer.SetAccountDailyControl(sendingModule.AccountTrustEnabled) api.SetEnforcer(enforcer) api.SetUsageStore(usageStore) api.SetSubscriberStore(subscriberStore) diff --git a/mcp/src/tools/agents.ts b/mcp/src/tools/agents.ts index ca7b07c2f..759d5e88e 100644 --- a/mcp/src/tools/agents.ts +++ b/mcp/src/tools/agents.ts @@ -51,7 +51,7 @@ export function registerAgentTools(server: McpServer, client: McpClient): void { title: "Get the authenticated account's identity", annotations: { readOnlyHint: true }, description: - "Use first when starting work on e2a to learn WHO you are: the authenticated user (email), the credential's scope (`account` or `agent`), and your plan + usage limits. For an agent-scoped credential it also returns `agent_email` — the single agent that credential IS. Account-scoped credentials own many agents; discover them with `list_agents`. This is identity, not an agent — it never guesses a 'default' agent. When the deployment restricts external sending, the optional `sending_access` object reports whether this account may email external recipients (`enforcement_applies`, `shared_external_approved`, `paid_external_sending_entitled`, `owner_recipient_verified`) and `available_unlocks` — what can lift the restriction on this deployment: `operator_approval` (always; file one request with `request_sending_access`), and only where listed `verified_domain` (send from your own verified domain) or `paid_entitlement` (a paid plan). A paid plan that is not listed there does NOT lift it. `read_only: true` means the account is frozen while its sending is paused for an abuse review: read tools keep working, and every tool that changes anything (send, create, update, delete, approve, …) fails with `account_read_only` — do not retry those; the account owner must contact support.", + "Use first when starting work on e2a to learn WHO you are: the authenticated user (email), the credential's scope (`account` or `agent`), and your plan + usage limits. When present, `daily_limit` reports today's external-recipient allowance, usage (including pending sends), and UTC reset time; shared-identity usage is a subset capped by `shared_limit`. Internal mail to own live agents and the verified owner mailbox does not consume it. For an agent-scoped credential it also returns `agent_email` — the single agent that credential IS. Account-scoped credentials own many agents; discover them with `list_agents`. This is identity, not an agent — it never guesses a 'default' agent. When the deployment restricts external sending, the optional `sending_access` object reports whether this account may email external recipients (`enforcement_applies`, `shared_external_approved`, `paid_external_sending_entitled`, `owner_recipient_verified`) and `available_unlocks` — what can lift the restriction on this deployment: `operator_approval` (always; file one request with `request_sending_access`), and only where listed `verified_domain` (send from your own verified domain) or `paid_entitlement` (a paid plan). A paid plan that is not listed there does NOT lift it. `read_only: true` means the account is frozen while its sending is paused for an abuse review: read tools keep working, and every tool that changes anything (send, create, update, delete, approve, …) fails with `account_read_only` — do not retry those; the account owner must contact support.", inputSchema: strictInputSchema({}), }, async () => runTool(() => client.whoami()), diff --git a/migrations/130_account_sending_trust.sql b/migrations/130_account_sending_trust.sql new file mode 100644 index 000000000..ba155afd4 --- /dev/null +++ b/migrations/130_account_sending_trust.sql @@ -0,0 +1,31 @@ +-- Account trust is opt-in. No existing account or domain is automatically exempted. +CREATE TABLE IF NOT EXISTS account_sending_trust ( + user_id TEXT PRIMARY KEY REFERENCES users(id) ON DELETE CASCADE, + archived_clean_days INTEGER NOT NULL DEFAULT 0 CHECK (archived_clean_days >= 0), + grandfather_daily INTEGER NOT NULL DEFAULT 0 CHECK (grandfather_daily BETWEEN 0 AND 2000) +); +CREATE TABLE IF NOT EXISTS account_send_days ( + user_id TEXT NOT NULL REFERENCES account_sending_trust(user_id) ON DELETE CASCADE, + day DATE NOT NULL, + reserved_count INTEGER NOT NULL DEFAULT 0 CHECK (reserved_count >= 0), + shared_count INTEGER NOT NULL DEFAULT 0 CHECK (shared_count >= 0), + confirmed_count INTEGER NOT NULL DEFAULT 0 CHECK (confirmed_count >= 0), + breached BOOLEAN NOT NULL DEFAULT false, + PRIMARY KEY (user_id, day) +); +CREATE TABLE IF NOT EXISTS account_send_reservations ( + message_id TEXT PRIMARY KEY REFERENCES messages(id) ON DELETE CASCADE, + user_id TEXT NOT NULL REFERENCES account_sending_trust(user_id) ON DELETE CASCADE, + day DATE NOT NULL, + units INTEGER NOT NULL CHECK (units > 0), + shared BOOLEAN NOT NULL, + state TEXT NOT NULL CHECK (state IN ('reserved','confirmed','released')), + FOREIGN KEY (user_id,day) REFERENCES account_send_days(user_id,day) +); +CREATE INDEX IF NOT EXISTS account_send_reservations_owner ON account_send_reservations(user_id,day); + +-- Internal-only customer messages still need an authorization token but expose +-- zero external recipients. Preserve all existing positive reservations. +ALTER TABLE sending_budget_reservations DROP CONSTRAINT IF EXISTS sending_budget_reservations_units_check; +ALTER TABLE sending_budget_reservations ADD CONSTRAINT sending_budget_reservations_units_check CHECK (units >= 0) NOT VALID; +ALTER TABLE sending_budget_reservations VALIDATE CONSTRAINT sending_budget_reservations_units_check; diff --git a/migrations/131_account_trust_attempt_provenance.sql b/migrations/131_account_trust_attempt_provenance.sql new file mode 100644 index 000000000..3531a8f83 --- /dev/null +++ b/migrations/131_account_trust_attempt_provenance.sql @@ -0,0 +1,5 @@ +-- Attempt units have different meanings across runtime-policy generations. +-- Only explicitly external-recipient units may earn account trust progress. +-- Existing attempts remain NULL and conservatively earn no account trust. +ALTER TABLE sending_budget_reservations ADD COLUMN IF NOT EXISTS account_trust_units INTEGER + CHECK (account_trust_units IS NULL OR account_trust_units >= 0); diff --git a/sdks/python/CHANGELOG.md b/sdks/python/CHANGELOG.md index 11a93877b..3d741a4b8 100644 --- a/sdks/python/CHANGELOG.md +++ b/sdks/python/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +- Account responses and daily quota refusals expose an optional daily-limit snapshot: external-recipient allowance, reserved usage, shared-identity subset, and UTC reset time. + +## Unreleased + ### Breaking - **``account.delete()`` now moves the account to the trash by default instead of erasing it immediately.** Every API key, OAuth grant, and diff --git a/sdks/python/src/e2a/v1/client.py b/sdks/python/src/e2a/v1/client.py index d930377af..28b7411fc 100644 --- a/sdks/python/src/e2a/v1/client.py +++ b/sdks/python/src/e2a/v1/client.py @@ -1409,6 +1409,9 @@ def __init__(self, api: AccountApi, client: AsyncE2AClient) -> None: async def get(self) -> AccountView: """The authenticated account (whoami). + ``daily_limit``, when present, reports today's external-recipient + allowance, reserved usage, shared-identity subset and UTC reset time. + On a deployment that restricts external sending, ``sending_access`` (beta) reports the account's state; its ``available_unlocks`` lists the routes that deployment accepts for lifting the restriction diff --git a/sdks/python/src/e2a/v1/generated/__init__.py b/sdks/python/src/e2a/v1/generated/__init__.py index eb67505f0..be3a4fe48 100644 --- a/sdks/python/src/e2a/v1/generated/__init__.py +++ b/sdks/python/src/e2a/v1/generated/__init__.py @@ -40,6 +40,7 @@ "ApiException", "APIKeyExportEntry", "APIKeyView", + "AccountDailyLimit", "AccountMetricsView", "AccountUserView", "AccountView", @@ -238,6 +239,7 @@ # import models into sdk package from e2a.v1.generated.models.api_key_export_entry import APIKeyExportEntry as APIKeyExportEntry from e2a.v1.generated.models.api_key_view import APIKeyView as APIKeyView +from e2a.v1.generated.models.account_daily_limit import AccountDailyLimit as AccountDailyLimit from e2a.v1.generated.models.account_metrics_view import AccountMetricsView as AccountMetricsView from e2a.v1.generated.models.account_user_view import AccountUserView as AccountUserView from e2a.v1.generated.models.account_view import AccountView as AccountView diff --git a/sdks/python/src/e2a/v1/generated/models/__init__.py b/sdks/python/src/e2a/v1/generated/models/__init__.py index 4156fbf64..991181499 100644 --- a/sdks/python/src/e2a/v1/generated/models/__init__.py +++ b/sdks/python/src/e2a/v1/generated/models/__init__.py @@ -15,6 +15,7 @@ # import models into model package from e2a.v1.generated.models.api_key_export_entry import APIKeyExportEntry from e2a.v1.generated.models.api_key_view import APIKeyView +from e2a.v1.generated.models.account_daily_limit import AccountDailyLimit from e2a.v1.generated.models.account_metrics_view import AccountMetricsView from e2a.v1.generated.models.account_user_view import AccountUserView from e2a.v1.generated.models.account_view import AccountView diff --git a/sdks/python/src/e2a/v1/generated/models/account_daily_limit.py b/sdks/python/src/e2a/v1/generated/models/account_daily_limit.py new file mode 100644 index 000000000..a23ab9e2d --- /dev/null +++ b/sdks/python/src/e2a/v1/generated/models/account_daily_limit.py @@ -0,0 +1,111 @@ +# coding: utf-8 + +""" + e2a API + + e2a — authenticated email gateway for AI agents. v1 contract. ## Stability policy The v1 surface is stable and evolves **additively only**: new endpoints, new optional request fields, new response fields, and new values in open string sets (event types, statuses) may appear at any time without a version bump. Clients MUST tolerate unknown response fields and unknown values in open string sets. This is machine-readable in the schemas: response schemas declare `additionalProperties: true`; request schemas stay strict (`additionalProperties: false` — an unknown request field is rejected with 422). Operations and schemas marked `x-stability-level: beta` are exempt from this freeze and may change or be removed without a major version. A field marked `x-experimental-values` is itself stable, but the listed values (and their event payloads) are experimental. Everything not marked beta, or enumerated as experimental, is stable. Removing or changing stable surface only happens on a new major version path (/v2); deprecations are announced ahead of time via `deprecated: true` in this document and keep working within v1. + + The version of the OpenAPI document: 1.0.0 + Generated by OpenAPI Generator (https://openapi-generator.tech) + + Do not edit the class manually. +""" # noqa: E501 + + +from __future__ import annotations +import pprint +import re # noqa: F401 +import json + +from datetime import datetime +from pydantic import BaseModel, ConfigDict, StrictInt +from typing import Any, ClassVar, Dict, List +from typing import Optional, Set +from typing_extensions import Self + +class AccountDailyLimit(BaseModel): + """ + AccountDailyLimit + """ # noqa: E501 + clean_active_days: StrictInt + limit: StrictInt + resets_at: datetime + shared_limit: StrictInt + shared_used: StrictInt + used: StrictInt + additional_properties: Dict[str, Any] = {} + __properties: ClassVar[List[str]] = ["clean_active_days", "limit", "resets_at", "shared_limit", "shared_used", "used"] + + model_config = ConfigDict( + populate_by_name=True, + validate_assignment=True, + protected_namespaces=(), + ) + + + def to_str(self) -> str: + """Returns the string representation of the model using alias""" + return pprint.pformat(self.model_dump(by_alias=True)) + + def to_json(self) -> str: + """Returns the JSON representation of the model using alias""" + # TODO: pydantic v2: use .model_dump_json(by_alias=True, exclude_unset=True) instead + return json.dumps(self.to_dict()) + + @classmethod + def from_json(cls, json_str: str) -> Optional[Self]: + """Create an instance of AccountDailyLimit from a JSON string""" + return cls.from_dict(json.loads(json_str)) + + def to_dict(self) -> Dict[str, Any]: + """Return the dictionary representation of the model using alias. + + This has the following differences from calling pydantic's + `self.model_dump(by_alias=True)`: + + * `None` is only added to the output dict for nullable fields that + were set at model initialization. Other fields with value `None` + are ignored. + * Fields in `self.additional_properties` are added to the output dict. + """ + excluded_fields: Set[str] = set([ + "additional_properties", + ]) + + _dict = self.model_dump( + by_alias=True, + exclude=excluded_fields, + exclude_none=True, + ) + # puts key-value pairs in additional_properties in the top level + if self.additional_properties is not None: + for _key, _value in self.additional_properties.items(): + _dict[_key] = _value + + return _dict + + @classmethod + def from_dict(cls, obj: Optional[Dict[str, Any]]) -> Optional[Self]: + """Create an instance of AccountDailyLimit from a dict""" + if obj is None: + return None + + if not isinstance(obj, dict): + return cls.model_validate(obj) + + _obj = cls.model_validate({ + "clean_active_days": obj.get("clean_active_days"), + "limit": obj.get("limit"), + "resets_at": obj.get("resets_at"), + "shared_limit": obj.get("shared_limit"), + "shared_used": obj.get("shared_used"), + "used": obj.get("used") + }) + # store additional fields in additional_properties + for _key in obj.keys(): + if _key not in cls.__properties: + _obj.additional_properties[_key] = obj.get(_key) + + return _obj + + diff --git a/sdks/python/src/e2a/v1/generated/models/account_view.py b/sdks/python/src/e2a/v1/generated/models/account_view.py index 1fad544c5..07ad4c918 100644 --- a/sdks/python/src/e2a/v1/generated/models/account_view.py +++ b/sdks/python/src/e2a/v1/generated/models/account_view.py @@ -20,6 +20,7 @@ from datetime import datetime from pydantic import BaseModel, ConfigDict, Field, StrictBool, StrictStr from typing import Any, ClassVar, Dict, List, Optional +from e2a.v1.generated.models.account_daily_limit import AccountDailyLimit from e2a.v1.generated.models.account_user_view import AccountUserView from e2a.v1.generated.models.limits_caps_view import LimitsCapsView from e2a.v1.generated.models.limits_usage_view import LimitsUsageView @@ -32,6 +33,7 @@ class AccountView(BaseModel): AccountView """ # noqa: E501 agent_email: Optional[StrictStr] = None + daily_limit: Optional[AccountDailyLimit] = Field(default=None, description="External-recipient allowance for this UTC day. Used includes pending or uncertain provider submissions. Shared-identity usage is included in total usage and also bounded by shared_limit. Internal recipients (own live agents and verified owner mailbox) do not count. Omitted when this deployment does not enable the account trust ladder.") deleted_at: Optional[datetime] = Field(default=None, description="When the account was moved to the trash. Absent for a live account.") limits: LimitsCapsView plan_code: StrictStr @@ -44,7 +46,7 @@ class AccountView(BaseModel): usage: LimitsUsageView user: AccountUserView additional_properties: Dict[str, Any] = {} - __properties: ClassVar[List[str]] = ["agent_email", "deleted_at", "limits", "plan_code", "purge_after", "read_only", "restored_at", "scope", "sending_access", "upgrade_url", "usage", "user"] + __properties: ClassVar[List[str]] = ["agent_email", "daily_limit", "deleted_at", "limits", "plan_code", "purge_after", "read_only", "restored_at", "scope", "sending_access", "upgrade_url", "usage", "user"] model_config = ConfigDict( populate_by_name=True, @@ -87,6 +89,9 @@ def to_dict(self) -> Dict[str, Any]: exclude=excluded_fields, exclude_none=True, ) + # override the default output from pydantic by calling `to_dict()` of daily_limit + if self.daily_limit: + _dict['daily_limit'] = self.daily_limit.to_dict() # override the default output from pydantic by calling `to_dict()` of limits if self.limits: _dict['limits'] = self.limits.to_dict() @@ -117,6 +122,7 @@ def from_dict(cls, obj: Optional[Dict[str, Any]]) -> Optional[Self]: _obj = cls.model_validate({ "agent_email": obj.get("agent_email"), + "daily_limit": AccountDailyLimit.from_dict(obj["daily_limit"]) if obj.get("daily_limit") is not None else None, "deleted_at": obj.get("deleted_at"), "limits": LimitsCapsView.from_dict(obj["limits"]) if obj.get("limits") is not None else None, "plan_code": obj.get("plan_code"), diff --git a/sdks/python/src/e2a/v1/generated/models/limit_exceeded_details.py b/sdks/python/src/e2a/v1/generated/models/limit_exceeded_details.py index c5a66fead..b80decbea 100644 --- a/sdks/python/src/e2a/v1/generated/models/limit_exceeded_details.py +++ b/sdks/python/src/e2a/v1/generated/models/limit_exceeded_details.py @@ -19,6 +19,7 @@ from pydantic import BaseModel, ConfigDict, Field, StrictInt, StrictStr from typing import Any, ClassVar, Dict, List, Optional +from e2a.v1.generated.models.account_daily_limit import AccountDailyLimit from typing import Optional, Set from typing_extensions import Self @@ -27,12 +28,13 @@ class LimitExceededDetails(BaseModel): LimitExceededDetails """ # noqa: E501 current: StrictInt = Field(description="The account's usage at the time the cap was hit (matches usage.).") + daily_limit: Optional[AccountDailyLimit] = Field(default=None, description="Current external-recipient allowance, usage and reset time when the account trust ladder refuses an immediate send.") limit: StrictInt = Field(description="The cap that was hit (matches limits.max_).") plan_code: Optional[StrictStr] = Field(default=None, description="The account's plan label.") - resource: StrictStr = Field(description="The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (per-UTC-day send cap; no AccountView field — resets at midnight UTC).") + resource: StrictStr = Field(description="The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (daily send cap; daily_limit reports external-recipient usage on deployments with the account trust ladder — resets at midnight UTC).") upgrade_url: Optional[StrictStr] = Field(default=None, description="An upgrade affordance URL, when the operator has configured one.") additional_properties: Dict[str, Any] = {} - __properties: ClassVar[List[str]] = ["current", "limit", "plan_code", "resource", "upgrade_url"] + __properties: ClassVar[List[str]] = ["current", "daily_limit", "limit", "plan_code", "resource", "upgrade_url"] model_config = ConfigDict( populate_by_name=True, @@ -75,6 +77,9 @@ def to_dict(self) -> Dict[str, Any]: exclude=excluded_fields, exclude_none=True, ) + # override the default output from pydantic by calling `to_dict()` of daily_limit + if self.daily_limit: + _dict['daily_limit'] = self.daily_limit.to_dict() # puts key-value pairs in additional_properties in the top level if self.additional_properties is not None: for _key, _value in self.additional_properties.items(): @@ -93,6 +98,7 @@ def from_dict(cls, obj: Optional[Dict[str, Any]]) -> Optional[Self]: _obj = cls.model_validate({ "current": obj.get("current"), + "daily_limit": AccountDailyLimit.from_dict(obj["daily_limit"]) if obj.get("daily_limit") is not None else None, "limit": obj.get("limit"), "plan_code": obj.get("plan_code"), "resource": obj.get("resource"), diff --git a/sdks/python/src/e2a/v1/generated/models/sending_ramp_view.py b/sdks/python/src/e2a/v1/generated/models/sending_ramp_view.py index 7fa480069..17986cb52 100644 --- a/sdks/python/src/e2a/v1/generated/models/sending_ramp_view.py +++ b/sdks/python/src/e2a/v1/generated/models/sending_ramp_view.py @@ -27,12 +27,12 @@ class SendingRampView(BaseModel): SendingRampView """ # noqa: E501 active_days: StrictInt = Field(description="UTC days that reached the provider-accepted volume threshold.") - daily_recipient_limit: StrictInt = Field(description="Current UTC-day recipient allowance. Zero means no ramp cap applies.") + daily_recipient_limit: StrictInt = Field(description="Current UTC-day recipient allowance. Zero means no per-domain ramp cap applies; account daily_limit can still apply.") estimated_completion_at: Optional[datetime] = Field(default=None, description="Earliest estimated completion assuming every remaining UTC day reaches the provider-accepted volume threshold.") ramp_days: StrictInt recipients_used_today: StrictInt = Field(description="Recipient capacity reserved for the current UTC day, including submissions whose provider outcome is still pending.") resets_at: Optional[datetime] = None - status: StrictStr = Field(description="Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt.") + status: StrictStr = Field(description="Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt, account_managed (the account daily_limit replaces the per-domain ramp).") additional_properties: Dict[str, Any] = {} __properties: ClassVar[List[str]] = ["active_days", "daily_recipient_limit", "estimated_completion_at", "ramp_days", "recipients_used_today", "resets_at", "status"] diff --git a/sdks/python/tests/test_v1_client.py b/sdks/python/tests/test_v1_client.py index 0eac54d61..2364adf17 100644 --- a/sdks/python/tests/test_v1_client.py +++ b/sdks/python/tests/test_v1_client.py @@ -2164,3 +2164,17 @@ async def test_message_delete_permanent_reports_a_deferral(httpx_mock): res = await c.messages.delete("bot@agents.localhost", "msg_sent", permanent=True) assert res.erase_deferred is True assert res.purge_after is not None + + +@pytest.mark.anyio +async def test_account_daily_limit_decodes_utc_reset(httpx_mock): + from e2a.v1.generated.models.account_view import AccountView + httpx_mock.add_response(json=_valid(AccountView, daily_limit={ + "limit": 88, "used": 17, "shared_limit": 50, "shared_used": 6, + "clean_active_days": 1, "resets_at": "2026-01-02T00:00:00Z", + })) + async with _client() as c: + account = await c.account.get() + assert account.daily_limit.limit == 88 + assert account.daily_limit.shared_used == 6 + assert account.daily_limit.resets_at.isoformat() == "2026-01-02T00:00:00+00:00" diff --git a/sdks/typescript/CHANGELOG.md b/sdks/typescript/CHANGELOG.md index b349adcc1..8a2467211 100644 --- a/sdks/typescript/CHANGELOG.md +++ b/sdks/typescript/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +- Account responses and daily quota refusals expose an optional daily-limit snapshot: external-recipient allowance, reserved usage, shared-identity subset, and UTC reset time. + +## Unreleased + ### Breaking - **`account.delete()` now moves the account to the trash by default instead of erasing it immediately.** Every API key, OAuth grant, and dashboard diff --git a/sdks/typescript/src/v1/client.ts b/sdks/typescript/src/v1/client.ts index 230d04e02..d786522b5 100644 --- a/sdks/typescript/src/v1/client.ts +++ b/sdks/typescript/src/v1/client.ts @@ -924,6 +924,7 @@ class AccountResource { * — is always among them; `"verified_domain"` and `"paid_entitlement"` only * where configured). Undefined from servers that predate the field. */ + /** Account identity and usage, including optional dailyLimit for external recipients. */ get(): Promise { return call(() => this.api.getAccount()); } diff --git a/sdks/typescript/src/v1/generated/.openapi-generator/FILES b/sdks/typescript/src/v1/generated/.openapi-generator/FILES index 212598ace..6bfbe7aee 100644 --- a/sdks/typescript/src/v1/generated/.openapi-generator/FILES +++ b/sdks/typescript/src/v1/generated/.openapi-generator/FILES @@ -19,6 +19,7 @@ index.ts middleware.ts models/APIKeyExportEntry.ts models/APIKeyView.ts +models/AccountDailyLimit.ts models/AccountMetricsView.ts models/AccountUserView.ts models/AccountView.ts diff --git a/sdks/typescript/src/v1/generated/models/AccountDailyLimit.ts b/sdks/typescript/src/v1/generated/models/AccountDailyLimit.ts new file mode 100644 index 000000000..bc0d1376d --- /dev/null +++ b/sdks/typescript/src/v1/generated/models/AccountDailyLimit.ts @@ -0,0 +1,71 @@ +/** + * e2a API + * e2a — authenticated email gateway for AI agents. v1 contract. ## Stability policy The v1 surface is stable and evolves **additively only**: new endpoints, new optional request fields, new response fields, and new values in open string sets (event types, statuses) may appear at any time without a version bump. Clients MUST tolerate unknown response fields and unknown values in open string sets. This is machine-readable in the schemas: response schemas declare `additionalProperties: true`; request schemas stay strict (`additionalProperties: false` — an unknown request field is rejected with 422). Operations and schemas marked `x-stability-level: beta` are exempt from this freeze and may change or be removed without a major version. A field marked `x-experimental-values` is itself stable, but the listed values (and their event payloads) are experimental. Everything not marked beta, or enumerated as experimental, is stable. Removing or changing stable surface only happens on a new major version path (/v2); deprecations are announced ahead of time via `deprecated: true` in this document and keep working within v1. + * + * OpenAPI spec version: 1.0.0 + * + * + * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). + * https://openapi-generator.tech + * Do not edit the class manually. + */ + +import { HttpFile } from '../http/http.js'; + +export class AccountDailyLimit { + 'cleanActiveDays': number; + 'limit': number; + 'resetsAt': Date; + 'sharedLimit': number; + 'sharedUsed': number; + 'used': number; + + static readonly discriminator: string | undefined = undefined; + + static readonly mapping: {[index: string]: string} | undefined = undefined; + + static readonly attributeTypeMap: Array<{name: string, baseName: string, type: string, format: string}> = [ + { + "name": "cleanActiveDays", + "baseName": "clean_active_days", + "type": "number", + "format": "int64" + }, + { + "name": "limit", + "baseName": "limit", + "type": "number", + "format": "int64" + }, + { + "name": "resetsAt", + "baseName": "resets_at", + "type": "Date", + "format": "date-time" + }, + { + "name": "sharedLimit", + "baseName": "shared_limit", + "type": "number", + "format": "int64" + }, + { + "name": "sharedUsed", + "baseName": "shared_used", + "type": "number", + "format": "int64" + }, + { + "name": "used", + "baseName": "used", + "type": "number", + "format": "int64" + } ]; + + static getAttributeTypeMap() { + return AccountDailyLimit.attributeTypeMap; + } + + public constructor() { + } +} diff --git a/sdks/typescript/src/v1/generated/models/AccountView.ts b/sdks/typescript/src/v1/generated/models/AccountView.ts index 77069564e..f83ef0f34 100644 --- a/sdks/typescript/src/v1/generated/models/AccountView.ts +++ b/sdks/typescript/src/v1/generated/models/AccountView.ts @@ -10,6 +10,7 @@ * Do not edit the class manually. */ +import { AccountDailyLimit } from '../models/AccountDailyLimit.js'; import { AccountUserView } from '../models/AccountUserView.js'; import { LimitsCapsView } from '../models/LimitsCapsView.js'; import { LimitsUsageView } from '../models/LimitsUsageView.js'; @@ -19,6 +20,10 @@ import { HttpFile } from '../http/http.js'; export class AccountView { 'agentEmail'?: string; /** + * External-recipient allowance for this UTC day. Used includes pending or uncertain provider submissions. Shared-identity usage is included in total usage and also bounded by shared_limit. Internal recipients (own live agents and verified owner mailbox) do not count. Omitted when this deployment does not enable the account trust ladder. + */ + 'dailyLimit'?: AccountDailyLimit; + /** * When the account was moved to the trash. Absent for a live account. */ 'deletedAt'?: Date; @@ -59,6 +64,12 @@ export class AccountView { "type": "string", "format": "" }, + { + "name": "dailyLimit", + "baseName": "daily_limit", + "type": "AccountDailyLimit", + "format": "" + }, { "name": "deletedAt", "baseName": "deleted_at", diff --git a/sdks/typescript/src/v1/generated/models/LimitExceededDetails.ts b/sdks/typescript/src/v1/generated/models/LimitExceededDetails.ts index b80694336..33fa4d4d0 100644 --- a/sdks/typescript/src/v1/generated/models/LimitExceededDetails.ts +++ b/sdks/typescript/src/v1/generated/models/LimitExceededDetails.ts @@ -10,6 +10,7 @@ * Do not edit the class manually. */ +import { AccountDailyLimit } from '../models/AccountDailyLimit.js'; import { HttpFile } from '../http/http.js'; export class LimitExceededDetails { @@ -18,6 +19,10 @@ export class LimitExceededDetails { */ 'current': number; /** + * Current external-recipient allowance, usage and reset time when the account trust ladder refuses an immediate send. + */ + 'dailyLimit'?: AccountDailyLimit; + /** * The cap that was hit (matches limits.max_). */ 'limit': number; @@ -26,7 +31,7 @@ export class LimitExceededDetails { */ 'planCode'?: string; /** - * The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (per-UTC-day send cap; no AccountView field — resets at midnight UTC). + * The capped resource stem. For stems with AccountView fields, key it to usage. and limits.max_. Open set: new values may be added over time, so treat these as strings and tolerate unknown values. Known values: agents, domains, messages_month, storage_bytes, messages_day (daily send cap; daily_limit reports external-recipient usage on deployments with the account trust ladder — resets at midnight UTC). */ 'resource': string; /** @@ -45,6 +50,12 @@ export class LimitExceededDetails { "type": "number", "format": "int64" }, + { + "name": "dailyLimit", + "baseName": "daily_limit", + "type": "AccountDailyLimit", + "format": "" + }, { "name": "limit", "baseName": "limit", diff --git a/sdks/typescript/src/v1/generated/models/ObjectSerializer.ts b/sdks/typescript/src/v1/generated/models/ObjectSerializer.ts index 8d8d1adb9..8e2cfcd88 100644 --- a/sdks/typescript/src/v1/generated/models/ObjectSerializer.ts +++ b/sdks/typescript/src/v1/generated/models/ObjectSerializer.ts @@ -1,5 +1,6 @@ export * from '../models/APIKeyExportEntry.js'; export * from '../models/APIKeyView.js'; +export * from '../models/AccountDailyLimit.js'; export * from '../models/AccountMetricsView.js'; export * from '../models/AccountUserView.js'; export * from '../models/AccountView.js'; @@ -172,6 +173,7 @@ export * from '../models/WebhookView.js'; import { APIKeyExportEntry } from '../models/APIKeyExportEntry.js'; import { APIKeyView } from '../models/APIKeyView.js'; +import { AccountDailyLimit } from '../models/AccountDailyLimit.js'; import { AccountMetricsView } from '../models/AccountMetricsView.js'; import { AccountUserView } from '../models/AccountUserView.js'; import { AccountView } from '../models/AccountView.js'; @@ -384,6 +386,7 @@ let enumsMap: Set = new Set([ let typeMap: {[index: string]: any} = { "APIKeyExportEntry": APIKeyExportEntry, "APIKeyView": APIKeyView, + "AccountDailyLimit": AccountDailyLimit, "AccountMetricsView": AccountMetricsView, "AccountUserView": AccountUserView, "AccountView": AccountView, diff --git a/sdks/typescript/src/v1/generated/models/SendingRampView.ts b/sdks/typescript/src/v1/generated/models/SendingRampView.ts index 603efd4f6..7a7355076 100644 --- a/sdks/typescript/src/v1/generated/models/SendingRampView.ts +++ b/sdks/typescript/src/v1/generated/models/SendingRampView.ts @@ -17,7 +17,7 @@ export class SendingRampView { */ 'activeDays': number; /** - * Current UTC-day recipient allowance. Zero means no ramp cap applies. + * Current UTC-day recipient allowance. Zero means no per-domain ramp cap applies; account daily_limit can still apply. */ 'dailyRecipientLimit': number; /** @@ -31,7 +31,7 @@ export class SendingRampView { 'recipientsUsedToday': number; 'resetsAt'?: Date; /** - * Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt. + * Platform-managed sending-ramp state. Open set; known values: inactive, ramping, complete, exempt, account_managed (the account daily_limit replaces the per-domain ramp). */ 'status': string; diff --git a/sdks/typescript/src/v1/generated/models/all.ts b/sdks/typescript/src/v1/generated/models/all.ts index c9952a07d..71678e622 100644 --- a/sdks/typescript/src/v1/generated/models/all.ts +++ b/sdks/typescript/src/v1/generated/models/all.ts @@ -1,5 +1,6 @@ export * from '../models/APIKeyExportEntry.js' export * from '../models/APIKeyView.js' +export * from '../models/AccountDailyLimit.js' export * from '../models/AccountMetricsView.js' export * from '../models/AccountUserView.js' export * from '../models/AccountView.js' diff --git a/sdks/typescript/src/v1/generated/types/ObjectParamAPI.ts b/sdks/typescript/src/v1/generated/types/ObjectParamAPI.ts index a01f327d7..c7f29873c 100644 --- a/sdks/typescript/src/v1/generated/types/ObjectParamAPI.ts +++ b/sdks/typescript/src/v1/generated/types/ObjectParamAPI.ts @@ -4,6 +4,7 @@ import type { Middleware } from '../middleware.js'; import { APIKeyExportEntry } from '../models/APIKeyExportEntry.js'; import { APIKeyView } from '../models/APIKeyView.js'; +import { AccountDailyLimit } from '../models/AccountDailyLimit.js'; import { AccountMetricsView } from '../models/AccountMetricsView.js'; import { AccountUserView } from '../models/AccountUserView.js'; import { AccountView } from '../models/AccountView.js'; diff --git a/sdks/typescript/src/v1/generated/types/ObservableAPI.ts b/sdks/typescript/src/v1/generated/types/ObservableAPI.ts index bf741fe5d..f0705998b 100644 --- a/sdks/typescript/src/v1/generated/types/ObservableAPI.ts +++ b/sdks/typescript/src/v1/generated/types/ObservableAPI.ts @@ -5,6 +5,7 @@ import { Observable, of, from } from '../rxjsStub.js'; import {mergeMap, map} from '../rxjsStub.js'; import { APIKeyExportEntry } from '../models/APIKeyExportEntry.js'; import { APIKeyView } from '../models/APIKeyView.js'; +import { AccountDailyLimit } from '../models/AccountDailyLimit.js'; import { AccountMetricsView } from '../models/AccountMetricsView.js'; import { AccountUserView } from '../models/AccountUserView.js'; import { AccountView } from '../models/AccountView.js'; diff --git a/sdks/typescript/src/v1/generated/types/PromiseAPI.ts b/sdks/typescript/src/v1/generated/types/PromiseAPI.ts index cced3e17e..b5effcfdb 100644 --- a/sdks/typescript/src/v1/generated/types/PromiseAPI.ts +++ b/sdks/typescript/src/v1/generated/types/PromiseAPI.ts @@ -3,6 +3,7 @@ import { Configuration, PromiseConfigurationOptions, wrapOptions } from '../conf import { PromiseMiddleware, Middleware, PromiseMiddlewareWrapper } from '../middleware.js'; import { APIKeyExportEntry } from '../models/APIKeyExportEntry.js'; import { APIKeyView } from '../models/APIKeyView.js'; +import { AccountDailyLimit } from '../models/AccountDailyLimit.js'; import { AccountMetricsView } from '../models/AccountMetricsView.js'; import { AccountUserView } from '../models/AccountUserView.js'; import { AccountView } from '../models/AccountView.js'; diff --git a/sdks/typescript/test/v1/client.test.ts b/sdks/typescript/test/v1/client.test.ts index fbb8018ae..d3265aea9 100644 --- a/sdks/typescript/test/v1/client.test.ts +++ b/sdks/typescript/test/v1/client.test.ts @@ -1289,6 +1289,14 @@ describe("E2AClient", () => { expect(lastCall().url).toContain("/v1/account/suppressions"); }); + it("account.get decodes daily usage and the UTC reset", async () => { + globalThis.fetch=mockFetch(200,{daily_limit:{limit:88,used:17,shared_limit:50,shared_used:6,clean_active_days:1,resets_at:"2026-01-02T00:00:00Z"}}); + const account=await client.account.get(); + expect(account.dailyLimit?.limit).toBe(88); + expect(account.dailyLimit?.sharedUsed).toBe(6); + expect(account.dailyLimit?.resetsAt.toISOString()).toBe("2026-01-02T00:00:00.000Z"); + }); + it("account.get decodes sending_access.available_unlocks (beta, additive)", async () => { globalThis.fetch = mockFetch(200, { plan: "free", diff --git a/tests/contract/scenarios.yaml b/tests/contract/scenarios.yaml index 1ddeff1f3..7c4c80661 100644 --- a/tests/contract/scenarios.yaml +++ b/tests/contract/scenarios.yaml @@ -3856,3 +3856,52 @@ scenarios: status: 401 body_match: "error.code": unauthorized + + - name: account_daily_trust_limit + description: > + Account daily limits count only external recipients. The real contract + service has the ladder enabled; workers are stopped so no external mail + is submitted. An immediate 21-recipient send exceeds the initial 20 + allowance and persists nothing. The same scenario runs in Go, TS and Python. + setup: + - register_agent: + email: "trust-{scenario_token}@agents.localhost" + steps: + - id: account_reports_daily_allowance + action: request + method: GET + path: /v1/account + expect: + status: 200 + body_match: + "daily_limit.limit": 20 + "daily_limit.used": 0 + "daily_limit.shared_limit": 20 + "daily_limit.shared_used": 0 + "daily_limit.clean_active_days": 0 + body_contains: ["daily_limit"] + capture: + daily_reset: daily_limit.resets_at + - id: daily_refusal_includes_snapshot + action: request + method: POST + path: /v1/agents/trust-{scenario_token}@agents.localhost/messages + body: + to: ["recipient-0@example.test", "recipient-1@example.test", "recipient-2@example.test", "recipient-3@example.test", "recipient-4@example.test", "recipient-5@example.test", "recipient-6@example.test", "recipient-7@example.test", "recipient-8@example.test", "recipient-9@example.test", "recipient-10@example.test", "recipient-11@example.test", "recipient-12@example.test", "recipient-13@example.test", "recipient-14@example.test", "recipient-15@example.test", "recipient-16@example.test", "recipient-17@example.test", "recipient-18@example.test", "recipient-19@example.test", "recipient-20@example.test"] + subject: Synthetic daily limit probe + text: This request must not queue mail. + expect: + status: 402 + body_match: + "error.code": limit_exceeded + "error.details.resource": messages_day + "error.details.daily_limit.limit": 20 + "error.details.daily_limit.used": 0 + "error.details.daily_limit.resets_at": "{daily_reset}" + cleanup: + - id: delete_trust_agent + action: request + method: DELETE + path: /v1/agents/trust-{scenario_token}@agents.localhost?confirm=DELETE&permanent=true + expect: + status: 200 diff --git a/web/src/app/(app)/billing/page.test.tsx b/web/src/app/(app)/billing/page.test.tsx index 7694764b0..650e5a0e1 100644 --- a/web/src/app/(app)/billing/page.test.tsx +++ b/web/src/app/(app)/billing/page.test.tsx @@ -53,6 +53,14 @@ function stageLimits(payload: unknown) { } describe("BillingPage", () => { + it("shows external daily usage, shared subset and UTC reset", async () => { + stageLimits({plan_code:"default",limits:{max_agents:10,max_domains:1,max_messages_month:1000,max_storage_bytes:1000},usage:{agents:1,domains:0,messages_month:17,storage_bytes:0},upgrade_url:"",daily_limit:{limit:88,used:17,shared_limit:50,shared_used:6,clean_active_days:1,resets_at:"2026-01-02T00:00:00Z"}}); + renderPage(); + expect(await screen.findByText("External recipients today")).toBeInTheDocument(); + expect(screen.getByText(/Shared identity: 6 of 50/)).toBeInTheDocument(); + expect(screen.getByText(/2026-01-02T00:00:00Z/)).toBeInTheDocument(); + }); + it("renders plan name and all four usage rows", async () => { stageLimits({ plan_code: "default", diff --git a/web/src/app/(app)/billing/page.tsx b/web/src/app/(app)/billing/page.tsx index 266ca2d05..fb1870476 100644 --- a/web/src/app/(app)/billing/page.tsx +++ b/web/src/app/(app)/billing/page.tsx @@ -4,13 +4,14 @@ import { useEffect, useRef, useState } from "react"; import useSWR from "swr"; import { billingPolling } from "../../../lib/livePolling"; import { limitsKey } from "../../../lib/swrKeys"; +import type { DailySendingLimit } from "../../components/onboarding/api"; import { PageShell } from "../../components/loft/PageShell"; // LimitsInfo matches the LimitsView shape returned by GET /v1/account. // Kept inline rather than imported from a generated client because the -// OSS SDK doesn't expose this endpoint yet (it's a dashboard-only -// surface — SDK consumers would call /agents and /messages directly). +// dashboard uses the HTTP wire names directly. type LimitsInfo = { + daily_limit?: DailySendingLimit; plan_code: string; limits: { max_agents: number; @@ -1106,6 +1107,13 @@ export default function BillingPage() { limit={formatNumber(data.limits.max_messages_month)} pct={pct(data.usage.messages_month, data.limits.max_messages_month)} /> + {data.daily_limit && ( +
+ +

Shared identity: {formatNumber(data.daily_limit.shared_used)} of {formatNumber(data.daily_limit.shared_limit)} (included above). Resets {data.daily_limit.resets_at} (UTC).

+

Pending sends count toward this allowance. Mail to your own inboxes and verified account email does not.

+
+ )}