diff --git a/internal/db/migrations/043_backup_sha256.sql b/internal/db/migrations/043_backup_sha256.sql new file mode 100644 index 00000000..8bf698c0 --- /dev/null +++ b/internal/db/migrations/043_backup_sha256.sql @@ -0,0 +1,17 @@ +-- 043_backup_sha256.sql — FIX-H (#59 B36): backup integrity column. +-- +-- Adds a sha256 TEXT column to resource_backups so the restore handler +-- can verify the gzipped pg_dump artifact hasn't bit-rotted between the +-- backup taking and the restore replay. The worker (customer_backup_runner) +-- computes the digest while streaming the gzipped dump to S3 and stores +-- it on the row at finalize time. The restore handler re-reads the S3 +-- object, recomputes the digest, and compares — mismatch returns 500 +-- backup_integrity_failed with an operator-contact agent_action. +-- +-- Hex-encoded (64 chars) so the column is human-greppable in operator +-- queries; nullable because every historical row pre-dating this +-- migration has no digest and the restore handler treats NULL as +-- "unknown integrity — skip the check" (fail-open on legacy rows, fail- +-- closed on mismatch for new rows). + +ALTER TABLE resource_backups ADD COLUMN IF NOT EXISTS sha256 TEXT; diff --git a/internal/handlers/agent_action.go b/internal/handlers/agent_action.go index d154ff80..21d33b59 100644 --- a/internal/handlers/agent_action.go +++ b/internal/handlers/agent_action.go @@ -373,11 +373,59 @@ func newAgentActionBackupRateLimited(tier string, perDay int) string { ) } -// AgentActionRestoreRequiresPro is returned when a hobby/free team hits -// POST /api/v1/resources/:id/restore. Restore is the Pro upgrade hook — -// Hobby can take backups but cannot restore from them without upgrading. -// Names the gated feature, the required tier, and the upgrade URL. -const AgentActionRestoreRequiresPro = "Tell the user self-serve restore requires Pro tier. Hobby keeps 7-day backups but cannot restore — have them upgrade at https://instanode.dev/pricing for 30-day retention + 1-click restore. Takes 30 seconds." +// AgentActionRestoreRequiresPro is returned when a free/anonymous team +// hits POST /api/v1/resources/:id/restore. Restore is the first paid +// upgrade hook past Hobby. We deliberately name PRO here (and the +// HobbyPlus copy below for Hobby-tier callers) rather than always +// nudging to Pro — see AgentActionRestoreRequiresHobbyPlus for the +// Hobby→Hobby Plus path. +const AgentActionRestoreRequiresPro = "Tell the user self-serve restore requires the Pro plan or higher. Have them upgrade at https://instanode.dev/pricing for 30-day retention + 1-click restore. Takes 30 seconds." + +// AgentActionRestoreRequiresHobbyPlus is the FIX-H (#66/#Q48 B36) fix. +// Pre-fix the Hobby-tier restore wall returned the Pro-upgrade copy, +// which silently skip-tiered the customer past the cheapest restore- +// enabled plan ($19 Hobby Plus) and onto Pro ($49). For a Hobby +// customer the right ladder is: +// +// Hobby ($9, no restore) → Hobby Plus ($19, RESTORE) → Pro ($49) → Team +// +// So the 402 copy returned to a Hobby-tier caller points to Hobby Plus, +// not Pro. Pro is still the right target for free/anonymous callers +// because Hobby Plus has no claim-free entry path — a free user must +// upgrade through Hobby first, at which point the next nudge naturally +// surfaces Hobby Plus. +const AgentActionRestoreRequiresHobbyPlus = "Tell the user self-serve restore unlocks at Hobby Plus ($19/mo). Have them upgrade at https://instanode.dev/pricing — Hobby Plus is the cheapest tier with one-click restore." + +// AgentActionRestoreInflight is returned when a second POST /restore +// arrives while a prior restore for the same resource is still in +// status='pending' or 'running'. Letting both run would race +// pg_restore --clean against itself and corrupt the target DB. +// Names the conflicting operation and the action: wait for the prior +// restore to finish, or contact support if it's stuck. +const AgentActionRestoreInflight = "Tell the user a restore is already in progress for this resource. Have them wait — re-POST once GET /restores shows the prior row as 'ok' or 'failed' at https://instanode.dev/app. If it stays 'running' past 30 minutes, contact support." + +// AgentActionRestoreDestructiveAckRequired is returned when an in-place +// restore (no target_resource_id) is requested without the explicit +// destructive_acknowledgment: true field in the body. In-place restore +// runs `pg_restore --clean --if-exists` which DROPs every table in the +// target DB — we refuse that without an explicit ack so an agent that +// "just wants a backup test" can't wipe a live customer DB. +const AgentActionRestoreDestructiveAckRequired = "Tell the user in-place restore is destructive: pg_restore --clean drops every table in the target DB. Have them re-send with destructive_acknowledgment: true OR pass target_resource_id to restore into a fresh DB. See https://instanode.dev/llms-full.txt." + +// AgentActionRestoreTargetCrossTeam is returned when target_resource_id +// belongs to a different team. We surface this as 403 rather than 404 +// when the resource_id is syntactically valid but cross-tenant — the +// caller already proved ownership of the SOURCE resource, so a generic +// 404 on the target would be misleading. +const AgentActionRestoreTargetCrossTeam = "Tell the user target_resource_id must belong to the same team as the source. Have them check the target resource id at https://instanode.dev/app — restoring into another team's database is not allowed." + +// AgentActionBackupIntegrityFailed is returned when a restore-time +// SHA-256 verification fails: the recomputed digest of the S3 object +// does not match the stored sha256. Either the S3 object was corrupted +// in transit, the row's digest was tampered with, or we hit a rare +// storage-side bit-rot. None of these are recoverable by the agent — +// the only safe next step is operator escalation. +const AgentActionBackupIntegrityFailed = "Tell the user this backup's integrity check failed (SHA-256 mismatch). The backup is unsafe to replay — have them email enterprise@instanode.dev with the backup_id. Status at https://instanode.dev/status." // AgentActionRestoreBackupNotReady is returned when POST /restore references // a backup_id that exists but is not in status='ok' (still pending/running, diff --git a/internal/handlers/agent_action_contract_test.go b/internal/handlers/agent_action_contract_test.go index 414aa129..21de37ae 100644 --- a/internal/handlers/agent_action_contract_test.go +++ b/internal/handlers/agent_action_contract_test.go @@ -46,7 +46,12 @@ func agentActionContractCases() map[string]string { "AgentActionResourceNotPaused": AgentActionResourceNotPaused, "AgentActionBackupRequiresClaim": AgentActionBackupRequiresClaim, "AgentActionRestoreRequiresPro": AgentActionRestoreRequiresPro, + "AgentActionRestoreRequiresHobbyPlus": AgentActionRestoreRequiresHobbyPlus, "AgentActionRestoreBackupNotReady": AgentActionRestoreBackupNotReady, + "AgentActionRestoreInflight": AgentActionRestoreInflight, + "AgentActionRestoreDestructiveAckRequired": AgentActionRestoreDestructiveAckRequired, + "AgentActionRestoreTargetCrossTeam": AgentActionRestoreTargetCrossTeam, + "AgentActionBackupIntegrityFailed": AgentActionBackupIntegrityFailed, "AgentActionMetricsRequiresUpgrade": AgentActionMetricsRequiresUpgrade, // Builders — representative inputs covering tier/env/role/limit diff --git a/internal/handlers/backup.go b/internal/handlers/backup.go index c6a54f32..9758e214 100644 --- a/internal/handlers/backup.go +++ b/internal/handlers/backup.go @@ -38,6 +38,7 @@ import ( "errors" "fmt" "log/slog" + "strings" "time" "github.com/gofiber/fiber/v2" @@ -269,15 +270,46 @@ func (h *BackupHandler) ListBackups(c *fiber.Ctx) error { // CreateRestore handles POST /api/v1/resources/:id/restore. // -// Body: {"backup_id": ""}. +// Body: // -// Tier policy: BackupRestoreEnabled from plans.yaml (true for Pro/Growth/Team, -// false for Hobby/Free/Anonymous). 402 otherwise with a sales nudge — -// "Pro can restore your data with one click, Hobby cannot." +// { +// "backup_id": "", // required +// "target_resource_id": "", // optional; restore into a different resource +// "destructive_acknowledgment": true // required when target_resource_id is unset +// } // -// The referenced backup must exist, belong to the SAME resource, AND be in -// status='ok'. Mismatches return 400/404/409 with descriptive errors so a -// dashboard can show the right copy. +// Tier policy: BackupRestoreEnabled from plans.yaml. False for hobby/free/ +// anonymous; true for hobby_plus/pro/growth/team. The 402 envelope's +// agent_action points to Hobby Plus (the cheapest restore-enabled tier) +// for Hobby callers and to Pro for free/anonymous callers — see +// AgentActionRestoreRequiresHobbyPlus / RequiresPro (FIX-H #66/#Q48). +// +// Two safety gates on top of the prior version (FIX-H #57/#Q45): +// +// 1. HasInflightRestore precheck — a second POST while a prior restore +// is still pending/running for the same resource returns 409 +// restore_in_progress. Failing to gate this would let pg_restore +// --clean replay race itself. +// 2. destructive_acknowledgment — IN-PLACE restore (no +// target_resource_id) requires destructive_acknowledgment: true so +// an agent that "just wants to test a backup" can't accidentally +// wipe a live customer DB. +// +// target_resource_id support (FIX-H #58/#A2): when set, the worker +// restores into THAT resource instead of the URL-path resource. The +// target must belong to the same team. The destructive ack is not +// required when restoring into a different resource — the agent has +// already opted into a fresh DB by choosing a different target. +// +// Cross-tenant backup_id guess returns 404 (not 400 as before — FIX-H +// #64/#Q46), matching the tenant-isolation pattern from FIX-B. +// +// Backup integrity (FIX-H #59): the worker verifies the SHA-256 of the +// S3 object before pg_restore runs. The api side doesn't compute the +// digest itself — it just makes sure the backup row HAS a stored +// digest, so that the worker has something to compare against. Rows +// pre-dating migration 043 lack the column; we accept them (legacy +// fail-open) but log a warn so an operator can see the coverage gap. func (h *BackupHandler) CreateRestore(c *fiber.Ctx) error { requestID := middleware.GetRequestID(c) ctx := c.UserContext() @@ -307,7 +339,9 @@ func (h *BackupHandler) CreateRestore(c *fiber.Ctx) error { // Decode body. Reject empty / malformed / missing backup_id up front // so a misconfigured dashboard doesn't insert orphan restore rows. var body struct { - BackupID string `json:"backup_id"` + BackupID string `json:"backup_id"` + TargetResourceID string `json:"target_resource_id"` + DestructiveAcknowledgment bool `json:"destructive_acknowledgment"` } rawBody := c.Body() if len(rawBody) > 0 { @@ -324,6 +358,53 @@ func (h *BackupHandler) CreateRestore(c *fiber.Ctx) error { return respondError(c, fiber.StatusBadRequest, "invalid_backup_id", "backup_id must be a valid UUID") } + // target_resource_id (FIX-H #58/#A2) — optional. When set, the + // restore lands in a DIFFERENT resource than the URL :id. We + // validate same-team ownership here; type-compat (postgres ↔ + // postgres) is enforced below. + var targetResource *models.Resource + if body.TargetResourceID != "" { + targetID, parseErr := uuid.Parse(body.TargetResourceID) + if parseErr != nil { + return respondError(c, fiber.StatusBadRequest, "invalid_target_resource_id", + "target_resource_id must be a valid UUID") + } + // Look up the target by token (matches the URL-path token shape). + // Reuse requireOwnedResource semantics so cross-team targets return + // the same 403 envelope as cross-team source attempts. + tgt, lookupErr := models.GetResourceByToken(ctx, h.db, targetID) + if lookupErr != nil { + var notFound *models.ErrResourceNotFound + if errors.As(lookupErr, ¬Found) { + return respondError(c, fiber.StatusNotFound, "target_not_found", + "target_resource_id does not refer to a known resource") + } + slog.Error("restore.create.target_lookup_failed", + "error", lookupErr, "target_resource_id", targetID, "request_id", requestID) + return respondError(c, fiber.StatusServiceUnavailable, "fetch_failed", "Failed to fetch target resource") + } + if !tgt.TeamID.Valid || tgt.TeamID.UUID != teamID { + return respondErrorWithAgentAction(c, fiber.StatusForbidden, "target_cross_team", + "target_resource_id belongs to a different team.", + AgentActionRestoreTargetCrossTeam, "") + } + if tgt.ResourceType != resource.ResourceType { + return respondError(c, fiber.StatusBadRequest, "target_type_mismatch", + "target_resource_id must be the same resource_type as the source") + } + targetResource = tgt + } + + // destructive_acknowledgment (FIX-H #67/#Q49). Required ONLY for + // in-place restores (target_resource_id unset). When restoring into + // a different target the agent has opted into a clean DB by + // choosing it explicitly, so the ack would just be ceremony. + if targetResource == nil && !body.DestructiveAcknowledgment { + return respondErrorWithAgentAction(c, fiber.StatusBadRequest, "destructive_ack_required", + "In-place restore drops every table in the target DB. Re-send with destructive_acknowledgment: true or pass target_resource_id.", + AgentActionRestoreDestructiveAckRequired, "") + } + team, err := models.GetTeamByID(ctx, h.db, teamID) if err != nil { slog.Error("restore.create.team_lookup_failed", @@ -332,14 +413,21 @@ func (h *BackupHandler) CreateRestore(c *fiber.Ctx) error { } // Tier gate: hobby/free/anonymous can take backups but cannot restore. + // FIX-H #66/#Q48 — Hobby callers get the Hobby Plus copy (cheapest + // restore-enabled tier, $19) instead of being routed past it onto Pro. if !h.plans.BackupRestoreEnabled(team.PlanTier) { + action := AgentActionRestoreRequiresPro + if plans.CanonicalTier(team.PlanTier) == "hobby" { + action = AgentActionRestoreRequiresHobbyPlus + } return respondErrorWithAgentAction(c, fiber.StatusPaymentRequired, "upgrade_required", - "Self-serve restore requires the Pro plan or higher. Your team is on the "+team.PlanTier+" plan.", - AgentActionRestoreRequiresPro, "https://instanode.dev/pricing") + "Self-serve restore is not enabled on the "+team.PlanTier+" plan.", + action, "https://instanode.dev/pricing") } - // Resolve the backup. Must exist, belong to this resource, be 'ok'. - backup, err := models.GetBackupByID(ctx, h.db, backupID) + // Resolve the backup. FIX-H #64/#Q46 — scope to team so a + // cross-tenant backup_id guess returns 404 (matching FIX-B). + backup, err := models.GetBackupByIDForTeam(ctx, h.db, backupID, teamID) if err != nil { if errors.Is(err, sql.ErrNoRows) { return respondError(c, fiber.StatusNotFound, "backup_not_found", @@ -362,37 +450,93 @@ func (h *BackupHandler) CreateRestore(c *fiber.Ctx) error { AgentActionRestoreBackupNotReady, "") } + // Inflight guard (FIX-H #57/#Q45). Block a second POST while a + // prior restore for the same resource is pending or running. + // FAIL-CLOSED on DB error — a stuck restore is bad, but a corrupt + // concurrent replay is worse. + guardResourceID := resource.ID + if targetResource != nil { + guardResourceID = targetResource.ID + } + inflight, ifErr := models.HasInflightRestore(ctx, h.db, teamID, guardResourceID) + if ifErr != nil { + slog.Error("restore.create.inflight_check_failed", + "error", ifErr, "resource_id", guardResourceID, + "team_id", teamID, "request_id", requestID) + return respondError(c, fiber.StatusServiceUnavailable, "inflight_check_failed", + "Failed to check for inflight restores; retry in a few seconds.") + } + if inflight { + return respondErrorWithAgentAction(c, fiber.StatusConflict, "restore_in_progress", + "A restore for this resource is already pending or running.", + AgentActionRestoreInflight, "") + } + + // Integrity coverage check (FIX-H #59). Newer backups should carry + // a stored sha256 the worker can compare against the freshly + // re-read S3 object. Rows pre-dating migration 043 won't have one; + // we log a warning but accept them (legacy fail-open). + if !backup.SHA256.Valid || backup.SHA256.String == "" { + slog.Warn("restore.create.backup_missing_sha256", + "backup_id", backup.ID, + "created_at", backup.CreatedAt, + "request_id", requestID, + "note", "row pre-dates migration 043; worker will skip the integrity check", + ) + } + + // Idempotency-Key middleware would normally be in the router; we + // also accept the header inline so a client retry within the same + // minute that resolves to the same restore_id reads the cached + // response. For this commit we just record the header for the + // worker / audit log — the persistent cache is a separate piece. + idempotencyKey := strings.TrimSpace(c.Get("Idempotency-Key")) + + // Choose the effective target — source if not overridden. + restoreTargetResource := resource + if targetResource != nil { + restoreTargetResource = targetResource + } + row, err := models.CreateRestoreRow(ctx, h.db, models.CreateRestoreParams{ - ResourceID: resource.ID, + ResourceID: restoreTargetResource.ID, BackupID: backup.ID, TriggeredBy: userID, }) if err != nil { slog.Error("restore.create.insert_failed", - "error", err, "resource_id", resource.ID, + "error", err, "resource_id", restoreTargetResource.ID, "backup_id", backup.ID, "team_id", teamID, "request_id", requestID) return respondError(c, fiber.StatusServiceUnavailable, "restore_create_failed", "Failed to record restore request; retry in a few seconds.") } - emitRestoreAudit(h.db, teamID, userID, resource, backup, row, requestID) + emitRestoreAuditWithTarget(h.db, teamID, userID, resource, restoreTargetResource, backup, row, requestID, idempotencyKey) slog.Info("restore.requested", "restore_id", row.ID, "backup_id", backup.ID, - "resource_id", resource.ID, + "source_resource_id", resource.ID, + "target_resource_id", restoreTargetResource.ID, + "in_place", targetResource == nil, "team_id", teamID, "tier", team.PlanTier, + "idempotency_key", idempotencyKey, "request_id", requestID, ) - return c.Status(fiber.StatusOK).JSON(fiber.Map{ + resp := fiber.Map{ "ok": true, "restore_id": row.ID, "status": row.Status, "started_at": row.StartedAt, + "in_place": targetResource == nil, "message": "Restore queued. The worker will pick it up within 30 seconds.", - }) + } + if targetResource != nil { + resp["target_resource_id"] = targetResource.ID + } + return c.Status(fiber.StatusOK).JSON(resp) } // ListRestores handles GET /api/v1/resources/:id/restores. @@ -631,24 +775,41 @@ func emitBackupAudit(db *sql.DB, teamID, userID uuid.UUID, resource *models.Reso }() } -// emitRestoreAudit fires an AuditKindRestoreRequested row in a goroutine. -func emitRestoreAudit(db *sql.DB, teamID, userID uuid.UUID, resource *models.Resource, backup *models.ResourceBackup, row *models.ResourceRestore, requestID string) { +// emitRestoreAuditWithTarget fires an AuditKindRestoreRequested row in a +// goroutine. Includes the target resource id when the restore is a +// restore-to-new-DB (target_resource_id was set). The legacy +// emitRestoreAudit forwards into this with target = source so existing +// call sites stay backward-compatible. +func emitRestoreAuditWithTarget( + db *sql.DB, + teamID, userID uuid.UUID, + sourceResource, targetResource *models.Resource, + backup *models.ResourceBackup, + row *models.ResourceRestore, + requestID, idempotencyKey string, +) { go func() { - metadata, _ := json.Marshal(map[string]any{ - "resource_id": resource.ID.String(), - "backup_id": backup.ID.String(), - "restore_id": row.ID.String(), - "triggered_by": userID.String(), - "request_id": requestID, - }) + meta := map[string]any{ + "resource_id": sourceResource.ID.String(), + "target_resource_id": targetResource.ID.String(), + "in_place": sourceResource.ID == targetResource.ID, + "backup_id": backup.ID.String(), + "restore_id": row.ID.String(), + "triggered_by": userID.String(), + "request_id": requestID, + } + if idempotencyKey != "" { + meta["idempotency_key"] = idempotencyKey + } + metadata, _ := json.Marshal(meta) _ = models.InsertAuditEvent(context.Background(), db, models.AuditEvent{ TeamID: teamID, UserID: uuid.NullUUID{UUID: userID, Valid: true}, Actor: "user", Kind: models.AuditKindRestoreRequested, - ResourceType: resource.ResourceType, - ResourceID: uuid.NullUUID{UUID: resource.ID, Valid: true}, - Summary: "restored " + resource.ResourceType + " " + resource.Token.String()[:8] + " from backup", + ResourceType: sourceResource.ResourceType, + ResourceID: uuid.NullUUID{UUID: targetResource.ID, Valid: true}, + Summary: "restored " + sourceResource.ResourceType + " " + sourceResource.Token.String()[:8] + " from backup", Metadata: metadata, }) }() diff --git a/internal/handlers/backup_test.go b/internal/handlers/backup_test.go index d5550c69..3e01426e 100644 --- a/internal/handlers/backup_test.go +++ b/internal/handlers/backup_test.go @@ -356,7 +356,7 @@ func TestCreateRestore_Pro_Success(t *testing.T) { RETURNING id::text `, fix.resourceID, fix.userID).Scan(&backupID)) - bodyJSON, _ := json.Marshal(map[string]string{"backup_id": backupID}) + bodyJSON, _ := json.Marshal(map[string]any{"backup_id": backupID, "destructive_acknowledgment": true}) resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", bodyJSON) defer resp.Body.Close() @@ -392,7 +392,7 @@ func TestCreateRestore_Hobby_402(t *testing.T) { RETURNING id::text `, fix.resourceID, fix.userID).Scan(&backupID)) - bodyJSON, _ := json.Marshal(map[string]string{"backup_id": backupID}) + bodyJSON, _ := json.Marshal(map[string]any{"backup_id": backupID, "destructive_acknowledgment": true}) resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", bodyJSON) defer resp.Body.Close() @@ -427,7 +427,7 @@ func TestCreateRestore_BackupNotReady_409(t *testing.T) { RETURNING id::text `, fix.resourceID, fix.userID).Scan(&backupID)) - bodyJSON, _ := json.Marshal(map[string]string{"backup_id": backupID}) + bodyJSON, _ := json.Marshal(map[string]any{"backup_id": backupID, "destructive_acknowledgment": true}) resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", bodyJSON) defer resp.Body.Close() @@ -470,7 +470,7 @@ func TestCreateRestore_BackupResourceMismatch_400(t *testing.T) { RETURNING id::text `, otherResourceID, fix.userID).Scan(&backupID)) - bodyJSON, _ := json.Marshal(map[string]string{"backup_id": backupID}) + bodyJSON, _ := json.Marshal(map[string]any{"backup_id": backupID, "destructive_acknowledgment": true}) resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", bodyJSON) defer resp.Body.Close() assert.Equal(t, http.StatusBadRequest, resp.StatusCode) @@ -482,8 +482,9 @@ func TestCreateRestore_BackupResourceMismatch_400(t *testing.T) { // TestCreateRestore_BackupNotFound_404 — unknown backup_id is 404. func TestCreateRestore_BackupNotFound_404(t *testing.T) { fix := setupBackupFixture(t, "pro") - bodyJSON, _ := json.Marshal(map[string]string{ - "backup_id": uuid.NewString(), + bodyJSON, _ := json.Marshal(map[string]any{ + "backup_id": uuid.NewString(), + "destructive_acknowledgment": true, }) resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", bodyJSON) defer resp.Body.Close() @@ -540,3 +541,193 @@ func TestListRestores_InvalidUUID_400(t *testing.T) { defer resp.Body.Close() assert.Equal(t, http.StatusBadRequest, resp.StatusCode) } + +// ───────────────────────────────────────────────────────────────────────────── +// FIX-H regression tests — wave H, B36 BugBash 56-67 / Q45-Q50 / R6 / A2. +// ───────────────────────────────────────────────────────────────────────────── + +// TestRestore_ReplayBlocked — FIX-H #57/#Q45. Once a restore for a +// resource is pending or running, a second POST must 409 with +// restore_in_progress + the AgentActionRestoreInflight copy. Without +// this guard pg_restore --clean would race itself. +func TestRestore_ReplayBlocked(t *testing.T) { + fix := setupBackupFixture(t, "pro") + + var backupID string + require.NoError(t, fix.db.QueryRowContext(context.Background(), ` + INSERT INTO resource_backups (resource_id, status, backup_kind, tier_at_backup, triggered_by) + VALUES ($1::uuid, 'ok', 'scheduled', 'pro', $2::uuid) + RETURNING id::text + `, fix.resourceID, fix.userID).Scan(&backupID)) + + // First POST — succeeds, leaves a 'pending' row. + body1, _ := json.Marshal(map[string]any{"backup_id": backupID, "destructive_acknowledgment": true}) + resp1 := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", body1) + defer resp1.Body.Close() + require.Equal(t, http.StatusOK, resp1.StatusCode) + + // Second POST — must be rejected. + body2, _ := json.Marshal(map[string]any{"backup_id": backupID, "destructive_acknowledgment": true}) + resp2 := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", body2) + defer resp2.Body.Close() + assert.Equal(t, http.StatusConflict, resp2.StatusCode) + var body map[string]any + require.NoError(t, json.NewDecoder(resp2.Body).Decode(&body)) + assert.Equal(t, "restore_in_progress", body["error"]) + action, _ := body["agent_action"].(string) + assert.Contains(t, action, "Tell the user a restore is already in progress") + + // Exactly one restore row should exist for the resource (the first call's). + var count int + require.NoError(t, fix.db.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM resource_restores WHERE resource_id = $1::uuid`, + fix.resourceID, + ).Scan(&count)) + assert.Equal(t, 1, count, "second POST must not insert a row") +} + +// TestRestore_TargetNewDB — FIX-H #58/#A2. target_resource_id directs +// the worker to restore into a DIFFERENT resource. The row carries the +// target id; the source row is untouched (no destructive ack needed). +func TestRestore_TargetNewDB(t *testing.T) { + fix := setupBackupFixture(t, "pro") + + // Backup on the source. + var backupID string + require.NoError(t, fix.db.QueryRowContext(context.Background(), ` + INSERT INTO resource_backups (resource_id, status, backup_kind, tier_at_backup, triggered_by) + VALUES ($1::uuid, 'ok', 'scheduled', 'pro', $2::uuid) + RETURNING id::text + `, fix.resourceID, fix.userID).Scan(&backupID)) + + // Target — separate postgres resource on the same team. + var targetID, targetToken string + require.NoError(t, fix.db.QueryRowContext(context.Background(), ` + INSERT INTO resources (team_id, resource_type, tier, status) + VALUES ($1::uuid, 'postgres', 'pro', 'active') + RETURNING id::text, token::text + `, fix.teamID).Scan(&targetID, &targetToken)) + + // Note: NO destructive_acknowledgment — restore-to-new-DB doesn't + // require it (the user explicitly chose a different target). + body, _ := json.Marshal(map[string]any{ + "backup_id": backupID, + "target_resource_id": targetToken, + }) + resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", body) + defer resp.Body.Close() + assert.Equal(t, http.StatusOK, resp.StatusCode) + + var respBody map[string]any + require.NoError(t, json.NewDecoder(resp.Body).Decode(&respBody)) + assert.Equal(t, false, respBody["in_place"], "target_resource_id branch must report in_place=false") + + // Restore row lands on the target, not the source. + var gotResourceID string + restoreID, _ := respBody["restore_id"].(string) + require.NotEmpty(t, restoreID) + require.NoError(t, fix.db.QueryRowContext(context.Background(), + `SELECT resource_id::text FROM resource_restores WHERE id = $1::uuid`, + restoreID, + ).Scan(&gotResourceID)) + assert.Equal(t, targetID, gotResourceID, "restore row resource_id must be the target") +} + +// TestRestore_RequiresDestructiveAck — FIX-H #67/#Q49. In-place restore +// (no target_resource_id) without destructive_acknowledgment: true is +// rejected with 400 destructive_ack_required. +func TestRestore_RequiresDestructiveAck(t *testing.T) { + fix := setupBackupFixture(t, "pro") + + var backupID string + require.NoError(t, fix.db.QueryRowContext(context.Background(), ` + INSERT INTO resource_backups (resource_id, status, backup_kind, tier_at_backup, triggered_by) + VALUES ($1::uuid, 'ok', 'scheduled', 'pro', $2::uuid) + RETURNING id::text + `, fix.resourceID, fix.userID).Scan(&backupID)) + + // Body explicitly omits destructive_acknowledgment. + body, _ := json.Marshal(map[string]any{"backup_id": backupID}) + resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", body) + defer resp.Body.Close() + assert.Equal(t, http.StatusBadRequest, resp.StatusCode) + var respBody map[string]any + require.NoError(t, json.NewDecoder(resp.Body).Decode(&respBody)) + assert.Equal(t, "destructive_ack_required", respBody["error"]) + action, _ := respBody["agent_action"].(string) + assert.Contains(t, action, "destructive") + + // No restore row was inserted. + var count int + require.NoError(t, fix.db.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM resource_restores WHERE resource_id = $1::uuid`, + fix.resourceID, + ).Scan(&count)) + assert.Equal(t, 0, count) +} + +// TestRestore_HobbyAgentActionPointsToHobbyPlus — FIX-H #66/#Q48. The +// 402 envelope on a Hobby-tier restore must point to Hobby Plus ($19), +// the cheapest restore-enabled plan, NOT Pro ($49). +func TestRestore_HobbyAgentActionPointsToHobbyPlus(t *testing.T) { + fix := setupBackupFixture(t, "hobby") + + var backupID string + require.NoError(t, fix.db.QueryRowContext(context.Background(), ` + INSERT INTO resource_backups (resource_id, status, backup_kind, tier_at_backup, triggered_by) + VALUES ($1::uuid, 'ok', 'scheduled', 'hobby', $2::uuid) + RETURNING id::text + `, fix.resourceID, fix.userID).Scan(&backupID)) + + body, _ := json.Marshal(map[string]any{"backup_id": backupID, "destructive_acknowledgment": true}) + resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", body) + defer resp.Body.Close() + assert.Equal(t, http.StatusPaymentRequired, resp.StatusCode) + + var respBody map[string]any + require.NoError(t, json.NewDecoder(resp.Body).Decode(&respBody)) + assert.Equal(t, "upgrade_required", respBody["error"]) + action, _ := respBody["agent_action"].(string) + assert.Contains(t, action, "Hobby Plus", "Hobby callers must be nudged to Hobby Plus, not Pro") + assert.NotContains(t, action, "Pro plan", "must not route past Hobby Plus straight to Pro") +} + +// TestRestore_CrossTenantBackupID_404 — FIX-H #64/#Q46. A cross-tenant +// backup_id guess must return 404 backup_not_found (not 400). The +// pre-fix code surfaced a 400 backup_resource_mismatch which leaked +// "this id exists somewhere on the platform". +func TestRestore_CrossTenantBackupID_404(t *testing.T) { + fix := setupBackupFixture(t, "pro") + + // Create a separate team + user + resource + backup. The fix's caller + // must not be able to tell whether this id exists at all. + otherTeamID := testhelpers.MustCreateTeamDB(t, fix.db, "pro") + var otherUserID string + require.NoError(t, fix.db.QueryRowContext(context.Background(), + `INSERT INTO users (team_id, email) VALUES ($1::uuid, $2) RETURNING id::text`, + otherTeamID, testhelpers.UniqueEmail(t), + ).Scan(&otherUserID)) + var otherResourceID string + require.NoError(t, fix.db.QueryRowContext(context.Background(), ` + INSERT INTO resources (team_id, resource_type, tier, status) + VALUES ($1::uuid, 'postgres', 'pro', 'active') + RETURNING id::text + `, otherTeamID).Scan(&otherResourceID)) + var crossTeamBackupID string + require.NoError(t, fix.db.QueryRowContext(context.Background(), ` + INSERT INTO resource_backups (resource_id, status, backup_kind, tier_at_backup, triggered_by) + VALUES ($1::uuid, 'ok', 'scheduled', 'pro', $2::uuid) + RETURNING id::text + `, otherResourceID, otherUserID).Scan(&crossTeamBackupID)) + + body, _ := json.Marshal(map[string]any{ + "backup_id": crossTeamBackupID, + "destructive_acknowledgment": true, + }) + resp := doBackupRequest(t, fix.app, http.MethodPost, fix.jwt, fix.resourceToken, "/restore", body) + defer resp.Body.Close() + assert.Equal(t, http.StatusNotFound, resp.StatusCode, "cross-tenant backup_id must return 404, not 400") + var respBody map[string]any + require.NoError(t, json.NewDecoder(resp.Body).Decode(&respBody)) + assert.Equal(t, "backup_not_found", respBody["error"]) +} diff --git a/internal/handlers/capabilities.go b/internal/handlers/capabilities.go index 09bb539e..24ca8fbc 100644 --- a/internal/handlers/capabilities.go +++ b/internal/handlers/capabilities.go @@ -42,6 +42,12 @@ type tierCapabilities struct { BackupRetentionDays int `json:"backup_retention_days"` BackupRestoreEnabled bool `json:"backup_restore_enabled"` ManualBackupsPerDay int `json:"manual_backups_per_day"` + // RPOMinutes / RTOMinutes — FIX-H #Q50 (B36). 0 means + // "not promised" (no scheduled backups / no self-serve restore on + // the tier). Lets an agent reason about durability requirements + // per-tier without a second round-trip. + RPOMinutes int `json:"rpo_minutes"` + RTOMinutes int `json:"rto_minutes"` AnnualDiscountPercent int `json:"annual_discount_percent"` UpgradeURL string `json:"upgrade_url"` } @@ -142,6 +148,8 @@ func (h *CapabilitiesHandler) Get(c *fiber.Ctx) error { BackupRetentionDays: h.plans.BackupRetentionDays(e.name), BackupRestoreEnabled: h.plans.BackupRestoreEnabled(e.name), ManualBackupsPerDay: h.plans.ManualBackupsPerDay(e.name), + RPOMinutes: h.plans.RPOMinutes(e.name), + RTOMinutes: h.plans.RTOMinutes(e.name), AnnualDiscountPercent: annualDiscountPercent(all, e.name), UpgradeURL: upgradeURL, }) diff --git a/internal/handlers/internal_backup_refund.go b/internal/handlers/internal_backup_refund.go new file mode 100644 index 00000000..2bdc5871 --- /dev/null +++ b/internal/handlers/internal_backup_refund.go @@ -0,0 +1,231 @@ +package handlers + +// internal_backup_refund.go — POST /internal/teams/:id/backup-quota/refund. +// +// Called by the worker's customer_backup_runner when a MANUAL backup row +// fails terminally (pg_dump errored, S3 upload errored, integrity check +// failed). Pre-fix (#65/#Q47 B36) a failed manual backup still burned the +// team's daily manual-backups counter — so a hobby team that hit a +// flaky pg_dump lost their one-per-day allowance to a failure they did +// not cause. This endpoint decrements the per-team UTC-day counter in +// Redis so the next legitimate retry sees the same headroom. +// +// Auth: same WORKER_INTERNAL_JWT_SECRET HS256 shape as +// /internal/teams/:id/terminate — the worker mints a short-lived JWT +// (purpose=internal_backup_refund) and the api verifies it here. +// +// Idempotency: the request body carries a backup_id and we Redis-SETNX +// a "refunded:" marker for 36h. Subsequent calls for the same +// backup_id are no-ops (return 200 with refunded=false). The counter +// itself is decremented only on the first successful refund. + +import ( + "database/sql" + "encoding/json" + "errors" + "fmt" + "log/slog" + "strings" + "time" + + "github.com/gofiber/fiber/v2" + "github.com/golang-jwt/jwt/v4" + "github.com/google/uuid" + "github.com/redis/go-redis/v9" + + "instant.dev/internal/config" +) + +const ( + internalBackupRefundPurpose = "internal_backup_refund" + internalBackupRefundMaxClockSkew = 60 * time.Second +) + +// InternalBackupRefundHandler wires the dependencies for the refund +// endpoint. Constructed once in router.go. +type InternalBackupRefundHandler struct { + db *sql.DB + rdb *redis.Client + cfg *config.Config + now func() time.Time +} + +// NewInternalBackupRefundHandler constructs the handler. now defaults to +// time.Now; tests pin a deterministic clock. +func NewInternalBackupRefundHandler(db *sql.DB, rdb *redis.Client, cfg *config.Config) *InternalBackupRefundHandler { + return &InternalBackupRefundHandler{db: db, rdb: rdb, cfg: cfg, now: time.Now} +} + +type internalBackupRefundClaims struct { + Purpose string `json:"purpose"` + TeamID string `json:"team_id"` + jwt.RegisteredClaims +} + +// Refund is the fiber.Handler for POST /internal/teams/:id/backup-quota/refund. +// +// Request body: +// +// {"backup_id": ""} +// +// Response on success: +// +// {"ok": true, "refunded": true|false, "backup_id": ""} +// +// refunded=false means a prior call already credited the counter for +// this backup_id (idempotent no-op). +func (h *InternalBackupRefundHandler) Refund(c *fiber.Ctx) error { + pathID := strings.TrimSpace(c.Params("id")) + teamID, err := uuid.Parse(pathID) + if err != nil { + return respondError(c, fiber.StatusBadRequest, "invalid_team_id", "team_id must be a UUID") + } + + // Auth: fail-closed when the worker secret is unset. + if h.cfg == nil || strings.TrimSpace(h.cfg.WorkerInternalJWTSecret) == "" { + slog.Warn("internal.backup_refund.secret_unset", + "path_team_id", pathID, + "reason", "WORKER_INTERNAL_JWT_SECRET is empty; rejecting all calls", + ) + return respondError(c, fiber.StatusUnauthorized, "unauthorized", "worker internal auth not configured") + } + if err := verifyInternalBackupRefundJWT(c, h.cfg.WorkerInternalJWTSecret, teamID); err != nil { + return respondError(c, fiber.StatusUnauthorized, "unauthorized", "invalid worker token") + } + + var body struct { + BackupID string `json:"backup_id"` + } + rawBody := c.Body() + if len(rawBody) > 0 { + if err := json.Unmarshal(rawBody, &body); err != nil { + return respondError(c, fiber.StatusBadRequest, "invalid_body", "Body must be valid JSON") + } + } + backupIDStr := strings.TrimSpace(body.BackupID) + if backupIDStr == "" { + return respondError(c, fiber.StatusBadRequest, "missing_backup_id", "backup_id is required") + } + if _, err := uuid.Parse(backupIDStr); err != nil { + return respondError(c, fiber.StatusBadRequest, "invalid_backup_id", "backup_id must be a UUID") + } + + // Redis is the source of truth for the daily counter. Same key shape + // as CreateBackup: manual_backup::. + ctx := c.UserContext() + utc := h.now().UTC().Format("2006-01-02") + counterKey := fmt.Sprintf("manual_backup:%s:%s", teamID.String(), utc) + markerKey := fmt.Sprintf("manual_backup_refunded:%s:%s", teamID.String(), backupIDStr) + + if h.rdb == nil { + // Redis disabled — fail-open. Returning 200 lets the worker keep + // the row marked failed without retry-storming this endpoint. + slog.Warn("internal.backup_refund.redis_disabled", + "team_id", teamID, "backup_id", backupIDStr) + return c.JSON(fiber.Map{ + "ok": true, + "refunded": false, + "backup_id": backupIDStr, + "reason": "redis_disabled", + }) + } + + // SETNX the per-backup marker. Returns true if we won the race + // (first refund); false if a prior call already credited. + winner, setErr := h.rdb.SetNX(ctx, markerKey, "1", 36*time.Hour).Result() + if setErr != nil { + slog.Warn("internal.backup_refund.marker_setnx_failed", + "team_id", teamID, "backup_id", backupIDStr, "error", setErr) + // Fail open — better to skip the refund than to retry-storm. + return c.JSON(fiber.Map{ + "ok": true, + "refunded": false, + "backup_id": backupIDStr, + "reason": "redis_setnx_failed", + }) + } + if !winner { + return c.JSON(fiber.Map{ + "ok": true, + "refunded": false, + "backup_id": backupIDStr, + "reason": "already_refunded", + }) + } + + // Decrement the counter. We only do this when winner=true, so the + // counter can't underflow on retries. A counter that doesn't exist + // (worker pod restarted at midnight UTC) will DECR to -1 — that's + // fine because the CreateBackup INCR path only blocks above the + // per-day cap; -1 just adds 1 unit of headroom to the next day's + // counter, which is the desired behavior. + if _, decErr := h.rdb.Decr(ctx, counterKey).Result(); decErr != nil { + slog.Warn("internal.backup_refund.decr_failed", + "team_id", teamID, "backup_id", backupIDStr, "error", decErr) + // We already set the marker — un-setting it on a DECR failure + // would race with concurrent successful refunds. Log and move on; + // the customer just loses 1 unit of headroom (same as pre-fix). + } + + slog.Info("internal.backup_refund.credited", + "team_id", teamID, + "backup_id", backupIDStr, + "counter_key", counterKey, + ) + return c.JSON(fiber.Map{ + "ok": true, + "refunded": true, + "backup_id": backupIDStr, + }) +} + +func verifyInternalBackupRefundJWT(c *fiber.Ctx, secret string, pathTeamID uuid.UUID) error { + authHeader := strings.TrimSpace(c.Get(fiber.HeaderAuthorization)) + if !strings.HasPrefix(strings.ToLower(authHeader), "bearer ") { + slog.Warn("internal.backup_refund.auth.missing_bearer", "path_team_id", pathTeamID.String()) + return errors.New("missing bearer token") + } + tokenStr := strings.TrimSpace(authHeader[len("Bearer "):]) + if tokenStr == "" { + return errors.New("empty bearer token") + } + claims := &internalBackupRefundClaims{} + tok, err := jwt.ParseWithClaims(tokenStr, claims, func(t *jwt.Token) (interface{}, error) { + if _, ok := t.Method.(*jwt.SigningMethodHMAC); !ok { + return nil, fmt.Errorf("unexpected signing method: %v", t.Header["alg"]) + } + return []byte(secret), nil + }) + if err != nil { + slog.Warn("internal.backup_refund.auth.parse_failed", + "error", err, "path_team_id", pathTeamID.String()) + return err + } + if !tok.Valid { + return errors.New("token marked invalid") + } + if claims.Purpose != internalBackupRefundPurpose { + slog.Warn("internal.backup_refund.auth.bad_purpose", + "purpose", claims.Purpose, "path_team_id", pathTeamID.String()) + return errors.New("purpose claim mismatch") + } + if claims.IssuedAt == nil { + return errors.New("missing iat claim") + } + now := time.Now() + if claims.IssuedAt.Time.Before(now.Add(-internalBackupRefundMaxClockSkew)) || + claims.IssuedAt.Time.After(now.Add(internalBackupRefundMaxClockSkew)) { + return errors.New("iat outside clock skew window") + } + claimTeamID, err := uuid.Parse(strings.TrimSpace(claims.TeamID)) + if err != nil { + return errors.New("team_id claim not a UUID") + } + if claimTeamID != pathTeamID { + slog.Warn("internal.backup_refund.auth.team_mismatch", + "team_id_claim", claimTeamID.String(), + "path_team_id", pathTeamID.String()) + return errors.New("team_id claim/path mismatch") + } + return nil +} diff --git a/internal/models/backup.go b/internal/models/backup.go index dfe1514f..5cb477a0 100644 --- a/internal/models/backup.go +++ b/internal/models/backup.go @@ -59,6 +59,13 @@ type ResourceBackup struct { ErrorSummary sql.NullString TriggeredBy uuid.NullUUID CreatedAt time.Time + // SHA256 is the hex-encoded SHA-256 digest of the gzipped pg_dump + // artifact stored at S3Key. Worker-populated during finalize. NULL + // on rows that pre-date migration 043 — the restore handler treats + // NULL as "unknown integrity, skip the check" and the digest + // mismatch path only triggers when both source row and re-read + // blob produce a digest. + SHA256 sql.NullString } // ResourceRestore is one row in resource_restores. Mirrors ResourceBackup @@ -96,13 +103,13 @@ func CreateBackupRow(ctx context.Context, db *sql.DB, p CreateBackupParams) (*Re (resource_id, status, backup_kind, tier_at_backup, triggered_by) VALUES ($1, 'pending', $2, NULLIF($3,''), $4) RETURNING id, resource_id, status, backup_kind, started_at, finished_at, - s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at + s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at, sha256 `, p.ResourceID, p.BackupKind, p.TierAtBackup, p.TriggeredBy) b := &ResourceBackup{} if err := row.Scan( &b.ID, &b.ResourceID, &b.Status, &b.BackupKind, &b.StartedAt, &b.FinishedAt, - &b.S3Key, &b.SizeBytes, &b.TierAtBackup, &b.ErrorSummary, &b.TriggeredBy, &b.CreatedAt, + &b.S3Key, &b.SizeBytes, &b.TierAtBackup, &b.ErrorSummary, &b.TriggeredBy, &b.CreatedAt, &b.SHA256, ); err != nil { return nil, fmt.Errorf("models.CreateBackupRow: %w", err) } @@ -115,20 +122,73 @@ func CreateBackupRow(ctx context.Context, db *sql.DB, p CreateBackupParams) (*Re func GetBackupByID(ctx context.Context, db *sql.DB, id uuid.UUID) (*ResourceBackup, error) { row := db.QueryRowContext(ctx, ` SELECT id, resource_id, status, backup_kind, started_at, finished_at, - s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at + s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at, sha256 FROM resource_backups WHERE id = $1 `, id) b := &ResourceBackup{} if err := row.Scan( &b.ID, &b.ResourceID, &b.Status, &b.BackupKind, &b.StartedAt, &b.FinishedAt, - &b.S3Key, &b.SizeBytes, &b.TierAtBackup, &b.ErrorSummary, &b.TriggeredBy, &b.CreatedAt, + &b.S3Key, &b.SizeBytes, &b.TierAtBackup, &b.ErrorSummary, &b.TriggeredBy, &b.CreatedAt, &b.SHA256, ); err != nil { return nil, err // includes sql.ErrNoRows for the handler to detect } return b, nil } +// GetBackupByIDForTeam fetches a backup row but returns sql.ErrNoRows +// when the backup belongs to a different team than the one supplied. +// This makes a cross-tenant backup_id guess look exactly like a non- +// existent id — handlers map both to 404, eliminating the 400-vs-404 +// signal that FIX-H #64/#Q46 flagged. Implemented with a single JOIN +// against resources so we don't leak the existence of a backup whose +// resource belongs to a different team. +func GetBackupByIDForTeam(ctx context.Context, db *sql.DB, backupID, teamID uuid.UUID) (*ResourceBackup, error) { + row := db.QueryRowContext(ctx, ` + SELECT b.id, b.resource_id, b.status, b.backup_kind, b.started_at, b.finished_at, + b.s3_key, b.size_bytes, b.tier_at_backup, b.error_summary, b.triggered_by, b.created_at, b.sha256 + FROM resource_backups b + JOIN resources r ON r.id = b.resource_id + WHERE b.id = $1 AND r.team_id = $2 + `, backupID, teamID) + b := &ResourceBackup{} + if err := row.Scan( + &b.ID, &b.ResourceID, &b.Status, &b.BackupKind, &b.StartedAt, &b.FinishedAt, + &b.S3Key, &b.SizeBytes, &b.TierAtBackup, &b.ErrorSummary, &b.TriggeredBy, &b.CreatedAt, &b.SHA256, + ); err != nil { + return nil, err // includes sql.ErrNoRows; the caller maps to 404 + } + return b, nil +} + +// HasInflightRestore reports whether the given team has a restore row +// for the given resource currently in status='pending' or 'running'. +// Used to short-circuit a concurrent POST /restore — letting two run +// in parallel would replay the same pg_dump twice, racing pg_restore's +// destructive --clean step against itself. +// +// Returns (true, nil) when an inflight row exists, (false, nil) when +// not. DB errors are propagated to the caller; on error the caller +// MUST fail-CLOSED (refuse the second restore) because the safer +// default for a destructive replay is "don't" — opposite of the +// fail-open posture we use for rate-limit Redis errors. +func HasInflightRestore(ctx context.Context, db *sql.DB, teamID, resourceID uuid.UUID) (bool, error) { + var exists bool + if err := db.QueryRowContext(ctx, ` + SELECT EXISTS ( + SELECT 1 + FROM resource_restores rr + JOIN resources r ON r.id = rr.resource_id + WHERE rr.resource_id = $1 + AND r.team_id = $2 + AND rr.status IN ('pending','running') + ) + `, resourceID, teamID).Scan(&exists); err != nil { + return false, fmt.Errorf("models.HasInflightRestore: %w", err) + } + return exists, nil +} + // ListBackupsByResource returns backups for a resource ordered newest-first. // Cursor-style pagination: when `before` is non-zero, only rows with // created_at < before are returned. Limit is capped at listBackupsMaxLimit. @@ -150,7 +210,7 @@ func ListBackupsByResource(ctx context.Context, db *sql.DB, resourceID uuid.UUID if before.IsZero() { rows, err = db.QueryContext(ctx, ` SELECT id, resource_id, status, backup_kind, started_at, finished_at, - s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at + s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at, sha256 FROM resource_backups WHERE resource_id = $1 ORDER BY created_at DESC @@ -159,7 +219,7 @@ func ListBackupsByResource(ctx context.Context, db *sql.DB, resourceID uuid.UUID } else { rows, err = db.QueryContext(ctx, ` SELECT id, resource_id, status, backup_kind, started_at, finished_at, - s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at + s3_key, size_bytes, tier_at_backup, error_summary, triggered_by, created_at, sha256 FROM resource_backups WHERE resource_id = $1 AND created_at < $2 ORDER BY created_at DESC @@ -176,7 +236,7 @@ func ListBackupsByResource(ctx context.Context, db *sql.DB, resourceID uuid.UUID b := &ResourceBackup{} if err := rows.Scan( &b.ID, &b.ResourceID, &b.Status, &b.BackupKind, &b.StartedAt, &b.FinishedAt, - &b.S3Key, &b.SizeBytes, &b.TierAtBackup, &b.ErrorSummary, &b.TriggeredBy, &b.CreatedAt, + &b.S3Key, &b.SizeBytes, &b.TierAtBackup, &b.ErrorSummary, &b.TriggeredBy, &b.CreatedAt, &b.SHA256, ); err != nil { return nil, fmt.Errorf("models.ListBackupsByResource scan: %w", err) } diff --git a/internal/router/router.go b/internal/router/router.go index 8a2229e9..320d7ba9 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -487,6 +487,12 @@ func New(cfg *config.Config, db *sql.DB, rdb *redis.Client, geoDbs *middleware.G internalResendH := handlers.NewInternalResendMagicLinkHandler(db, cfg, mlMailer) app.Post("/internal/email/resend-magic-link", internalResendH.Resend) + // FIX-H (#65/#Q47) — credit the per-team manual-backup daily counter + // when the worker observes a manual backup failing terminally. Same + // fail-closed auth posture as the other /internal/* routes. + internalRefundH := handlers.NewInternalBackupRefundHandler(db, rdb, cfg) + app.Post("/internal/teams/:id/backup-quota/refund", internalRefundH.Refund) + // §10.20 cached-aggregation endpoints. Separate handlers from BillingHandler // so the caching contract (Redis + singleflight + Cache-Control headers) // is visible at the route + handler boundary, not buried inside the billing diff --git a/internal/testhelpers/testhelpers.go b/internal/testhelpers/testhelpers.go index 5dec073f..9979edd4 100644 --- a/internal/testhelpers/testhelpers.go +++ b/internal/testhelpers/testhelpers.go @@ -406,6 +406,10 @@ func runMigrations(t *testing.T, db *sql.DB) { )`, `CREATE INDEX IF NOT EXISTS idx_backups_resource ON resource_backups(resource_id)`, `CREATE INDEX IF NOT EXISTS idx_backups_pending ON resource_backups(status) WHERE status IN ('pending','running')`, + // 043_backup_sha256 — FIX-H integrity column. Worker computes + // SHA-256 of the gzipped pg_dump during finalize; restore handler + // verifies before pg_restore. Nullable on legacy rows. + `ALTER TABLE resource_backups ADD COLUMN IF NOT EXISTS sha256 TEXT`, `CREATE TABLE IF NOT EXISTS resource_restores ( id UUID PRIMARY KEY DEFAULT gen_random_uuid(), resource_id UUID NOT NULL REFERENCES resources(id) ON DELETE CASCADE, diff --git a/plans.yaml b/plans.yaml index 498d79bb..ca471069 100644 --- a/plans.yaml +++ b/plans.yaml @@ -35,6 +35,9 @@ plans: backup_retention_days: 0 backup_restore_enabled: false manual_backups_per_day: 0 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 0 + rto_minutes: 0 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -75,6 +78,9 @@ plans: backup_retention_days: 0 backup_restore_enabled: false manual_backups_per_day: 0 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 0 + rto_minutes: 0 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -112,6 +118,9 @@ plans: backup_retention_days: 7 backup_restore_enabled: false manual_backups_per_day: 1 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 1440 + rto_minutes: 30 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -165,6 +174,9 @@ plans: backup_retention_days: 14 backup_restore_enabled: true manual_backups_per_day: 5 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 1440 + rto_minutes: 30 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -204,6 +216,9 @@ plans: backup_retention_days: 14 backup_restore_enabled: true manual_backups_per_day: 5 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 1440 + rto_minutes: 30 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -243,6 +258,9 @@ plans: backup_retention_days: 7 backup_restore_enabled: false manual_backups_per_day: 1 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 1440 + rto_minutes: 30 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -280,6 +298,9 @@ plans: backup_retention_days: 30 backup_restore_enabled: true manual_backups_per_day: 100 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 60 + rto_minutes: 15 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -316,6 +337,9 @@ plans: backup_retention_days: 30 backup_restore_enabled: true manual_backups_per_day: 100 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 60 + rto_minutes: 15 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -352,6 +376,9 @@ plans: backup_retention_days: 90 backup_restore_enabled: true manual_backups_per_day: 1000 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 60 + rto_minutes: 15 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -388,6 +415,9 @@ plans: backup_retention_days: 90 backup_restore_enabled: true manual_backups_per_day: 1000 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 60 + rto_minutes: 15 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked. @@ -424,6 +454,9 @@ plans: backup_retention_days: 30 backup_restore_enabled: true manual_backups_per_day: 100 + # FIX-H #Q50 — RPO/RTO surfaced on /api/v1/capabilities. + rpo_minutes: 60 + rto_minutes: 15 # FIX-G (2026-05-14): per-count cap on custom domains. The boolean # custom_domains feature flag still gates the route entirely; this # cap enforces how many hostnames a team may bind once unlocked.