package server import ( "bytes" "context" "crypto/subtle" "encoding/base64" "encoding/json" "errors" "fmt" "io" "io/fs" "log/slog" "net" "net/http" "net/netip" "net/url" "os" "runtime/debug" "strconv" "strings" "sync" "sync/atomic" "time" "github.com/go-chi/chi/v5" chimiddleware "github.com/go-chi/chi/v5/middleware" "github.com/go-chi/cors" "github.com/prometheus/client_golang/prometheus/promhttp" "github.com/PerpetualSoftware/pad/internal/attachments" "github.com/PerpetualSoftware/pad/internal/billing" "github.com/PerpetualSoftware/pad/internal/collab" "github.com/PerpetualSoftware/pad/internal/email" "github.com/PerpetualSoftware/pad/internal/events" "github.com/PerpetualSoftware/pad/internal/metrics" "github.com/PerpetualSoftware/pad/internal/models" "github.com/PerpetualSoftware/pad/internal/oauth" "github.com/PerpetualSoftware/pad/internal/store" "github.com/PerpetualSoftware/pad/internal/textguard" "github.com/PerpetualSoftware/pad/internal/watchevents" "github.com/PerpetualSoftware/pad/internal/webhooks" ) type Server struct { // afterItemPreRead is a TEST-ONLY seam, nil in production. When set, the // item-update handler calls it after loading its own copy of the item and // before handing the write to the store — the window in which another // writer can commit a change this request did not make (BUG-2776). A hook // makes that interleaving deterministic; two real goroutines produce it // only sometimes, and a test that reproduces a race sometimes is a // detector with an unknown rate rather than a regression test. // // Set it only while no other request is in flight against this Server: it // is read on a request path with no synchronisation. // // It fires once per item-update REQUEST, so a hook that issues its own // update must nil the field for the duration of that nested call or it // recurses without end. Every test here does exactly that; the pattern is // load-bearing, not incidental (codex round 3). afterItemPreRead func(itemID string) store *store.Store router *chi.Mux routerOnce sync.Once // ensures setupRouter runs once, after all config admitOnce sync.Once // lazily builds streamAdmit for servers that never call SetSSELimits streamGaugeFor *metrics.Metrics // the metrics instance pad_stream_connections_active is registered on (BUG-2726) httpServer *http.Server // underlying HTTP server (set during ListenAndServe) webFS fs.FS // embedded web UI static files (optional) events events.EventBus // real-time event bus (optional) watchEvents watchevents.Bus // watch/nudge notification bus (optional, TASK-2533) sessionPresence SessionPresence // live event-stream connections per user (optional, PLAN-2558 S1) redisHealth *RedisHealth // cached Redis reachability, reported by /api/v1/health/ready and pad_redis_up (optional, BUG-2727) collab *collab.RoomManager // Yjs collab room manager (PLAN-1248); optional webhooks *webhooks.Dispatcher // webhook dispatcher (optional) email *email.Sender // transactional email sender (optional) emailAPIKey string // Maileroo API key (used for unsubscribe HMAC) emailEnvConfigured bool // email was wired from env vars (SetEmailSender); reconfigureEmail must not tear this down when platform settings clear the key rateLimiters *RateLimiters // per-endpoint rate limiters baseURL string // public base URL for generating links (e.g. invite URLs) corsOrigins string // comma-separated CORS origins (empty = localhost defaults) secureCookies bool // set Secure flag on cookies (for TLS deployments) metrics *metrics.Metrics // Prometheus metrics (optional) metricsToken string // shared bearer token for /metrics scrapes ("" = loopback-only) trustedProxyCIDRs []*net.IPNet // CIDRs allowed to set X-Forwarded-For (nil = proxy headers untrusted) ipChangeEnforceStrict bool // when true, revoke+reject sessions whose client IP OR User-Agent hash differs from the one recorded at session creation sseMaxConnections int // global SSE connection limit (0 = unlimited) sseMaxPerWorkspace int // per-workspace SSE connection limit (0 = unlimited) // midStreamGapCooldownOverride shortens the mid-stream gap rate limit. // Zero (production) means midStreamGapCooldown. Tests only; set before // serving, read per connection. midStreamGapCooldownOverride time.Duration sseMaxPerUser int // per-user SSE connection limit across BOTH stream endpoints (0 = unlimited, BUG-2726) streamAdmit *streamAdmission // shared admission gate for both stream endpoints (BUG-2726) cloudMode bool // true when running as Pad Cloud (PAD_CLOUD=true or PAD_MODE=cloud) cloudSecrets []string // shared secrets for sidecar ↔ pad communication (supports rotation) cloudSidecar CloudSidecar // reverse pad → pad-cloud client (e.g. Stripe cancel on account delete); nil = not configured billingAvailable bool // true when PAD_BILLING_AVAILABLE=true — gates Stripe Checkout CTAs in the web UI (TASK-800) version string // release version (e.g. "dev", "1.2.3") commit string // git commit hash buildTime string // build timestamp twoFAChallengeSecret []byte // HMAC key for 2FA challenge tokens // Attachments storage. Wired via SetAttachments at startup; nil-checked // by handlers so a server constructed for a test that doesn't need // uploads still compiles and serves every other endpoint. attachments *attachments.Registry attachmentMaxBytes int64 // per-file upload cap; 0 = use defaultAttachmentMaxBytes // Image processor used by the upload handler to derive thumbnail // variants (TASK-878) and by the editor's rotate / crop tools // (TASK-879/880). Wired via SetImageProcessor; nil-checked by // callers so a server without image processing — e.g. a self-host // build that doesn't want the dependency — still serves every // other endpoint and stores originals untouched. imageProcessor attachments.Processor // MCP Streamable HTTP transport (PLAN-943 TASK-950). Wired via // SetMCPTransport at startup when the deployment is in cloud mode. // nil on self-hosted deployments and on any cloud build that hasn't // constructed the MCP server yet — registerMCPRoutes nil-checks so // the routes don't mount in either case. See handlers_mcp.go. mcpTransport http.Handler mcpPublicURL string // canonical public URL of the MCP vhost (e.g. https://mcp.getpad.dev) mcpAuthServerURL string // canonical URL of the OAuth auth server (e.g. https://app.getpad.dev), TASK-951 // MCP tool-surface descriptor source (PLAN-1888 / TASK-1891). Wired // via SetToolSurfaceHandler at startup from mcp.ToolSurfaceJSON. The // injection mirrors SetMCPTransport and exists for the same reason: // internal/mcp imports internal/server (dispatch_http.go), so this // package CANNOT import internal/mcp to build the catalog JSON // itself. cmd/pad/main.go imports both and injects the serializer. // nil → GET /api/v1/mcp/tool-surface returns 404 (handler not wired). toolSurfaceJSON func() ([]byte, error) // OAuth 2.1 authorization server (PLAN-943 TASK-1024 sub-PR B, // HTTP handlers in TASK-1025 sub-PR C). Wired via SetOAuthServer // at startup when the deployment is in cloud mode + has the // fosite-backed server constructed. nil disables the OAuth // surface — registerOAuthRoutes nil-checks so the routes don't // mount on self-hosted deployments. See handlers_oauth.go. oauthServer *oauth.Server // claimSecret is the HMAC key for stateless 6-digit claim codes // (PLAN-1519 / TASK-1521 / IDEA-1517 §4). Wired by SetClaimSecret // at startup — production reuses the deployment's 32-byte // encryption key (cfg.EncryptionKey) since both are server- // stable secrets with equivalent rotation cadence. nil/short → // /api/v1/oauth/claim returns 412 "claim_disabled" on every // request, surfacing a clear misconfiguration signal rather than // silently accepting forgeable codes. claimSecret []byte // oauthMetricsWired records whether wireOAuthMetricsObserver has // already attached the active-tokens callback collector. Re- // registering would panic via prometheus.MustRegister, so the flag // guards the one-shot registration. The TTL observer side is // idempotent (just a function-pointer set) and runs unconditionally. oauthMetricsWired bool // MCP audit log async writer (PLAN-943 TASK-960). Spawned by // startMCPAuditWriter at startup when MCP is wired; shut down // from Server.Stop. nil-safe: every audit-emitting code path // nil-checks so MCP-less builds + tests that don't start the // writer still work. See middleware_mcp_audit.go. mcpAudit *mcpAuditWriter // MCP session tracker (PLAN-943 TASK-1120). Replaces the naive // +1/-1 active-sessions accounting from TASK-961. Wired by // startMCPSessionTracker (called from SetMCPTransport in cloud // mode); shut down from Server.Stop alongside the audit writer. // nil-safe: trackMCPSession + the gauge-update path both // nil-check so non-cloud builds + tests run without the tracker. // See middleware_mcp_session.go. mcpSessions *mcpSessionTracker mcpSessionTTL time.Duration // 0 → defaultMCPSessionTTL mcpSessionSweepInterval time.Duration // 0 → defaultMCPSessionSweepInterval // storageInfoCache memoizes per-workspace storage usage summaries // behind a short TTL (storageInfoTTL). Reduces DB load on the // Settings → Storage page and quota-aware UI surfaces. Initialized // in newServer; never nil so handlers can call get/set without // guarding. storageInfoCache *storageInfoCache // copyItemFn indirects Store.CopyItemAcrossWorkspaces for the // cross-workspace copy endpoint (PLAN-2357 / TASK-2365). // // It exists because DR-13 forbids the endpoint from ever transparently // retrying a mutating copy — there is no idempotency key in v1, so a // retry after a post-commit failure duplicates the item — and the only // falsifiable way to assert "called exactly once" is to count the calls. // Several tests also use it to inject the store's typed errors and to // land a concurrent mutation deterministically between the handler's // authorization and the store call. See handleCopyItem and // TestCopyEndpoint_DoesNotRetryOnAmbiguousError. // // It is nil on every production path — nothing outside package server can // set it, no constructor or setter assigns it, and only _test.go files do // (Codex round 6: it is compiled into the binary, so "test-only" describes // the convention, not a compiler-enforced guarantee). copyItemFn func(store.CrossWorkspaceCopyRequest) (*store.CrossWorkspaceCopyResult, error) // importBundleMaxBytes caps a single workspace import bundle. // 0 → defaultImportBundleMaxBytes (2 GiB). Set via // SetImportBundleMaxBytes from cmd/pad/main.go using the // PAD_IMPORT_BUNDLE_MAX_BYTES env var so operators with larger // exports can opt in without recompiling. importBundleMaxBytes int64 // importArtifactMaxBytes caps a single playbook/convention artifact // import (POST /workspaces/{ws}/import-artifact). 0 → // defaultImportArtifactMaxBytes (1 MiB). A single artifact is tiny; // the cap is the first line of defense against an oversized body // being materialized before the YAML-bomb guard runs. Set via // SetImportArtifactMaxBytes from cmd/pad/main.go. importArtifactMaxBytes int64 // orphanGC holds the periodic-sweep config + lifecycle for the // attachment orphan garbage collector (TASK-886). Configured via // SetOrphanGCConfig and started via StartOrphanGC. Stop() signals // the loop to exit and waits for it via the bg WaitGroup. orphanGC orphanGCConfig // opLogGC holds the periodic-sweep config + lifecycle for the // Yjs op-log prune sweeper (TASK-1309). Mirrors orphanGC's // pattern. Configured via SetOpLogGCConfig + started via // StartOpLogGC; Stop() signals the loop via stopOpLogGC. opLogGC opLogGCConfig // tokenReaper holds the periodic-sweep config + lifecycle for the // short-lived-credential reaper (PLAN-1933 DR-5 / TASK-1936). // Mirrors orphanGC/opLogGC. Configured via SetTokenReaperConfig + // started via StartTokenReaper; Stop() signals the loop via // stopTokenReaper. tokenReaper tokenReaperConfig // workspacePurge holds the periodic-sweep config + lifecycle for the // soft-deleted-workspace hard-purge sweeper (TASK-1966 — the 30-day // GDPR erasure SLA). Mirrors orphanGC. Configured via // SetWorkspacePurgeConfig + started via StartWorkspacePurgeSweeper; // Stop() signals the loop via stopWorkspacePurgeSweeper. workspacePurge workspacePurgeConfig // outboxDrain holds the periodic config + lifecycle for the SPEC-3 event // outbox drain (TASK-2714). Mirrors orphanGC. Configured via // SetOutboxDrainConfig + started via StartOutboxDrain; Stop() signals the // loop via stopOutboxDrain. outboxDrain outboxDrainConfig // reminderTick holds the periodic config + lifecycle for the item-reminder // scheduler (IDEA-2641). Mirrors outboxDrain. Configured via // SetReminderTickConfig + started via StartReminderTick; Stop() signals // the loop via stopReminderTick. reminderTick reminderTickConfig // inFlightUploadHashes tracks content_hash values for uploads // that have called AttachmentStore.Put but not yet inserted the // attachments row. Without this, the orphan GC could delete a // blob between Put and CreateAttachment, leaving a live row that // references a missing blob (Codex P2 on PR #307 round 1). // // A plain map + mutex rather than sync.Map: counters need // atomic-with-delete semantics (decrement-then-delete-if-zero // must be one critical section, not two — sync.Map.CompareAndDelete // addresses the entry but not the inc/dec interleaving). Codex // P1 round 2 caught the prior sync.Map version racing on // release-vs-reload of the same hash. inFlightHashesMu sync.Mutex inFlightHashes map[string]int64 // rowlessNoListerOnce gates the once-per-process notice that a // registered attachment backend lacks the Lister capability, leaving // the rowless-blob sweep (BUG-2406) inert for it. Logged rather than // silently skipped so an operator can tell the leak class is // unguarded on that backend; once, so a 24h-cadence sweep doesn't // turn it into log spam. rowlessNoListerOnce sync.Once // rowlessPreDeleteHook, when non-nil, runs inside the rowless // sweep's in-flight critical section immediately BEFORE the // delete-time row re-check. Test seam only (injectedStageFailure // precedent): it lets a test commit a row for the hash at exactly // the point that distinguishes the batched subtraction from the // re-check, making the TOCTOU leg deterministic. rowlessPreDeleteHook func(hash string) // bg tracks fire-and-forget goroutines spawned by request handlers // (TouchUserActivity in middleware_auth, async email sends, etc.) so // the server can drain them before shutdown / test cleanup. Without // this, tests using t.TempDir() race the still-running goroutine's // SQLite WAL write against TempDir RemoveAll, leaving "directory not // empty" cleanup errors in CI. See BUG-842. bg sync.WaitGroup // First-run bootstrap token (TASK-1167 / PLAN-1166). When non-empty, // handleBootstrap accepts the value via the X-Bootstrap-Token header // from non-loopback peers (self-host mode only — cloud mode never // loads or honors a token, D2/D10). Wired at startup via // SetBootstrapToken; cleared by consumeBootstrapToken after the first // admin is created. // // The mutex protects the token field AND the entire validate-token → // check-UserCount → CreateUser → consume sequence in handleBootstrap. // Two simultaneous valid-token requests with different emails would // otherwise create two admins from one token (F5). Bootstrap happens // once per install, so the contention window is irrelevant. bootstrapMu sync.Mutex bootstrapToken string bootstrapTokenPath string // bypassSetupToken, when true, allows the first-admin bootstrap POST to // succeed from any IP without an X-Bootstrap-Token header — i.e. the // /setup form on the web UI works directly, without the operator having // to copy a token out of `docker logs`. Wired from PAD_BYPASS_SETUP_TOKEN // at startup via SetBypassSetupToken (cmd/pad/main.go). // // Self-host only — cloud mode IGNORES this flag entirely (D2/D10 from // the original logs-token design: cloud bootstrap stays loopback-only). // The UserCount==0 gate in handleBootstrap is unchanged: once the first // admin exists, the bootstrap endpoint returns 409 "already initialized" // regardless of bypass. This matches the operator's mental model — the // flag opens up the *first-run* surface, not registration in general. // // Operators on trusted networks (Unraid LAN, Tailscale-only deployments, // homelabs behind a firewall) typically prefer this; operators with // public exposure should leave it off and use the logs-token path. bypassSetupToken bool // restoreAckFault is a TEST SEAM (always nil in production). When non-nil, // handleRestoreItemVersion's collab commit closure invokes it AFTER the restore // transaction has durably committed; a non-nil return simulates a Postgres commit // whose acknowledgement was lost at the connection boundary (the tx landed, but // the driver surfaces an error), exercising BUG-2276 residual 1's commit-outcome // reconciliation end-to-end through the real handler. restoreAckFault func() error // watchPredicatesLoadFault is a TEST SEAM (always nil in production, // TASK-2533). When non-nil, loadWatchPredicates calls it before // touching the store; a non-nil return short-circuits the real // ListWatchesForUser call and is returned as the reload error — // exercising GET /api/v1/events/stream's reval-tick error path // (codex round 4: a watch-list reload failure must not also skip // the identity/visibility refresh) deterministically, without // needing to actually break the DB connection mid-test. // // atomic.Pointer, not a plain func field (codex round 5 finding 2): // unlike restoreAckFault — set once, synchronously, before the single // HTTP request that will read it, so goroutine-creation's own // happens-before edge makes a plain field safe there — this seam is // set by a test AFTER the SSE stream's background goroutine is // already running and reading it on every reval tick. A plain field // written from the test's goroutine while that goroutine reads it // concurrently is a genuine, if timing-dependent, data race // (verified: restoreAckFault's OWN usage doesn't share this flaw, // since it's never touched after the goroutine that reads it starts, // so it was intentionally left as a plain field rather than changed // too). watchPredicatesLoadFault atomic.Pointer[func() error] // watchRevalTickOverride is a TEST SEAM (always nil in production, // BUG-2570). When non-nil, GET /api/v1/events/stream selects reval // ticks from this channel instead of the interval ticker, letting a // test drive each revalidation tick explicitly. That is the only way // to pin assertions to a SPECIFIC tick: with a free-running ticker, // no test ordering can guarantee that an unwanted extra tick doesn't // fire between two test steps — an extra SUCCESSFUL tick resets the // visibility cache and reloads the watch list, masking exactly the // reset-skipped-on-fault regression the reval-fault test guards, and // enough extra FAULTING ticks clear the watch set (both observed as // codex-round findings on BUG-2570's first fix attempt). // // Same atomic.Pointer rationale as watchPredicatesLoadFault above: // read by the stream's background goroutine (once, at stream setup) // while tests may write it. Tests must set it BEFORE connecting the // stream they want to drive — a write after setup is not observed. watchRevalTickOverride atomic.Pointer[chan time.Time] } // goAsync spawns fn in a goroutine that's tracked by s.bg, so Stop() can // wait for in-flight background work to finish. Use this for any // fire-and-forget work that touches the database, filesystem, or external // services from inside a request handler — never bare `go func() {...}()`. func (s *Server) goAsync(fn func()) { s.bg.Add(1) go func() { defer s.bg.Done() // Recover from panics in fn so a single bad background task // (e.g. deriveThumbnails hitting a Go image-decoder panic on a // crafted upload, or an email send) can't crash the whole // single-binary server for every tenant. chi's Recoverer only // covers request goroutines, not these detached ones. The // deferred Done() above still fires because recover() keeps the // goroutine from unwinding past this point. defer func() { if r := recover(); r != nil { slog.Error("background task panicked", "panic", r, "stack", string(debug.Stack())) } }() fn() }() } // recoverSweeper is the panic firewall for the long-running background // sweeper loops (orphan GC, op-log GC, token reaper, workspace purge). // Each of those manages its own s.bg.Add/Done + stop-channel lifecycle, // so — unlike fire-and-forget work — they can't just route through // goAsync without breaking their shutdown handling or double-counting // s.bg. Instead each spawns `defer s.recoverSweeper("")` as a // deferred call INSIDE its goroutine (BUG-2071): a panic in the loop // body is logged with a stack (matching goAsync's style) and unwinds // cleanly, the goroutine's own deferred s.bg.Done() still fires because // recover() stops the unwind here, and Stop() still returns. Without it a // panic in any sweeper takes down the whole single-binary server for // every tenant. Must be `defer`-called directly in the goroutine for // recover() to catch the panic. func (s *Server) recoverSweeper(name string) { if r := recover(); r != nil { slog.Error("background sweeper panicked", "sweeper", name, "panic", r, "stack", string(debug.Stack())) } } // Stop waits for all background goroutines started via goAsync to finish // AND drains the rate-limiter cleanup goroutines spawned at construction // time (BUG-851). Safe to call multiple times. Should be called before // Store.Close() so in-flight DB writes don't race a closed connection // (or worse, the SQLite -wal/-shm file removal in t.TempDir cleanup). func (s *Server) Stop() { // Signal long-running background loops (orphan GC, etc.) to exit. // Each loop registers itself on s.bg, so the Wait() below blocks // until they actually finish and any in-flight goroutines drain. s.stopOrphanGC() // Yjs op-log prune sweeper (TASK-1309). Same lifecycle pattern; // signals BEFORE Wait() so the goroutine sees the close and exits. s.stopOpLogGC() // Short-lived-credential reaper (PLAN-1933 DR-5 / TASK-1936). Same // lifecycle pattern; signal BEFORE Wait() so the goroutine exits. s.stopTokenReaper() // Soft-deleted-workspace hard-purge sweeper (TASK-1966). Same // lifecycle pattern; signal BEFORE Wait() so the goroutine exits. s.stopWorkspacePurgeSweeper() // SPEC-3 event outbox drain (TASK-2714). Same lifecycle pattern; an // in-flight delivery is tracked on s.bg and awaited below. s.stopOutboxDrain() // Item reminder tick (IDEA-2641). Same lifecycle pattern; an in-flight // pass is tracked on s.bg and awaited below. s.stopReminderTick() // MCP audit writer / sweeper run on s.bg too. Signal first so // the workers see the close BEFORE Wait() blocks; without the // signal Wait would hang forever on the writer's blocking // queue receive. s.stopMCPAuditWriter() // MCP session tracker (TASK-1120) runs its sweeper on s.bg too. // Order with the audit writer doesn't matter — both are // independent goroutines; we just need the close BEFORE Wait(). s.stopMCPSessionTracker() // Close the collab room manager BEFORE bg.Wait() so any in-flight // op-log GC sweep (TASK-1309) blocked on a per-item lock behind // an active Join can drain. collab.Close() tears down the Joins // (their WS readLoops return, runConn unwinds, itemLocks // release), which unblocks the GC's per-item PruneItemOpLogIfDormantBefore // call. Without this ordering, Stop() can deadlock: GC waits on // itemLock; Join holds itemLock until WS closes; WS only closes // when collab.Close() runs; collab.Close() only runs after // bg.Wait(); bg.Wait() never returns because GC is stuck. // Per Codex review of TASK-1309 [P2]. nil-safe: collab is optional. if s.collab != nil { s.collab.Close() } s.bg.Wait() // Watch/nudge bus (BUG-2651). Closed AFTER bg.Wait() so a background // producer cannot publish into a bus that is already tearing down; both // implementations are safe if one does anyway (MemoryBus finds no // subscribers, RedisBus finds a cancelled context and fails closed). // // This matters more than it did for MemoryBus, whose Close only dropped // channels: RedisBus holds a receive goroutine and a Redis subscription // from construction, so skipping it leaks both for the process's life // (Codex round 1 P2). Closing also closes every subscriber channel, which // is how a still-open SSE stream learns to unwind. if s.watchEvents != nil { s.watchEvents.Close() } s.rateLimiters.Stop() // nil-safe via the RateLimiters receiver guard } func New(s *store.Store) *Server { rl := NewRateLimiters() // PAD_DISABLE_RATE_LIMITS turns off ALL HTTP rate limiting when set to a // truthy value. It exists ONLY for the E2E harness (BUG-2089): every // Playwright test shares one loopback IP (127.0.0.1), so the auth limiter // (5 logins/min/IP) trips the moment a spec logs in a couple of browser // clients — collab-persistence.spec.ts logs in two per test. A nil // rateLimiters makes RateLimit() a pass-through (see middleware_ratelimit.go // line 379), and Stop() + the MCP path are already nil-safe. Never set this // in production or self-host; it's an explicit opt-in so it can't flip on // by accident. if disabled, _ := strconv.ParseBool(os.Getenv("PAD_DISABLE_RATE_LIMITS")); disabled { rl = nil } return &Server{ store: s, rateLimiters: rl, storageInfoCache: newStorageInfoCache(storageInfoTTL), } } // Init2FASecret loads the 2FA challenge signing key from platform_settings. // If no key exists (first run), a new random key is generated and persisted. // This must be called before the server handles requests so that challenge // tokens survive process restarts and work across multiple instances. func (s *Server) Init2FASecret() error { const settingKey = "2fa_challenge_secret" existing, err := s.store.GetPlatformSetting(settingKey) if err != nil { return fmt.Errorf("load 2FA secret: %w", err) } if existing != "" { decoded, err := base64.StdEncoding.DecodeString(existing) if err != nil { return fmt.Errorf("decode 2FA secret: %w", err) } s.twoFAChallengeSecret = decoded return nil } // First run — generate and persist a new secret. // Multiple instances may race here on a fresh database; after persisting, // re-read the winning value so all instances converge on the same key. secret, err := generateTwoFASecret() if err != nil { return err } encoded := base64.StdEncoding.EncodeToString(secret) if err := s.store.SetPlatformSetting(settingKey, encoded); err != nil { return fmt.Errorf("persist 2FA secret: %w", err) } // Re-read to pick up whichever instance won the race (upsert may have // been overwritten by a concurrent instance between our check and write). final, err := s.store.GetPlatformSetting(settingKey) if err != nil { return fmt.Errorf("re-read 2FA secret: %w", err) } decoded, err := base64.StdEncoding.DecodeString(final) if err != nil { return fmt.Errorf("decode 2FA secret after re-read: %w", err) } s.twoFAChallengeSecret = decoded slog.Info("initialized 2FA challenge signing key") return nil } // SetCloudMode enables cloud mode with the shared sidecar secret(s). // Accepts a comma-separated list of secrets for rotation support: // "new-key,old-key" — both are accepted for INBOUND calls from pad-cloud. // The OUTBOUND direction (pad → pad-cloud, see SetCloudSidecar) is // configured separately via PAD_CLOUD_OUTBOUND_SECRET or derived from the // last entry of this list — see cmd/pad/main.go for the resolution order. func (s *Server) SetCloudMode(secret string) { s.cloudMode = true for _, k := range strings.Split(secret, ",") { k = strings.TrimSpace(k) if k != "" { s.cloudSecrets = append(s.cloudSecrets, k) } } // Propagate to the email sender so transactional emails carry the // getpad.dev marketing footer (docs/brand.md §7) on Cloud installs. // Self-hosted deployments leave cloudMode false on the sender, keeping // outgoing mail neutral so operators can ship under their own brand. if s.email != nil { s.email.SetCloudMode(true) } } // CloudSidecar is the reverse pad → pad-cloud client interface. Concrete // implementation lives in internal/billing so server has no direct Stripe // dependency. Kept as an interface so tests can inject fakes without // spinning up a real HTTP server or touching Stripe. type CloudSidecar interface { // CancelCustomer asks pad-cloud to cancel every active Stripe subscription // for customerID and then delete the Stripe customer object. Used by // handleDeleteAccount to cascade account deletion through to Stripe billing // (TASK-690). // // Failure contract: any non-nil error means the caller MUST abort the // local delete. pad-cloud normalizes Stripe's "already gone" cases to a // 200 on its side (see pad-cloud stripe.go isStripeAlreadyGone), so // every error we see here is a real failure — transport, 4xx (ops // misconfig), or 5xx (upstream breakage). Continuing after an error // would wipe the user's StripeCustomerID while leaving the subscription // billing, which is exactly the regression TASK-690 exists to prevent. CancelCustomer(customerID string) error // GetBillingMetrics fetches an aggregated Stripe-derived snapshot from // pad-cloud's /admin/metrics/billing endpoint (active subs, MRR, ARR, // churn, cancellations). Used by handleAdminBillingStats to power the // admin Billing dashboard (TASK-827 / PLAN-825). // // Failure contract: returns an error on transport failure or non-200 // status. The admin handler treats any error as "degrade to local-only" // and surfaces the distinction in its response via cloud_unreachable — // it never propagates the upstream failure to the operator's browser. GetBillingMetrics() (*billing.BillingMetricsResponse, error) } // SetCloudSidecar installs the reverse pad → pad-cloud client. Called from // cmd/pad/main.go when PAD_CLOUD_SIDECAR_URL + PAD_CLOUD_SECRET are set. // When unset, handleDeleteAccount skips the Stripe cancel step (self-hosted // deploys that don't run a Stripe-backed sidecar have nothing to cascade). func (s *Server) SetCloudSidecar(c CloudSidecar) { s.cloudSidecar = c } // SetBillingAvailable marks this deployment as having Stripe Checkout wired // up. Called from cmd/pad/main.go when PAD_BILLING_AVAILABLE=true is set. // When false (the default), the session payload advertises billing_available=false // so the web UI hides Stripe CTAs rather than dead-ending at a 503. TASK-800. func (s *Server) SetBillingAvailable(v bool) { s.billingAvailable = v } // IsCloud reports whether the server is running in cloud mode. func (s *Server) IsCloud() bool { return s.cloudMode } // SetVersion stores the build version info for the health endpoint. func (s *Server) SetVersion(version, commit, buildTime string) { s.version = version s.commit = commit s.buildTime = buildTime } // SetBaseURL sets the public base URL used for generating shareable links. // // If the supplied URL has an unspecified bind-all host ("0.0.0.0", "::", // "[::]"), this logs a WARN: such a URL is the right thing to *bind* to // but the wrong thing to *send* to a recipient (their browser cannot // resolve 0.0.0.0 / :: as a connect target). Callers shipping email // links from such a deployment should set PAD_URL or PUBLIC_URL to the // real public hostname (e.g. https://app.getpad.dev). See BUG-899. func (s *Server) SetBaseURL(rawURL string) { s.baseURL = strings.TrimRight(rawURL, "/") if s.baseURL == "" { return } if u, err := url.Parse(s.baseURL); err == nil { switch u.Hostname() { case "", "0.0.0.0", "::": slog.Warn("server base URL has an unspecified host; emailed links (password reset, invites, share links) will not be reachable. Set PAD_URL or PUBLIC_URL to the deployment's public URL (e.g. https://app.getpad.dev).", "base_url", s.baseURL) } } } // specialUseTLDs is the (finite) set of reserved top-level names from the IANA // Special-Use Domain Names registry + related RFCs that never resolve to a // public web host, so an emailed link using one is undeliverable. Matched on // the final label so example.com (public) is allowed while foo.example / // foo.test / bar.internal / pad.home.arpa / x.onion / y.alt (non-public) are // rejected. Sources: RFC 6761 (localhost/invalid/test/example), RFC 6762 // (local), RFC 8375 (home.arpa) + the "arpa" infrastructure TLD, RFC 7686 // (onion), RFC 9476 (alt), and ICANN-reserved "internal". var specialUseTLDs = map[string]bool{ "localhost": true, "local": true, "internal": true, "invalid": true, "test": true, "example": true, "arpa": true, "onion": true, "alt": true, } // hasUsableBaseURL reports whether s.baseURL is a URL a verification-email // recipient on the public internet can actually reach. It supersets the // unreachable-host warning in SetBaseURL (BUG-899): the bind-all hosts // 0.0.0.0 / :: are the right thing to bind() to but the wrong thing to email, // and so is every other host an external recipient can't resolve or route to. // // This gates the cloud self-serve signup path (PLAN-1933 DR-6): creating an // UNVERIFIED user whose only way out of the write-lock is an emailed link we // can't deliver would strand them permanently. So "usable" is conservative — // anything not clearly a public web endpoint disqualifies self-serve signup: // // - scheme must be http/https (a browser can't follow ftp://, file://, …); // - no query/fragment (the link is built by concatenation, so either would // push the /verify-email/ route into the query/fragment); // - a present port must be a valid TCP port (1–65535); // - the host must be a valid public DNS FQDN, NOT a literal IP. A real cloud // verification endpoint is a hostname (Pad Cloud is app.getpad.dev); // bare-IP base URLs aren't used for emailed links, and exhaustively // enumerating every non-public IP range (loopback / private / CGNAT / // TEST-NET / 6to4 / benchmarking / reserved / … across IPv4 and IPv6) is a // losing game — so we require a hostname and fail closed on any IP literal. // // No usable base URL → no self-serve signup (registration stays closed rather // than minting a write-locked user). Self-host and admin/invitation signup are // unaffected either way. func (s *Server) hasUsableBaseURL() bool { if s.baseURL == "" { return false } u, err := url.Parse(s.baseURL) if err != nil { return false } if u.Scheme != "http" && u.Scheme != "https" { return false } // The verification link is built by concatenation (baseURL + // "/verify-email/" + token), so a base URL carrying a query or fragment // would push the route into the query/fragment and break the link. if u.RawQuery != "" || u.Fragment != "" { return false } host := u.Hostname() if host == "" { return false } // A present port must be a valid TCP port (1–65535); url.Parse accepts // out-of-range numeric ports that no client can actually connect to. if p := u.Port(); p != "" { n, perr := strconv.Atoi(p) if perr != nil || n < 1 || n > 65535 { return false } } // Reject literal IPs outright — a usable public verification endpoint is a // DNS hostname, and "is this IP publicly reachable" is not decidable from a // finite denylist. Fail closed on any IP. if _, aerr := netip.ParseAddr(host); aerr == nil { return false } // The host must be a syntactically-valid, multi-label public FQDN (rejects // malformed hosts like ".com", "foo..com", "-a.com" and special-use TLDs). return isPublicDNSName(host) } // isPublicDNSName reports whether host is a syntactically-valid public FQDN // (RFC 1123 labels) whose TLD is neither a special-use reserved name nor // all-numeric. Empty labels, over-length labels, and invalid characters are // rejected. Assumes host is not an IP literal (that's handled before this is // called). Punycode/IDN TLDs (xn--…) are accepted since they aren't all-digit. func isPublicDNSName(host string) bool { host = strings.ToLower(strings.TrimSuffix(host, ".")) if host == "" || len(host) > 253 { return false } labels := strings.Split(host, ".") if len(labels) < 2 { return false } for _, l := range labels { if !isValidDNSLabel(l) { return false } } tld := labels[len(labels)-1] if specialUseTLDs[tld] || isAllDigits(tld) { return false } return true } // isValidDNSLabel reports whether l is a valid RFC 1123 hostname label: // 1–63 chars of [a-z0-9-], not starting or ending with a hyphen. func isValidDNSLabel(l string) bool { if len(l) == 0 || len(l) > 63 { return false } if l[0] == '-' || l[len(l)-1] == '-' { return false } for i := 0; i < len(l); i++ { c := l[i] if !((c >= 'a' && c <= 'z') || (c >= '0' && c <= '9') || c == '-') { return false } } return true } // isAllDigits reports whether s is non-empty and entirely ASCII digits. A // public TLD is never all-numeric (RFC 3696), so an all-digit final label // signals a malformed host rather than a reachable name. func isAllDigits(s string) bool { if s == "" { return false } for i := 0; i < len(s); i++ { if s[i] < '0' || s[i] > '9' { return false } } return true } // emailConfigured reports whether this instance can actually SEND an emailed // link — a sender is wired AND the public base URL is usable. This is the // DR-6 gate for cloud email self-registration: the sender-only check (s.email // != nil) is insufficient because link generation also needs a reachable // public base URL (see handleRegister's verification email + BUG-899). func (s *Server) emailConfigured() bool { return s.email != nil && s.hasUsableBaseURL() } // SetEventBus attaches an event bus for real-time SSE streaming. func (s *Server) SetEventBus(bus events.EventBus) { s.events = bus } // SetWatchEventsBus attaches the watch/nudge notification bus consumed by // GET /api/v1/events/stream (TASK-2533). Nil-checked by every producer and // by the stream handler, so a server constructed without one (e.g. a test // that doesn't exercise watches) still serves every other endpoint. func (s *Server) SetWatchEventsBus(bus watchevents.Bus) { s.watchEvents = bus } // SetSessionPresence attaches the live-session registry read by // GET /api/v1/sessions and written by GET /api/v1/events/stream // (PLAN-2558 S1). Nil-checked at both ends, so a server constructed // without one still streams events — it just can't answer "who is // listening?", and says so with a 503 rather than an empty list (see // handleListSessions). func (s *Server) SetSessionPresence(p SessionPresence) { s.sessionPresence = p } // SetRedisHealth attaches the Redis reachability prober (BUG-2727). Nil // on a deployment with no Redis, in which case /api/v1/health/ready reports no // redis block at all — the honest shape, since "no Redis configured" and // "Redis configured and down" are different states and a `false` would // merge them. // // It does NOT make readiness depend on Redis; see RedisHealth's doc // comment for why that is deliberate. func (s *Server) SetRedisHealth(h *RedisHealth) { s.redisHealth = h } // SetCollabRoomManager attaches a Yjs collab RoomManager (PLAN-1248). // When set, the /api/v1/collab/{itemID} WebSocket endpoint hands new // connections to the manager for op-log replay + fan-out. When nil, // the endpoint exists but answers 503 — that's intentional so a // self-host build that wants the editor without collab can leave // this unwired without surfacing surprise behaviour. func (s *Server) SetCollabRoomManager(rm *collab.RoomManager) { s.collab = rm } // SetWebhookDispatcher attaches a webhook dispatcher for outgoing // notifications. Delivery goroutines are routed through s.goAsync so they're // tracked on s.bg — Server.Stop() waits for in-flight deliveries (closing the // BUG-842 shutdown race where a detached delivery writes to a closed store) // and inherits goAsync's panic recovery (BUG-2011). func (s *Server) SetWebhookDispatcher(d *webhooks.Dispatcher) { if d != nil { d.SetSpawn(s.goAsync) } s.webhooks = d } // SetEmailSender attaches a transactional email sender. // The apiKey is stored separately for deriving the unsubscribe HMAC secret. // // If the server is already in cloud mode when this is called (i.e. // SetCloudMode ran before email config arrived from main.go), propagate // the flag so the new sender adds the getpad.dev marketing footer to // outgoing emails. Without this, the cloud-mode flag would silently // fail to take effect when callers wired email and cloud mode in // either order. func (s *Server) SetEmailSender(e *email.Sender, apiKey ...string) { s.email = e // A sender wired here comes from an out-of-band source (env vars at startup). // Mark it so reconfigureEmail leaves it in place when platform settings carry // no key — env is the deployment baseline, not something the admin UI disables. s.emailEnvConfigured = e != nil if len(apiKey) > 0 { s.emailAPIKey = apiKey[0] } if s.cloudMode && s.email != nil { s.email.SetCloudMode(true) } } // SetCORSOrigins configures allowed CORS origins (comma-separated). func (s *Server) SetCORSOrigins(origins string) { s.corsOrigins = origins } // SetAttachments wires the attachment storage Registry that the upload // and download handlers use. Pass maxBytes = 0 to keep the // defaultAttachmentMaxBytes ceiling (25 MiB). func (s *Server) SetAttachments(reg *attachments.Registry, maxBytes int64) { s.attachments = reg s.attachmentMaxBytes = maxBytes } // SetImageProcessor wires the image processor that the upload handler // uses to derive thumbnail variants (TASK-878). Optional — without it // uploads still succeed but no thumbnails are generated; the // download handler's variant fallback path returns the original blob. // The capabilities endpoint reflects whichever processor is wired. func (s *Server) SetImageProcessor(p attachments.Processor) { s.imageProcessor = p } // markUploadInFlight increments the in-flight counter for a content // hash. Returns a release func the caller MUST defer; the release // decrements and removes the entry once it hits zero. Used by the // upload handler to fence Put + CreateAttachment against orphan-GC // blob deletions of the same hash. // // Increment + map-store + decrement + delete all run under one // mutex so a concurrent uploadInFlight call can't observe a stale // "0" between the last release-decrement and the next-upload // increment. The earlier sync.Map version split increment from // LoadOrStore-then-atomic-add and missed that window (Codex P1 on // PR #307 round 2). func (s *Server) markUploadInFlight(hash string) func() { s.inFlightHashesMu.Lock() if s.inFlightHashes == nil { s.inFlightHashes = make(map[string]int64) } s.inFlightHashes[hash]++ s.inFlightHashesMu.Unlock() return func() { s.inFlightHashesMu.Lock() defer s.inFlightHashesMu.Unlock() s.inFlightHashes[hash]-- if s.inFlightHashes[hash] <= 0 { delete(s.inFlightHashes, hash) } } } // uploadInFlight reports whether any upload is currently materializing // a blob with the given hash. The orphan GC consults this before // deleting a blob — if an upload just finished Put but hasn't // inserted the row yet, GC must NOT reclaim the blob. func (s *Server) uploadInFlight(hash string) bool { s.inFlightHashesMu.Lock() defer s.inFlightHashesMu.Unlock() return s.inFlightHashes[hash] > 0 } // SetImportBundleMaxBytes overrides the default 2 GiB cap on a // single workspace import bundle. Set to 0 to fall back to the // default. Wired from PAD_IMPORT_BUNDLE_MAX_BYTES in cmd/pad/main.go // so operators with workspaces over 2 GiB can opt in without // recompiling. Larger caps trade memory headroom (one blob in // flight at a time, ≤25 MiB) for a longer import wall-clock. func (s *Server) SetImportBundleMaxBytes(n int64) { s.importBundleMaxBytes = n } // SetImportArtifactMaxBytes overrides the default 1 MiB cap on a single // playbook/convention artifact import. Set to 0 to fall back to the // default. Wired from PAD_IMPORT_ARTIFACT_MAX_BYTES in cmd/pad/main.go. func (s *Server) SetImportArtifactMaxBytes(n int64) { s.importArtifactMaxBytes = n } // SetSecureCookies enables the Secure flag on all cookies. func (s *Server) SetSecureCookies(secure bool) { s.secureCookies = secure } // countResumeGap records a resume this instance refused to serve, for the gaps // the BUS never sees (BUG-2731, codex round 12). // // The buses count their own coverage-based refusals, which is where most of // them happen. A cursor we could not even PARSE never reaches a bus, so // without this the counters undercount exactly the resyncs an operator is most // likely to be asked about — a client stuck in a resync loop because it keeps // sending a cursor nobody can read. No double counting: these paths return // before any bus call. // // Split by stream, matching the two counters' separate identities. func (s *Server) countResumeGap(activity bool) { if s.metrics == nil { return } if activity { s.metrics.EventResumeGapsTotal.Inc() return } s.metrics.WatchResumeGapsTotal.Inc() } // countMidStreamResync records a client told MID-STREAM that it missed // events, on a connection that stayed open (BUG-2730). // // Deliberately NOT countResumeGap: that counter's population is resumes this // instance could not serve, and existing alerts are written against it. // Widening it in place would have changed what those alerts measure without // changing their name, and a mixed-version fleet would report two populations // under one metric for the length of a rollout. func (s *Server) countMidStreamResync(activity bool) { if s.metrics == nil { return } if activity { s.metrics.EventMidstreamResyncsTotal.Inc() return } s.metrics.WatchMidstreamResyncsTotal.Inc() } // SetMetrics attaches Prometheus metrics to the server. // Must be called before the first request is served. // // Side effect (TASK-961): when both metrics AND the OAuth server are // wired, this also attaches the OAuth-active-tokens callback collector // and the revocation TTL observer. Order-independent — both // SetMetrics and SetOAuthServer call wireOAuthMetricsObserver, which // no-ops until both prerequisites are present. func (s *Server) SetMetrics(m *metrics.Metrics) { s.metrics = m s.wireOAuthMetricsObserver() s.wireStreamGauge() } // wireOAuthMetricsObserver attaches the OAuth metrics that need both // the metrics registry AND the OAuth server: the active-tokens // callback collector (reads via the store) and the per-revocation // TTL observer (fires from internal/oauth/storage.go on every // access-token family revocation). // // Idempotent — re-registering the same collector would panic via // prometheus.MustRegister, so we guard with a flag. Setting the // observer multiple times is harmless (just replaces the function // pointer). // // Why this lives on Server rather than in cmd/pad: it composes two // optional Server fields whose set-order isn't guaranteed by the // boot sequence, and centralizing the wiring here keeps the cmd/pad // startup path declarative ("set X, set Y") without an explicit // "now wire the cross-cut" call. func (s *Server) wireOAuthMetricsObserver() { if s.metrics == nil || s.oauthServer == nil { return } if !s.oauthMetricsWired { s.metrics.RegisterOAuthActiveTokensCollector(s.store.CountActiveOAuthAccessTokens) s.oauthMetricsWired = true } s.oauthServer.Storage().SetRevocationObserver(func(kind string, ttl time.Duration) { s.metrics.OAuthTokenRevocationsTotal.WithLabelValues(kind).Inc() s.metrics.OAuthTokenTTLSeconds.Observe(ttl.Seconds()) }) } // SetMetricsToken configures the static bearer token required to scrape // /metrics. When empty (the default), /metrics is exposed only to loopback // callers so a self-hosted Prometheus on the same host keeps working // without config — but LAN/internet scrapes are refused. A non-empty // token requires "Authorization: Bearer " regardless of source. func (s *Server) SetMetricsToken(token string) { s.metricsToken = strings.TrimSpace(token) } // metricsAuth gates the /metrics endpoint. See SetMetricsToken for the // policy. Uses constant-time comparison to avoid leaking the configured // token via response timing. func (s *Server) metricsAuth(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { if s.metricsToken == "" { // No token configured → loopback-only access. if !requestIsLoopback(r) { writeError(w, http.StatusForbidden, "forbidden", "/metrics is restricted to loopback when PAD_METRICS_TOKEN is unset") return } next.ServeHTTP(w, r) return } const prefix = "Bearer " authHeader := r.Header.Get("Authorization") if !strings.HasPrefix(authHeader, prefix) { w.Header().Set("WWW-Authenticate", `Bearer realm="metrics"`) writeError(w, http.StatusUnauthorized, "unauthorized", "Missing Bearer token for /metrics") return } given := strings.TrimSpace(strings.TrimPrefix(authHeader, prefix)) if subtle.ConstantTimeCompare([]byte(given), []byte(s.metricsToken)) != 1 { w.Header().Set("WWW-Authenticate", `Bearer realm="metrics"`) writeError(w, http.StatusUnauthorized, "unauthorized", "Invalid Bearer token for /metrics") return } next.ServeHTTP(w, r) }) } // SetSSELimits configures the streaming connection limits. A value of 0 // means unlimited for any of them. // // `global` and `perUser` bound BOTH stream endpoints together — // /api/v1/events and /api/v1/events/stream — through one admission gate, // because a held connection costs the same process resources whichever // one opened it (BUG-2726). `perWorkspace` bounds only /api/v1/events, // which is the only workspace-scoped one. func (s *Server) SetSSELimits(global, perWorkspace, perUser int) { s.sseMaxConnections = global s.sseMaxPerWorkspace = perWorkspace s.sseMaxPerUser = perUser // Updates the EXISTING gate rather than replacing it — see setLimits // for why replacing strands held slots and drifts the gauge. // admission() builds one on first use with the values just stored. s.admission().setLimits(global, perUser) s.wireStreamGauge() } // wireStreamGauge registers pad_stream_connections_active as a // scrape-time collector over the admission gate's total. Called from both // SetSSELimits and SetMetrics because either can land first. // // Registered ONCE PER METRICS INSTANCE. MustRegister panics on a // duplicate and both callers can fire, so it cannot register every time; // a once-per-SERVER guard would be wrong in the other direction, since // SetMetrics can install a different registry and would leave it // silently missing the series. // // The closure reads s.admission() at scrape time rather than capturing // the gate, so it stays correct if the gate is ever rebuilt. func (s *Server) wireStreamGauge() { if s.metrics == nil || s.streamGaugeFor == s.metrics { return } s.streamGaugeFor = s.metrics s.metrics.RegisterStreamConnectionsCollector(func() int { return s.admission().heldTotal() }) } // admission returns the shared stream admission gate, constructing an // unbounded one on first use so a Server built without SetSSELimits (every // test that does not care about limits) still has a working gate rather // than a nil check at each call site. // // SetSSELimits reconfigures this gate in place rather than replacing it, // so a late call neither strands held slots nor over-grants capacity. // // That is damage limitation, NOT a live-reconfiguration feature (codex // round 9). SetSSELimits also writes plain Server fields that request // handlers read — sseMaxPerWorkspace among them — so calling it while // requests are in flight is a data race regardless of how careful this // gate is. It is config-time: call it before ListenAndServe. The in-place // update exists so that a test, or an embedder that does call it late, // does not silently blow past its own limit. func (s *Server) admission() *streamAdmission { s.admitOnce.Do(func() { if s.streamAdmit == nil { s.streamAdmit = newStreamAdmission(s.sseMaxConnections, s.sseMaxPerUser) // The lazily-built gate needs the gauge too, or a server // wired with metrics but never given explicit limits would // export a permanently-zero pad_stream_connections_active // while serving streams — the same shape of lie as an // unregistered metric reporting 0. s.wireStreamGauge() } }) return s.streamAdmit } // SetTrustedProxies configures which direct TCP peers are allowed to set // X-Real-IP / X-Forwarded-For on incoming requests. Accepts a comma- // separated list of CIDRs or bare IPs (e.g. "10.0.0.0/8, 172.16.0.0/12"). // When empty (the default), proxy headers are ignored entirely — the // actual TCP peer address is used for rate limiting, the bootstrap // loopback check, and audit logging. func (s *Server) SetTrustedProxies(spec string) { s.trustedProxyCIDRs = ParseTrustedProxyCIDRs(spec) } // SetIPChangeEnforce controls how the auth middleware reacts when a // session's binding (client IP OR User-Agent hash) changes mid-lifetime: // - mode == "strict": revoke the session and reject the request (the token // is treated as possibly stolen). Covers BOTH the IP and the UA signal — // one flag arms the whole session-binding enforcement. // - anything else (default): log to the audit log, update the stored IP, // and let the request through. Strict mode breaks legitimate mobility // (mobile roaming, VPN toggles for IP; browser/WebView updates for UA) so // it is opt-in for high-sensitivity deployments via the // PAD_IP_CHANGE_ENFORCE env var. See handleSessionIPChange / // handleSessionUAChange for the per-signal semantics. func (s *Server) SetIPChangeEnforce(mode string) { s.ipChangeEnforceStrict = strings.EqualFold(strings.TrimSpace(mode), "strict") } // reconfigureEmail reads email settings from the platform_settings table // and updates (or creates) the email sender. Called after admin settings change. func (s *Server) reconfigureEmail() { apiKey, _ := s.store.GetPlatformSetting(settingMailerooAPIKey) fromAddr, _ := s.store.GetPlatformSetting(settingEmailFrom) fromName, _ := s.store.GetPlatformSetting(settingEmailFromName) if apiKey == "" { // No platform-settings key. If email was wired from env vars, leave that // sender in place — env is the deployment baseline. Otherwise the admin // cleared the only email config (e.g. selecting provider "None"), so tear // down the live sender: clearing the DB key alone left the running process // sending mail until restart (BUG-1890). if !s.emailEnvConfigured { s.email = nil s.emailAPIKey = "" } return } s.emailAPIKey = apiKey if s.email == nil { // Create a new sender from platform settings s.email = email.NewSender(apiKey, fromAddr, fromName, s.baseURL) } else { // Update existing sender s.email.Configure(apiKey, fromAddr, fromName, s.baseURL) } // Propagate cloud mode whichever way email was wired — see SetEmailSender // for the matching note. Configure() preserves cloudMode on existing // senders since SetCloudMode is independent; this branch covers the // fresh-NewSender path. if s.cloudMode { s.email.SetCloudMode(true) } } // InitEmailFromSettings loads email config from platform settings on startup, // merging with any env-var-based sender that was already attached. func (s *Server) InitEmailFromSettings() { s.reconfigureEmail() } func (s *Server) setupRouter() { r := chi.NewRouter() // Infrastructure middleware (applies to all routes including /metrics) // CapturePeerAddr MUST run before TrustedProxyRealIP so downstream code // that needs to verify the real TCP peer (e.g. the bootstrap loopback // check) can read the untampered value from request context even on // deployments with a trusted reverse proxy in front. r.Use(CapturePeerAddr) // RealIP is gated on PAD_TRUSTED_PROXIES. When unset (the default), proxy // headers are ignored and the real TCP peer address is used everywhere. // This prevents X-Forwarded-For spoofing from bypassing rate limits, the // bootstrap loopback check, or audit logs on direct-exposed deployments. r.Use(TrustedProxyRealIP(s.trustedProxyCIDRs)) r.Use(chimiddleware.RequestID) r.Use(StructuredLogger) if s.metrics != nil { r.Use(MetricsMiddleware(s.metrics)) } r.Use(chimiddleware.Recoverer) // Security headers (applies to all routes) r.Use(SecurityHeaders) if s.secureCookies { r.Use(StrictTransportSecurity) } // One CORS handler, shared. The /api/v1 group mounts it as ordinary // middleware; ValidatePath serves its REJECTION through the same // instance, because that rejection short-circuits above the group and // would otherwise answer a cross-origin caller without the CORS // headers its siblings carry — see the middleware's doc comment. corsMW := cors.Handler(cors.Options{ AllowedOrigins: parseCORSOrigins(s.corsOrigins), AllowedMethods: []string{"GET", "POST", "PATCH", "PUT", "DELETE", "OPTIONS"}, AllowedHeaders: []string{"Accept", "Authorization", "Content-Type", "X-CSRF-Token", "X-Share-Password", "X-Bootstrap-Token"}, // Credentials flag is gated on an operator explicitly listing // PAD_CORS_ORIGINS. The CLI uses Bearer tokens so the default // "no CORS_ORIGINS set" path doesn't need credential sharing; // leaving it off by default prevents cross-origin fetches from // a browser on a different site from piggy-backing cookies // on the victim's session. AllowCredentials: corsAllowCredentials(s.corsOrigins), MaxAge: 300, }) // Reject a request whose decoded path is not valid UTF-8 (or contains a // NUL) before anything routes it, so no handler can hand such a segment // to the store (BUG-2782). Placed at the END of the infrastructure block // deliberately: after StructuredLogger and MetricsMiddleware so a // rejection is still logged with a request id and still counted (as // route "unmatched", status 400) rather than being invisible to an // operator watching for a flood, and after SecurityHeaders so the 400 // carries them like every other response. r.Use(ValidatePath(corsMW)) // The query-string half of the same rule (BUG-2784). Separate middleware // rather than one combined check so the two rejections carry distinct // error codes; ordered after ValidatePath so a request that is bad in // both places is answered for its path, which is the more specific fault. r.Use(ValidateQuery(corsMW)) // MCP Streamable HTTP transport + OAuth discovery endpoints // (PLAN-943 TASK-950). Mounted outside the standard /api/v1 // auth-required group because: // // - /mcp uses Bearer auth via its own MCPBearerAuth middleware, // producing the spec-shape 401 + WWW-Authenticate that MCP // clients expect (the API-stack 401 envelope is JSON-only and // would fail Claude Desktop's discovery handshake). // - /.well-known/oauth-protected-resource and // /.well-known/oauth-authorization-server are public discovery // documents (RFC 9728 / RFC 8414); routing them through // TokenAuth+SessionAuth+RequireAuth would 401 unauth probes. // // No-op when SetMCPTransport hasn't been called or cloud mode is // off — see registerMCPRoutes for the gating. s.registerMCPRoutes(r) // OAuth 2.1 authorization-server flow endpoints (PLAN-943 // TASK-1025 sub-PR C). /oauth/{register,authorize,token, // authorize/decide} mounted alongside /mcp + /.well-known/*, // outside /api/v1's auth-required group. CSRF middleware runs // only on /api/* paths so /oauth/* is naturally exempt; the // consent-decision endpoint adds its own form-token check // using the existing __Host-pad_csrf cookie. // // SessionAuth runs in this group so /oauth/authorize can detect // whether the user is logged in via the __Host-pad_session // cookie. SessionAuth falls through gracefully when no cookie // is present (handlers see currentUser(r)==nil and redirect to // /login). RequireAuth is intentionally NOT used — /oauth/authorize // must be reachable anonymously to trigger the login redirect. // // RateLimit gates /oauth/register specifically (per Codex review // #372 round 2 — the DCR endpoint is open by RFC 7591 design, // but unlimited writes to oauth_clients are an obvious DoS // surface). The middleware short-circuits other /oauth/* paths // because they're either session-bound or PKCE-bound; explicit // per-endpoint limits arrive with TASK-959. // // No-op when SetOAuthServer hasn't been called or cloud mode is off. r.Group(func(r chi.Router) { r.Use(s.requireCloudMode) r.Use(s.SessionAuth) r.Use(s.RateLimit) s.registerOAuthRoutes(r) }) // Prometheus scrape endpoint — exempt from the standard auth/CSRF stack // (Prometheus can't present a session cookie or pass a CSRF header), but // gated by a dedicated static bearer token. Without the gate, any // unauthenticated caller on the network can read workspace counts, API // usage patterns, and — via label enumeration — user/workspace IDs. // // The gate runs in three layers: // 1. No PAD_METRICS_TOKEN → endpoint is open ONLY to loopback. Safe // default for self-hosters running Prometheus on the same box. // 2. PAD_METRICS_TOKEN set → "Authorization: Bearer " required. // Compared in constant time; empty/missing header → 401. // 3. In either case the SecurityHeaders / rate-limit / logging chain // already wraps this group from the outer r.Use() calls above. if s.metrics != nil { r.Group(func(r chi.Router) { r.Use(s.metricsAuth) r.Handle("/metrics", promhttp.HandlerFor(s.metrics.Registry, promhttp.HandlerOpts{})) }) } // All other routes — full middleware stack r.Group(func(r chi.Router) { r.Use(corsMW) r.Use(s.TokenAuth) r.Use(s.SessionAuth) r.Use(s.RateLimit) r.Use(s.CSRFProtect) r.Use(s.RequireAuth) // PLAN-1933 DR-4: block content-mutating requests from an // authenticated cloud user whose email is unverified. Mounted // AFTER RequireAuth so currentUser is already resolved; a no-op // on self-host and for verified / unauthenticated callers. The // method gate here covers the /api/v1 surface (session + PAT); // the collab GET-upgrade, the OAuth-provider flow, and the MCP // write path are gated at their own out-of-band mounts. r.Use(s.RequireVerifiedEmail) r.Use(jsonContentType) // SSE endpoint (outside jsonContentType middleware — but inherits auth) r.Get("/api/v1/events", s.handleSSE) // User-scoped watch/nudge event stream (TASK-2533, DOC-2479). // Unlike /api/v1/events above, this is NOT workspace-scoped — a // caller's watches and addressed pushes can span // every workspace they belong to. Lives alongside the other SSE // endpoint for the same "outside jsonContentType, inherits auth" // reason. r.Get("/api/v1/events/stream", s.handleWatchEventsStream) // Live-session presence (PLAN-2558 S1) — the READ side of the // registry the stream above writes. Mounted here beside the // endpoint it reports on rather than in the /api/v1 Route block // below: the two are one feature, and a reader asking "what // fills this list?" should find the answer on the adjacent // line. Self-scoped; see handleListSessions on why there is // deliberately no admin view. r.Get("/api/v1/sessions", s.handleListSessions) // WebSocket endpoint for Yjs-based collaborative editing on a // single item (PLAN-1248). Lives outside jsonContentType for // the same reason as SSE: the response is a WS upgrade, not // JSON. Inherits the auth middleware chain — handleCollab // then re-checks workspace access keyed on the item's // workspace ID (the URL only carries itemID). r.Get("/api/v1/collab/{itemID}", s.handleCollab) // API routes r.Route("/api/v1", func(r chi.Router) { r.Get("/health", s.handleHealth) r.Get("/health/live", s.handleHealthLive) r.Get("/health/ready", s.handleHealthReady) r.Get("/plan-limits", s.handleGetPlanLimits) // Public: billing page reads plan limits r.Get("/unsubscribe", s.handleUnsubscribe) // Public: email opt-out (HMAC-signed) // Server capabilities — public so the editor can fetch it // pre-login and gate per-format rotate / crop UI on the // processor's reach (TASK-878). The response is static for // the lifetime of the binary; clients can cache freely. r.Get("/server/capabilities", s.handleServerCapabilities) // Auth endpoints (exempt from auth middleware) r.Route("/auth", func(r chi.Router) { r.Get("/session", s.handleSessionCheck) r.Post("/bootstrap", s.handleBootstrap) r.Post("/register", s.handleRegister) r.Get("/check-username", s.handleCheckUsername) r.Post("/login", s.handleLogin) r.Post("/logout", s.handleLogout) r.Get("/me", s.handleGetCurrentUser) r.Patch("/me", s.handleUpdateCurrentUser) // Password reset r.Post("/forgot-password", s.handleForgotPassword) r.Post("/reset-password", s.handleResetPassword) // Localhost-only recovery escape hatch (self-host, non-cloud). r.Post("/local-reset", s.handleLocalReset) // Email verification (PLAN-1933 Wave 3b). Both are // enumeration-safe and rate-limited (middleware_ratelimit.go // reuses the PasswordReset bucket) and are already in // RequireVerifiedEmail's exempt list so an unverified user // can reach them to clear their own unverified state. r.Post("/verify-email", s.handleVerifyEmail) r.Post("/resend-verification", s.handleResendVerification) // Two-factor authentication r.Post("/2fa/setup", s.handleTOTPSetup) r.Post("/2fa/verify", s.handleTOTPVerify) r.Post("/2fa/disable", s.handleTOTPDisable) r.Post("/2fa/login-verify", s.handleTOTPLoginVerify) // Account management (GDPR) r.Post("/delete-account", s.handleDeleteAccount) r.Get("/export", s.handleExportAccount) // User-scoped API tokens r.Get("/tokens", s.handleListUserTokens) r.Post("/tokens", s.handleCreateUserToken) r.Delete("/tokens/{tokenID}", s.handleDeleteUserToken) r.Post("/tokens/{tokenID}/rotate", s.handleRotateUserToken) // Cloud: OAuth login/linking (called by pad-cloud sidecar, protected by cloud secret) r.Post("/oauth-login", s.handleOAuthLogin) r.Post("/oauth-link", s.handleOAuthLink) r.Post("/oauth-unlink", s.handleOAuthUnlink) // CLI browser-based auth flow r.Post("/cli/sessions", s.handleCreateCLIAuthSession) r.Get("/cli/sessions/{code}", s.handlePollCLIAuthSession) r.Post("/cli/sessions/{code}/approve", s.handleApproveCLIAuthSession) }) // Admin endpoints (admin-only, handlers check role internally) r.Route("/admin", func(r chi.Router) { r.Get("/settings", s.handleGetPlatformSettings) r.Patch("/settings", s.handleUpdatePlatformSettings) r.Post("/test-email", s.handleTestEmail) // Cloud sidecar endpoints — only exist in cloud mode. requireCloudMode // returns 404 outside cloud mode so a self-hosted deployment doesn't // expose "Cloud mode not configured" to unauthenticated probes. r.Group(func(r chi.Router) { r.Use(s.requireCloudMode) r.Post("/plan", s.handleSetPlan) // Cloud: sidecar sets user plans; also accessible to admins r.Post("/stripe-customer-id", s.handleSetStripeCustomerID) // Cloud: sidecar stores Stripe customer ID after checkout r.Get("/user-by-customer", s.handleGetUserByCustomerID) // Cloud: sidecar looks up user by Stripe customer ID r.Post("/stripe-event-processed", s.handleStripeEventProcessed) // Cloud: sidecar webhook idempotency (TASK-696) r.Post("/stripe-event-unmark", s.handleStripeEventUnmark) // Cloud: sidecar handler-failure rollback (TASK-736) r.Post("/payment-failed", s.handlePaymentFailed) // Cloud: sidecar forwards invoice.payment_failed to trigger email (TASK-712) // Admin Billing dashboard data (TASK-827 / PLAN-825). Proxies // pad-cloud's /admin/metrics/billing for Stripe-derived stats // (active subs, MRR, ARR, churn) and merges with local // users-table aggregates (customers_by_plan, new_signups_30d). // Always returns 200; degraded states (sidecar unreachable, // Stripe not configured) are surfaced as flags in the body. r.Get("/billing-stats", s.handleAdminBillingStats) }) // User management r.Get("/users", s.handleAdminListUsers) r.Get("/users/{userID}", s.handleAdminGetUser) r.Patch("/users/{userID}", s.handleAdminUpdateUser) r.Post("/users/{userID}/reset-password", s.handleAdminResetPassword) r.Get("/users/{userID}/workspaces", s.handleAdminGetUserWorkspaces) r.Get("/users/{userID}/detail", s.handleAdminGetUserDetail) r.Get("/users/{userID}/activity", s.handleAdminGetUserActivity) r.Get("/users/{userID}/metrics", s.handleAdminGetUserMetrics) r.Post("/users/{userID}/disable", s.handleAdminDisableUser) r.Post("/users/{userID}/enable", s.handleAdminEnableUser) r.Post("/users/{userID}/verify-email", s.handleAdminVerifyEmail) // Invitations r.Get("/invitations", s.handleAdminListInvitations) r.Post("/invitations/{invID}/resend", s.handleAdminResendInvitation) r.Delete("/invitations/{invID}", s.handleAdminDeleteInvitation) // Plan limits r.Get("/limits", s.handleAdminGetLimits) r.Patch("/limits", s.handleAdminUpdateLimits) // Platform stats r.Get("/stats", s.handleAdminStats) // MCP audit log — admin-only full-table view (TASK-960). // Powers /console/admin/mcp-audit. Per-connection // drilldown that users see for their own connections // lives at /api/v1/connected-apps/{id}/audit (registered // outside the admin group so non-admin users can read // their own). r.Get("/mcp-audit", s.handleAdminMCPAudit) }) // Audit log (admin-only) r.Get("/audit-log", s.handleAuditLog) // MCP per-connection audit (TASK-960). Owner-only via the // store query (user_id is one of the WHERE clauses); // returns the requesting user's own MCP activity for one // connection. The handler runs inside the standard // /api/v1 auth-required group, so unauthenticated callers // 401 here just like every other API endpoint. r.Get("/connected-apps/{id}/audit", s.handleMCPConnectionAudit) // Connected-apps management (TASK-954). Lists every // active OAuth grant chain the user has authorized // (Claude Desktop, Cursor, …) and lets them revoke one. // Cloud-mode-gated because OAuth is a cloud-only // surface — self-hosted deployments would always see // an empty list. r.Group(func(r chi.Router) { r.Use(s.requireCloudMode) r.Get("/connected-apps", s.handleListConnectedApps) r.Delete("/connected-apps/{id}", s.handleRevokeConnectedApp) // PLAN-1519 / TASK-1524 / IDEA-1517 §3: mutation // endpoints for the connections-page UI. Per-field // patches rather than a general PATCH for cleaner // error envelopes + audit shape. r.Patch("/connected-apps/{id}/name", s.handleRenameConnectedApp) r.Patch("/connected-apps/{id}/flags", s.handleUpdateConnectedAppFlags) r.Post("/connected-apps/{id}/workspaces", s.handleAddConnectedAppWorkspace) r.Delete("/connected-apps/{id}/workspaces/{slug}", s.handleRemoveConnectedAppWorkspace) }) // Templates r.Get("/templates", s.handleListTemplates) // Convention Library r.Get("/convention-library", s.handleConventionLibrary) // Playbook Library r.Get("/playbook-library", s.handlePlaybookLibrary) // Single library entry by title (conventions first, then playbooks). // TASK-1561 / PLAN-1560. r.Get("/library/entry", s.handleLibraryEntry) // URL import — fetch a remote page and return markdown. // Side-effect-free; the client decides what to do with the // markdown. See PLAN-1467 / TASK-1472 / internal/urlimport. r.Post("/import/url", s.handleImportURL) // Invitations (outside workspace scope) r.Post("/invitations/{code}/accept", s.handleAcceptInvitation) // Non-consuming invitation preview (BUG-1934). Public/pre-auth // (exempted in isPublicAPIPath) so the logged-out /join page can // prefill the invited email read-only and pick register-vs-login // mode. Always HTTP 200 + rate limited (see middleware_ratelimit.go) // so it can't be used to enumerate invite codes. r.Get("/invitations/{code}/preview", s.handlePreviewInvitation) // OAuth client public-info (PLAN-943 TASK-1027 sub-PR E). // Read-only consent-screen support for OAuth clients // registered via /oauth/register. Auth-required (inherits // RequireAuth from the parent group); cloud-mode-gated so // self-hosted deployments without an OAuth server don't // expose a hollow endpoint. Returns four non-sensitive // fields (client_id, client_name, logo_uri, redirect_uris) // — see handlers_oauth_clients.go for the full leak-surface // rationale. r.Group(func(r chi.Router) { r.Use(s.requireCloudMode) r.Get("/oauth/clients/{id}/public-info", s.handleOAuthClientPublicInfo) }) // Share link resolution (outside workspace scope, no auth required) r.Get("/s/{token}", s.handleResolveShareLink) // Share-link asset bytes (BUG-2389 2b / TASK-2637): rendered image // VARIANTS for attachments embedded in the shared content. Same // public/no-auth group; protected links gate on a short-lived // signed ref minted by handleResolveShareLink. Originals and // file downloads are out of scope by authorization. r.Get("/s/{token}/attachments/{attachmentID}", s.handleGetShareLinkAttachment) // Claim-code redemption (PLAN-1519 / TASK-1521 / IDEA-1517 §4). // POST /api/v1/oauth/claim with body {workspace, code} grants // the calling OAuth connection access to one workspace via a // stateless 6-digit HMAC code the user generated in the web // UI's "Connect project" modal. Auth: standard /api/v1 chain // (TokenAuth + RequireAuth); the handler itself short-circuits // the side effect when the caller isn't an OAuth grant (PAT / // CLI session) and 412s when the claim secret isn't wired. r.Post("/oauth/claim", s.handleOAuthClaim) // Workspaces r.Route("/workspaces", func(r chi.Router) { r.Get("/", s.handleListWorkspaces) r.Post("/", s.handleCreateWorkspace) r.Post("/import", s.handleImportWorkspace) r.Put("/reorder", s.handleReorderWorkspaces) // Soft-delete recovery (PLAN-1969 / TASK-1970). Both live // OUTSIDE the /{slug} RequireWorkspaceAccess subrouter // because that middleware resolves only LIVE workspaces // (deleted_at IS NULL) and would 404 a soft-deleted one // before the handler ran. The static "/deleted" segment is // registered before the /{slug} param route so chi matches // it exactly (static beats param); it lists the caller's own // deleted-but-restorable workspaces. "/{slug}/restore" // resolves the soft-deleted row itself and enforces // owner-only authz inside the handler. r.Get("/deleted", s.handleListDeletedWorkspaces) r.Post("/{slug}/restore", s.handleRestoreWorkspace) r.Route("/{slug}", func(r chi.Router) { r.Use(s.RequireWorkspaceAccess) r.Get("/", s.handleGetWorkspace) r.Patch("/", s.handleUpdateWorkspace) r.Delete("/", s.handleDeleteWorkspace) r.Get("/export", s.handleExportWorkspace) // Import a single playbook/convention artifact (Markdown // + YAML frontmatter) into this workspace. Editor+ gate // is enforced inside the handler against the destination // collection. r.Post("/import-artifact", s.handleImportArtifact) // Activity (workspace level) r.Get("/activity", s.handleListWorkspaceActivity) // Claim-code generation + smart suppression (PLAN-1519 // / TASK-1525 / IDEA-1517 §4). Inherits // RequireWorkspaceAccess so any member can pull a code // for any workspace they belong to — membership IS // the consent. See handlers_claim_code.go. r.Get("/claim-code", s.handleWorkspaceClaimCode) // Documents (v1 — will be replaced by items in Phase 2) r.Route("/documents", func(r chi.Router) { r.Get("/", s.handleListDocuments) r.Post("/", s.handleCreateDocument) r.Route("/{docID}", func(r chi.Router) { r.Get("/", s.handleGetDocument) r.Patch("/", s.handleUpdateDocument) r.Delete("/", s.handleDeleteDocument) r.Post("/restore", s.handleRestoreDocument) // Versions r.Get("/versions", s.handleListVersions) r.Get("/versions/{versionID}", s.handleGetVersion) // Activity (document level) r.Get("/activity", s.handleListDocumentActivity) }) }) // Collections (v2) r.Route("/collections", func(r chi.Router) { r.Get("/", s.handleListCollections) r.Post("/", s.handleCreateCollection) r.Route("/{collSlug}", func(r chi.Router) { r.Get("/", s.handleGetCollection) r.Patch("/", s.handleUpdateCollection) r.Delete("/", s.handleDeleteCollection) // Items within collection r.Get("/items", s.handleListCollectionItems) r.Post("/items", s.handleCreateItem) // Pairs with /items-index — server-side checkbox // progress so the collection page can render // list/board/table progress badges without // fetching item content (TASK-1349). r.Get("/checkbox-progress", s.handleCollectionCheckboxProgress) // Child-item completion progress for any collection // (BUG-1509). Same visibility/guest-grant semantics // as /plans-progress but collection-generic. r.Get("/child-progress", s.handleCollectionChildrenProgress) // Collection grants r.Get("/grants", s.handleListCollectionGrants) r.Post("/grants", s.handleCreateCollectionGrant) r.Delete("/grants/{grantID}", s.handleDeleteCollectionGrant) r.Get("/share-links", s.handleListCollectionShareLinks) r.Post("/share-links", s.handleCreateCollectionShareLink) // Saved views within collection r.Get("/views", s.handleListViews) r.Post("/views", s.handleCreateView) r.Route("/views/{viewID}", func(r chi.Router) { r.Patch("/", s.handleUpdateView) r.Delete("/", s.handleDeleteView) }) }) }) // Plans progress r.Get("/plans-progress", s.handlePlansProgress) // Skinny-projection cross-collection items list for the // local-first read model bootstrap (PLAN-1343 / TASK-1344). // Lives at workspace level — sibling to /plans-progress // and /starred — so the path can't ever collide with an // item slug under /items/{itemSlug}. r.Get("/items-index", s.handleListItemsIndex) // Delta-fetch sibling of /items-index: returns rows // where seq > since, including tombstones, so a // local-first read-model client can resume without // re-downloading the whole index (PLAN-1343 / TASK-1354). r.Get("/items-changes", s.handleListItemsChanges) // User grants (all grants for a specific user in this workspace) r.Get("/users/{userID}/grants", s.handleListUserGrants) // Starred items r.Get("/starred", s.handleListStarredItems) // Distinct tags across the workspace (with item counts) r.Get("/tags", s.handleListTags) // Items (cross-collection, v2) r.Get("/items", s.handleListItems) // Bulk mutation (TASK-1668). Static segment must be // registered before the /items/{itemSlug} param route // so "bulk" isn't captured as an item slug. r.Post("/items/bulk", s.handleBulkItems) r.Route("/items/{itemSlug}", func(r chi.Router) { r.Get("/", s.handleGetItem) r.Patch("/", s.handleUpdateItem) r.Delete("/", s.handleDeleteItem) r.Post("/restore", s.handleRestoreItem) r.Post("/move", s.handleMoveItem) // Cross-workspace copy PREFLIGHT (PLAN-2357 / // TASK-2364). Reports what a copy into another // workspace would carry, drop and need, and // leaves no trace a copy would have left — see // handlers_items_copy_preflight.go for the exact // scope of that guarantee. POST because the // request carries a body (destination + override // map), not because it mutates. The mutating // sibling lands at /copy in TASK-2365. r.Post("/copy/preflight", s.handleCopyItemPreflight) // Cross-workspace copy, the MUTATION (PLAN-2357 / // TASK-2365). Same request shape as the preflight // above; with archive_source it is the move. Post- // commit fanout is asymmetric — see // handlers_items_copy.go. Registered after the more // specific /copy/preflight, though chi's trie makes // the order immaterial. r.Post("/copy", s.handleCopyItem) // Export a single playbook/convention item as a // portable artifact (Markdown + YAML frontmatter). // Gated by per-item visibility, not the workspace- // export owner gate — a viewer who can see the item // may export it. r.Get("/export", s.handleExportItemArtifact) r.Get("/versions", s.handleListItemVersions) r.Get("/versions/{versionID}", s.handleGetItemVersion) r.Post("/versions/{versionID}/restore", s.handleRestoreItemVersion) r.Get("/activity", s.handleListItemActivity) r.Get("/links", s.handleGetItemLinks) r.Post("/links", s.handleCreateItemLink) r.Get("/comments", s.handleListComments) r.Post("/comments", s.handleCreateComment) r.Get("/timeline", s.handleListItemTimeline) r.Get("/children", s.handleGetItemChildren) r.Get("/progress", s.handleGetItemProgress) r.Get("/backlinks", s.handleGetItemBacklinks) r.Get("/tasks", s.handleGetItemChildren) // deprecated alias r.Get("/grants", s.handleListItemGrants) r.Post("/grants", s.handleCreateItemGrant) r.Delete("/grants/{grantID}", s.handleDeleteItemGrant) r.Get("/share-links", s.handleListItemShareLinks) r.Post("/share-links", s.handleCreateItemShareLink) // Stars r.Get("/star", s.handleGetItemStarStatus) r.Post("/star", s.handleStarItem) r.Delete("/star", s.handleUnstarItem) // Watches (TASK-2533): durable per-item subscriptions // for the padd event-stream / plugin-monitor nudge // pipeline. `pad watch ` / `pad watch remove `. r.Post("/watch", s.handleCreateWatch) r.Delete("/watch", s.handleDeleteWatch) // Push (IDEA-2544 Phase 1): transient, self-addressed // human→harness dispatch over the SAME watch-events // bus/stream — no durable row, see handlePushToItem's // doc comment. `pad push -m "message"`. r.Post("/push", s.handlePushToItem) // Reminders (IDEA-2641 / GitHub #1010): the // fire-at-an-instant primitive. Arming lives under // the item because a reminder is meaningless without // one; the lifecycle verbs live at the workspace // level below, addressed by reminder id, because an // acknowledgement is about the reminder rather than // about the item it names. r.Get("/reminders", s.handleListItemReminders) r.Post("/reminders", s.handleCreateItemReminder) }) // Links (v2) r.Delete("/links/{linkID}", s.handleDeleteItemLink) // Share links (workspace-scoped management) r.Delete("/share-links/{linkID}", s.handleDeleteShareLink) r.Get("/share-links/{linkID}/views", s.handleShareLinkViews) // Comments (v2) r.Route("/comments/{commentID}", func(r chi.Router) { r.Patch("/", s.handleUpdateComment) r.Delete("/", s.handleDeleteComment) r.Post("/replies", s.handleCreateReply) r.Post("/reactions", s.handleAddReaction) r.Delete("/reactions/{emoji}", s.handleRemoveReaction) }) // Role Board (cross-collection role-based view) r.Get("/roles/board", s.handleRoleBoard) r.Put("/roles/board/reorder", s.handleRoleBoardReorder) r.Put("/roles/board/lane-order", s.handleRoleBoardLaneReorder) // Agent Roles r.Route("/agent-roles", func(r chi.Router) { r.Get("/", s.handleListAgentRoles) r.Post("/", s.handleCreateAgentRole) r.Route("/{roleID}", func(r chi.Router) { r.Get("/", s.handleGetAgentRole) r.Patch("/", s.handleUpdateAgentRole) r.Delete("/", s.handleDeleteAgentRole) }) }) // Attachments // POST /attachments — upload (TASK-871) // GET /attachments/{attachmentID} — serve blob (TASK-872, supports ?variant=) // HEAD /attachments/{attachmentID} — metadata only (TASK-877 file-chip enrichment) // POST /attachments/{attachmentID}/transform — server-side rotate/crop (TASK-879/880) // // chi does not auto-route HEAD to the GET handler, so the // editor's HEAD probe for size + MIME has to be registered // explicitly. The handler short-circuits the streaming // path on HEAD; http.ServeContent already strips the body // on the seekable path. r.Post("/attachments", s.handleUploadAttachment) r.Get("/attachments", s.handleListWorkspaceAttachments) r.Get("/attachments/{attachmentID}", s.handleGetAttachment) r.Head("/attachments/{attachmentID}", s.handleGetAttachment) r.Post("/attachments/{attachmentID}/transform", s.handleTransformAttachment) r.Delete("/attachments/{attachmentID}", s.handleDeleteWorkspaceAttachment) // Storage usage summary for Settings → Storage and other // quota-aware UI surfaces (TASK-881). Cached behind a // short TTL — see handleGetWorkspaceStorageUsage. r.Get("/storage/usage", s.handleGetWorkspaceStorageUsage) // Webhooks r.Route("/webhooks", func(r chi.Router) { r.Get("/", s.handleListWebhooks) r.Post("/", s.handleCreateWebhook) r.Route("/{webhookID}", func(r chi.Router) { r.Delete("/", s.handleDeleteWebhook) r.Post("/test", s.handleTestWebhook) }) }) // API Tokens r.Route("/tokens", func(r chi.Router) { r.Get("/", s.handleListTokens) r.Post("/", s.handleCreateToken) r.Delete("/{tokenID}", s.handleDeleteToken) }) // Members r.Route("/members", func(r chi.Router) { r.Get("/", s.handleListMembers) r.Post("/invite", s.handleInviteMember) r.Delete("/invitations/{invID}", s.handleCancelInvitation) r.Delete("/{userID}", s.handleRemoveMember) r.Patch("/{userID}", s.handleUpdateMemberRole) r.Get("/{userID}/collection-access", s.handleGetMemberCollectionAccess) r.Put("/{userID}/collection-access", s.handleSetMemberCollectionAccess) }) // Me — current user's effective workspace context (role, // collection access, grants). Open to any principal admitted // by RequireWorkspaceAccess (members + guests). r.Get("/me", s.handleGetMe) // Dashboard (v2) // Reminder lifecycle, addressed by reminder id rather // than by item: an acknowledgement is about the reminder, // and an item can carry several. Permission is still the // ITEM's — see resolveReminderForWrite. r.Patch("/reminders/{reminderID}", s.handleRearmReminder) r.Post("/reminders/{reminderID}/ack", s.handleAckReminder) r.Delete("/reminders/{reminderID}", s.handleDeleteReminder) r.Get("/dashboard", s.handleGetDashboard) // Workspace graph — {nodes, edges} for the 3D // graph view (PLAN-1730 / TASK-1731). Active // items by default; ?include_terminal=true for // the full history. r.Get("/graph", s.handleGetWorkspaceGraph) // Project report — windowed throughput/flow/status // stats (PLAN-1628 / TASK-1630). r.Get("/report", s.handleGetReport) // Per-user Insights layout prefs (PLAN-1628 / TASK-1634). r.Get("/report/layout", s.handleGetReportLayout) r.Put("/report/layout", s.handleSaveReportLayout) // Project intelligence reads — next/standup/changelog // (PLAN-1888 / TASK-1894). Mirror `pad project // next|standup|changelog` (cmd/pad/main.go) — KEEP IN // SYNC, see handlers_project_intel.go's doc comments. // The MCP HTTP transport's dispatchProjectNext/Standup/ // Changelog (internal/mcp/dispatch_http_project.go) // proxy directly to these three handlers (TASK-1916), // so they need no separate sync-keeping. r.Get("/next", s.handleGetProjectNext) r.Get("/standup", s.handleGetProjectStandup) r.Get("/changelog", s.handleGetProjectChangelog) // Agent bootstrap (PLAN-1377 / TASK-1379) — single // round-trip that returns workspace + user + // collections + always-on conventions + roles + // playbook metadata + dashboard + recent activity. // Replaces the four /pad context-loading calls the // skill used to make. Same shape via the MCP // surfaces in TASK-1380. r.Get("/agent/bootstrap", s.handleGetBootstrap) // Playbook surface (PLAN-1377 / TASK-1382) — list / // show / run for first-class invokable procedures. // run is side-effect-free: it parses args per the // playbook's declared spec and returns the body + // bound args. The agent (skill or MCP-driven) // executes the body; the server does not. r.Get("/playbooks", s.handleListPlaybooks) r.Get("/playbooks/{ref}", s.handleShowPlaybook) r.Post("/playbooks/{ref}/run", s.handleRunPlaybook) // Incremental sync — returns items changed since a timestamp r.Get("/changes", s.handleGetChanges) }) }) // Search r.Get("/search", s.handleSearch) // My watches (TASK-2533), cross-workspace — mirrors // /auth/tokens' shape for a user-scoped-not-workspace-scoped // resource. `pad watch list`. Create/delete are per-item and // live under /workspaces/{ws}/items/{itemSlug}/watch instead // (they need the item's workspace context to resolve the // ref/slug the CLI's positional arg names). r.Get("/watches", s.handleListWatches) // MCP tool-surface descriptor (PLAN-1888 / TASK-1891). Serves // the catalog JSON (the nine env.Catalog tools + per-action // read_only flags) for the browser-side WebMCP layer to build // tool descriptors. Inside the authed group so it inherits // TokenAuth/SessionAuth/CSRFProtect/RequireAuth — same-origin // session/token only, NOT the bearer-gated /mcp infra path. // The handler nil-checks toolSurfaceJSON: 404 when the // serializer hasn't been injected (mirrors the SetMCPTransport // gating). Exposes only catalog descriptors — no route table, // handler internals, or other server state. r.Get("/mcp/tool-surface", s.handleMCPToolSurface) }) // Cross-workspace wiki-link resolver (IDEA-1492). Resolves // `[[workspace::REF]]` links emitted by the markdown renderer to // the canonical item URL via a 302 redirect. Lives outside /api/v1 // because rendered HTML hrefs target user-facing paths, not API // endpoints. Registered at the outer group level so chi matches // these URLs ahead of the catch-all SPA handler. ACL check matches // existing workspace-access semantics — 404 (not 403) on no-access // so we don't leak whether a workspace exists. // // URL shape: `/-/r/{workspace}/{ref}` — the leading `-/r/` prefix // is structurally impossible to collide with any user-namespace // URL because username slugs require a leading letter (slugify // rule), so no existing or future page route under // /{username}/... can shadow this resolver, and no collection // slug under /{u}/{ws}/{coll}/... can intercept it // (slug grammar also requires letter-led). This replaces the // earlier `/{username}/{workspace}/ref/{ref}` shape that risked // collision with collection slugs named "ref" on pre-existing // data (Codex round-2 P1.4 — picked Option B over a migration // because the feature is unshipped, the new shape is more // defensive, and the only cost is a frontend emit-shape change). r.Get("/-/r/{workspace}/{ref}", s.handleResolveCrossWorkspaceRef) }) // end r.Group (full middleware stack) s.router = r } // SetWebUI sets the embedded web UI filesystem for serving the SPA. func (s *Server) SetWebUI(fsys fs.FS) { s.webFS = fsys s.ensureRouter() s.router.Handle("/*", s.spaHandler()) } func (s *Server) spaHandler() http.Handler { fileServer := http.FileServer(http.FS(s.webFS)) indexHTML, err := fs.ReadFile(s.webFS, "index.html") if err != nil { // Embedded web UI is missing — fail fast instead of silently // serving blank HTML to every request. This indicates a broken // build, so the server should refuse to start. panic(fmt.Sprintf("spaHandler: failed to read embedded index.html: %v", err)) } return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { path := r.URL.Path if strings.HasPrefix(path, "/api/") { http.NotFound(w, r) return } cleanPath := strings.TrimPrefix(path, "/") if cleanPath != "" { if _, err := fs.Stat(s.webFS, cleanPath); err == nil { if strings.Contains(path, "/immutable/") { w.Header().Set("Cache-Control", "public, max-age=31536000, immutable") } else { w.Header().Set("Cache-Control", "no-cache") } fileServer.ServeHTTP(w, r) return } } // Generate per-request nonce for inline script CSP nonce := generateCSPNonce() // Inject nonce into inline