diff --git a/docs/production-incident-baseline.md b/docs/production-incident-baseline.md index 628ef0b818..45c28cde14 100644 --- a/docs/production-incident-baseline.md +++ b/docs/production-incident-baseline.md @@ -18986,3 +18986,59 @@ Drizzle schema and runtime Zod schemas. Findings and remediations: deployed release is observed; if the token-limit branch occurs, confirm the UI returns the documented clean-grant retry guidance before resolving DOFEK-SERVER-5H. + +## 2026-07-27 — Global activity overlap expansion exhausted web DB pools + +- **Status:** Root cause confirmed from production query plans and direct + source fix prepared; merge, deployment, and production validation pending. +- **Symptoms:** Dashboard bursts reported + [DOFEK-SERVER-5K](https://east-bay-software.sentry.io/issues/DOFEK-SERVER-5K), + [DOFEK-SERVER-5M](https://east-bay-software.sentry.io/issues/DOFEK-SERVER-5M), + and + [DOFEK-SERVER-5N](https://east-bay-software.sentry.io/issues/DOFEK-SERVER-5N) + as failures of session, insights, and PMC SQL statements. +- **User impact:** Five session validations plus the insights and PMC SQL + calls could not acquire a database connection at 17:46 UTC. At 18:35 UTC, + insights and PMC failed while the surrounding dashboard batch took 47.7 + seconds to finish. +- **Runbook classification:** Following the + [loading-performance evidence gate](./performance/loading-performance-runbook.md#evidence-gate), + this was a request-time query-shape slowdown in PostgreSQL. The first fatal + line was `timeout exceeded when trying to connect`; the named slow family was + `fitness.v_activity`, and the evidence below excludes client blanking, + ClickHouse queueing, stale statistics, disk I/O, and lock contention. +- **Evidence:** Every Sentry event's causal error was + `timeout exceeded when trying to connect` at the web process's unchanged + ten-second pool acquisition boundary; the displayed SQL had not started. + Web request logs and `pg_stat_statements` show the pool was instead occupied + by `fitness.v_activity` queries that took 13.94–16.29 seconds. A read-only + `EXPLAIN (ANALYZE, BUFFERS)` of the exact insights query took 5.08 seconds + from shared buffers: the recursive view built 2,635 activity rows globally, + compared the same-user cross-product, rejected 6,941,004 pairs, and produced + 2,221 overlapping pairs before applying the request's user and date filters. + Current statistics were fresh, there were no lock waits or idle + transactions, and the server had 11 of 40 connections in use. +- **Root cause:** `fitness.v_activity` evaluated the expensive overlap-ratio + arithmetic for every same-user activity pair, including millions of + time-disjoint pairs. The view's recursive grouping prevents the outer + request predicates from bounding that global work. PostgreSQL documents CTE + evaluation and recursive query behavior: + . +- **Fix / mitigation:** Require strict positive interval overlap in the pair + join alongside the two 80% ratios. Both ratio branches already require + positive overlap, so the guards preserve deduplication, contained-activity, + cross-provider, and boundary-touching semantics while giving PostgreSQL + cheap necessary conditions for disjoint windows. The production read-only + benchmark returned the same 210 rows in 1.84 seconds, 64% faster. No pool + size, acquisition timeout, retry, cache, or database resource limit changed. +- **Validation:** A real-PostgreSQL regression fixture covers a contained pair + and a boundary-touching non-pair, then inspects the database's executable + pairs plan for both positive-overlap guards. Local execution is blocked + because Docker Desktop's filesystem is full and PostgreSQL cannot create an + isolated test database; GitHub integration CI remains the executable green + gate. PostgreSQL documents executable plan inspection with `EXPLAIN`: + . +- **Remaining risk / follow-up:** Merge and deploy through the normal release + path, repeat the exact production plan and dashboard request, confirm all + three Sentry issues remain quiet, and resolve them only after the deployed + query no longer exhausts a web process's five-connection pool. diff --git a/drizzle/0059_v_activity_positive_overlap.sql b/drizzle/0059_v_activity_positive_overlap.sql new file mode 100644 index 0000000000..8ae5c97789 --- /dev/null +++ b/drizzle/0059_v_activity_positive_overlap.sql @@ -0,0 +1,378 @@ +-- Immutable forward migration for positive-overlap guards in fitness.v_activity. +-- For future changes, edit drizzle/_views/01_v_activity.sql and add a new +-- forward migration. + +CREATE OR REPLACE VIEW fitness.v_activity AS +WITH RECURSIVE ranked AS ( + SELECT + a.*, + COALESCE(dp.priority, pp.priority, 100) AS prio + FROM fitness.activity AS a + LEFT JOIN fitness.provider_priority AS pp ON a.provider_id = pp.provider_id + LEFT JOIN LATERAL ( + SELECT dp2.priority + FROM fitness.device_priority AS dp2 + WHERE + dp2.provider_id = a.provider_id + AND a.source_name LIKE dp2.source_name_pattern + ORDER BY LENGTH(dp2.source_name_pattern) DESC + LIMIT 1 + ) AS dp ON true + WHERE + a.provider_absent_at IS null + AND a.deleted_at IS null +), + +tombstoned AS ( + SELECT + a.id, + a.user_id, + a.provider_id, + a.activity_type, + a.external_id, + a.started_at, + a.ended_at, + a.provider_absent_at, + COALESCE( + NULLIF(TRIM(a.raw ->> 'sourceName'), ''), + NULLIF(TRIM(a.source_name), '') + ) AS subsource + FROM fitness.activity AS a + WHERE + a.provider_absent_at IS NOT null + AND a.deleted_at IS null + AND a.external_id IS NOT null + AND a.external_id <> '' +), + +effective_tombstoned AS ( + SELECT + t.id, + t.user_id, + t.provider_id, + t.activity_type, + t.external_id, + t.started_at, + t.ended_at, + t.provider_absent_at, + t.subsource + FROM tombstoned AS t + WHERE t.provider_id <> 'apple_health' + UNION ALL + SELECT + t.id, + t.user_id, + t.provider_id, + t.activity_type, + t.external_id, + t.started_at, + t.ended_at, + t.provider_absent_at, + t.subsource + FROM tombstoned AS t + INNER JOIN fitness.activity AS a ON t.id = a.id + WHERE + t.provider_id = 'apple_health' + AND NOT EXISTS ( + SELECT 1 + FROM fitness.activity AS sib + WHERE + sib.user_id = a.user_id + AND sib.provider_id = 'apple_health' + AND sib.deleted_at IS null + AND sib.id <> a.id + AND COALESCE( + NULLIF(TRIM(sib.raw -> 'metadata' ->> 'HKMetadataKeySyncIdentifier'), ''), + 'time:' || sib.started_at::text || ':' || COALESCE(sib.ended_at::text, '') || ':' || COALESCE( + NULLIF(TRIM(sib.raw ->> 'sourceName'), ''), + NULLIF(TRIM(sib.source_name), ''), + '' + ) + ) = COALESCE( + NULLIF(TRIM(a.raw -> 'metadata' ->> 'HKMetadataKeySyncIdentifier'), ''), + 'time:' || a.started_at::text || ':' || COALESCE(a.ended_at::text, '') || ':' || COALESCE( + NULLIF(TRIM(a.raw ->> 'sourceName'), ''), + NULLIF(TRIM(a.source_name), ''), + '' + ) + ) + AND ( + sib.provider_absent_at IS null AND sib.deleted_at IS null + OR COALESCE( + CASE + WHEN (sib.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion') ~ '^[0-9]+$' + THEN (sib.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion')::bigint + END, + 0 + ) > COALESCE( + CASE + WHEN (a.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion') ~ '^[0-9]+$' + THEN (a.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion')::bigint + END, + 0 + ) + OR ( + COALESCE( + CASE + WHEN (sib.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion') ~ '^[0-9]+$' + THEN (sib.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion')::bigint + END, + 0 + ) = COALESCE( + CASE + WHEN (a.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion') ~ '^[0-9]+$' + THEN (a.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion')::bigint + END, + 0 + ) + AND sib.created_at > a.created_at + ) + ) + ) +), + +clusterable AS ( + SELECT + r.id, + r.user_id, + r.provider_id, + r.activity_type, + r.started_at, + COALESCE(r.ended_at, r.started_at + interval '1 hour') AS ended_at + FROM ranked AS r + UNION ALL + SELECT + t.id, + t.user_id, + t.provider_id, + t.activity_type, + t.started_at, + COALESCE(t.ended_at, t.started_at + interval '1 hour') AS ended_at + FROM effective_tombstoned AS t +), + +pair_metrics AS ( + SELECT + c1.id AS id1, + c2.id AS id2, + c1.provider_id AS provider_id1, + c2.provider_id AS provider_id2, + c1.activity_type AS activity_type1, + c2.activity_type AS activity_type2, + EXTRACT(EPOCH FROM ( + LEAST(c1.ended_at, c2.ended_at) - GREATEST(c1.started_at, c2.started_at) + )) AS overlap_seconds, + EXTRACT(EPOCH FROM ( + GREATEST(c1.ended_at, c2.ended_at) - LEAST(c1.started_at, c2.started_at) + )) AS union_seconds, + LEAST( + EXTRACT(EPOCH FROM (c1.ended_at - c1.started_at)), + EXTRACT(EPOCH FROM (c2.ended_at - c2.started_at)) + ) AS shorter_duration_seconds + FROM clusterable AS c1 + INNER JOIN clusterable AS c2 + ON + c1.user_id = c2.user_id + AND c1.id < c2.id + AND c1.started_at < c2.ended_at + AND c1.ended_at > c2.started_at +), + +pairs AS ( + SELECT + id1, + id2 + FROM pair_metrics + WHERE + overlap_seconds / NULLIF(union_seconds, 0) > 0.8 + OR ( + provider_id1 <> provider_id2 + AND activity_type1 = activity_type2 + AND overlap_seconds / NULLIF(shorter_duration_seconds, 0) > 0.8 + ) +), + +edges AS ( + SELECT + id1 AS a, + id2 AS b + FROM pairs + UNION ALL + SELECT + id2 AS a, + id1 AS b + FROM pairs +), + +clusters (activity_id, group_id, depth) AS ( + SELECT + id, + id::text, + 0 + FROM clusterable + UNION + SELECT + e.b, + c.group_id, + c.depth + 1 + FROM edges AS e + INNER JOIN clusters AS c ON e.a = c.activity_id + WHERE c.depth < 2 +), + +final_groups AS ( + SELECT + activity_id, + MIN(group_id) AS group_id + FROM clusters + GROUP BY activity_id +), + +best_per_group AS ( + SELECT DISTINCT ON (fg.group_id) + fg.group_id, + r.id AS canonical_id, + r.provider_id, + r.user_id, + r.activity_type, + r.started_at, + r.ended_at, + r.source_name, + r.prio + FROM final_groups AS fg + INNER JOIN ranked AS r ON fg.activity_id = r.id + ORDER BY fg.group_id ASC, r.prio ASC, r.id ASC +), + +group_bounds AS ( + SELECT + fg.group_id, + MIN(r.started_at) AS started_at, + MAX(r.ended_at) AS ended_at + FROM final_groups AS fg + INNER JOIN ranked AS r ON fg.activity_id = r.id + GROUP BY fg.group_id +), + +absent_source_links AS ( + SELECT + fg.group_id, + JSONB_AGG( + JSONB_BUILD_OBJECT( + 'providerId', t.provider_id, + 'externalId', t.external_id, + 'memberActivityId', t.id::text, + 'providerAbsentAt', t.provider_absent_at, + 'subsource', t.subsource + ) + ORDER BY t.provider_id, t.id + ) AS absent_source_external_ids + FROM final_groups AS fg + INNER JOIN effective_tombstoned AS t ON fg.activity_id = t.id + GROUP BY fg.group_id +), + +tombstoned_groups AS ( + SELECT DISTINCT fg.group_id + FROM final_groups AS fg + INNER JOIN effective_tombstoned AS t ON fg.activity_id = t.id +), + +merged AS ( + SELECT + b.canonical_id, + b.provider_id, + b.user_id, + b.activity_type, + bounds.started_at, + bounds.ended_at, + b.source_name, + ( + SELECT r.name FROM final_groups AS fg2 INNER JOIN ranked AS r ON fg2.activity_id = r.id + WHERE fg2.group_id = b.group_id AND r.name IS NOT null + ORDER BY r.prio ASC LIMIT 1 + ) AS name, + ( + SELECT r.notes FROM final_groups AS fg2 INNER JOIN ranked AS r ON fg2.activity_id = r.id + WHERE fg2.group_id = b.group_id AND r.notes IS NOT null + ORDER BY r.prio ASC LIMIT 1 + ) AS notes, + ( + SELECT r.timezone FROM final_groups AS fg2 INNER JOIN ranked AS r ON fg2.activity_id = r.id + WHERE fg2.group_id = b.group_id AND r.timezone IS NOT null + ORDER BY r.prio ASC LIMIT 1 + ) AS timezone, + ( + SELECT JSONB_OBJECT_AGG(sub.key, sub.value) + FROM ( + SELECT + raw_entry.key, + raw_entry.value, + ROW_NUMBER() OVER (PARTITION BY raw_entry.key ORDER BY r.prio ASC) AS rn + FROM final_groups AS fg2 + INNER JOIN ranked AS r ON fg2.activity_id = r.id, + LATERAL JSONB_EACH(COALESCE(r.raw, '{}'::jsonb)) AS raw_entry + WHERE fg2.group_id = b.group_id + ) AS sub + WHERE sub.rn = 1 + ) AS raw, + ( + SELECT ARRAY_AGG(DISTINCT r.provider_id ORDER BY r.provider_id) + FROM final_groups AS fg2 INNER JOIN ranked AS r ON fg2.activity_id = r.id + WHERE fg2.group_id = b.group_id + ) AS source_providers, + ( + SELECT + JSONB_AGG( + JSONB_BUILD_OBJECT( + 'providerId', r.provider_id, + 'externalId', r.external_id, + 'memberActivityId', r.id::text, + -- Preserve the per-member upstream app for grouped Apple Health rows. + 'subsource', COALESCE( + NULLIF(TRIM(r.raw ->> 'sourceName'), ''), + NULLIF(TRIM(r.source_name), '') + ) + ) + ORDER BY r.provider_id + ) + FROM final_groups AS fg2 INNER JOIN ranked AS r ON fg2.activity_id = r.id + WHERE + fg2.group_id = b.group_id + AND r.external_id IS NOT null + AND r.external_id <> '' + ) AS source_external_ids, + ( + SELECT ARRAY_AGG(fg2.activity_id ORDER BY fg2.activity_id) + FROM final_groups AS fg2 + WHERE fg2.group_id = b.group_id + ) AS member_activity_ids, + absent_source_links.absent_source_external_ids + FROM best_per_group AS b + INNER JOIN group_bounds AS bounds ON b.group_id = bounds.group_id + LEFT JOIN absent_source_links ON b.group_id = absent_source_links.group_id + WHERE NOT EXISTS ( + SELECT 1 FROM tombstoned_groups AS tg + WHERE tg.group_id = b.group_id + ) +) + +SELECT + m.canonical_id AS id, + m.provider_id, + m.user_id, + m.canonical_id AS primary_activity_id, + m.activity_type, + m.started_at, + m.ended_at, + m.source_name, + m.name, + m.notes, + m.timezone, + m.raw, + m.source_providers, + m.source_external_ids, + m.member_activity_ids, + m.absent_source_external_ids +FROM merged AS m +ORDER BY m.started_at DESC; diff --git a/drizzle/_views/01_v_activity.sql b/drizzle/_views/01_v_activity.sql index dc9655b47d..a59d54c4b0 100644 --- a/drizzle/_views/01_v_activity.sql +++ b/drizzle/_views/01_v_activity.sql @@ -146,30 +146,39 @@ clusterable AS ( COALESCE(t.ended_at, t.started_at + interval '1 hour') AS ended_at FROM effective_tombstoned t ), -pairs AS ( - SELECT c1.id AS id1, c2.id AS id2 +pair_metrics AS ( + SELECT + c1.id AS id1, + c2.id AS id2, + c1.provider_id AS provider_id1, + c2.provider_id AS provider_id2, + c1.activity_type AS activity_type1, + c2.activity_type AS activity_type2, + EXTRACT(EPOCH FROM ( + LEAST(c1.ended_at, c2.ended_at) - GREATEST(c1.started_at, c2.started_at) + )) AS overlap_seconds, + EXTRACT(EPOCH FROM ( + GREATEST(c1.ended_at, c2.ended_at) - LEAST(c1.started_at, c2.started_at) + )) AS union_seconds, + LEAST( + EXTRACT(EPOCH FROM (c1.ended_at - c1.started_at)), + EXTRACT(EPOCH FROM (c2.ended_at - c2.started_at)) + ) AS shorter_duration_seconds FROM clusterable c1 JOIN clusterable c2 ON c1.user_id = c2.user_id AND c1.id < c2.id - CROSS JOIN LATERAL ( - SELECT - EXTRACT(EPOCH FROM ( - LEAST(c1.ended_at, c2.ended_at) - GREATEST(c1.started_at, c2.started_at) - )) AS overlap_seconds, - EXTRACT(EPOCH FROM ( - GREATEST(c1.ended_at, c2.ended_at) - LEAST(c1.started_at, c2.started_at) - )) AS union_seconds, - LEAST( - EXTRACT(EPOCH FROM (c1.ended_at - c1.started_at)), - EXTRACT(EPOCH FROM (c2.ended_at - c2.started_at)) - ) AS shorter_duration_seconds - ) o - WHERE o.overlap_seconds / NULLIF(o.union_seconds, 0) > 0.8 + AND c1.started_at < c2.ended_at + AND c1.ended_at > c2.started_at +), +pairs AS ( + SELECT id1, id2 + FROM pair_metrics + WHERE overlap_seconds / NULLIF(union_seconds, 0) > 0.8 OR ( - c1.provider_id <> c2.provider_id - AND c1.activity_type = c2.activity_type - AND o.overlap_seconds / NULLIF(o.shorter_duration_seconds, 0) > 0.8 + provider_id1 <> provider_id2 + AND activity_type1 = activity_type2 + AND overlap_seconds / NULLIF(shorter_duration_seconds, 0) > 0.8 ) ), edges AS ( diff --git a/drizzle/meta/_journal.json b/drizzle/meta/_journal.json index a6b8e8f13e..057b436b1d 100644 --- a/drizzle/meta/_journal.json +++ b/drizzle/meta/_journal.json @@ -484,6 +484,13 @@ "when": 1785100000000, "tag": "0058_personal_experiment", "breakpoints": true + }, + { + "idx": 69, + "version": "7", + "when": 1785180698000, + "tag": "0059_v_activity_positive_overlap", + "breakpoints": true } ] } diff --git a/src/db/activity-overlap-plan.integration.test.ts b/src/db/activity-overlap-plan.integration.test.ts new file mode 100644 index 0000000000..e8f2f76d01 --- /dev/null +++ b/src/db/activity-overlap-plan.integration.test.ts @@ -0,0 +1,132 @@ +import { sql } from "drizzle-orm"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; +import { z } from "zod"; +import { TEST_USER_ID } from "./schema/core.ts"; +import { setupTestDatabase, type TestContext } from "./test-helpers.ts"; +import { executeWithSchema } from "./typed-sql.ts"; + +const activityIds = [ + "00000000-0000-4000-8000-000000000201", + "00000000-0000-4000-8000-000000000202", + "00000000-0000-4000-8000-000000000203", +] as const; + +function isRecord(value: unknown): value is Record { + return value !== null && typeof value === "object" && !Array.isArray(value); +} + +function collectPlanStrings(value: unknown): string[] { + if (Array.isArray(value)) { + return value.flatMap(collectPlanStrings); + } + + if (isRecord(value)) { + return Object.values(value).flatMap(collectPlanStrings); + } + + return typeof value === "string" ? [value] : []; +} + +function hasBidirectionalOverlapGuards(planText: string): boolean { + const directions: Array<{ endAlias: string; startAlias: string }> = []; + const comparisonPattern = + /\b([a-zA-Z_]\w*)\.(started_at|ended_at)\s*([<>])\s*([a-zA-Z_]\w*)\.(started_at|ended_at)\b/g; + + for (const match of planText.matchAll(comparisonPattern)) { + const leftAlias = match[1]; + const leftColumn = match[2]; + const operator = match[3]; + const rightAlias = match[4]; + const rightColumn = match[5]; + if (!(leftAlias && rightAlias) || leftAlias === rightAlias) continue; + + if (leftColumn === "started_at" && operator === "<" && rightColumn === "ended_at") { + directions.push({ startAlias: leftAlias, endAlias: rightAlias }); + } else if (leftColumn === "ended_at" && operator === ">" && rightColumn === "started_at") { + directions.push({ startAlias: rightAlias, endAlias: leftAlias }); + } + } + + return directions.some(({ startAlias, endAlias }) => + directions.some( + (candidate) => candidate.startAlias === endAlias && candidate.endAlias === startAlias, + ), + ); +} + +describe("activity overlap query plan", () => { + let testCtx: TestContext | undefined; + + beforeAll(async () => { + testCtx = await setupTestDatabase(); + await executeWithSchema( + testCtx.db, + z.object({}), + sql`INSERT INTO fitness.provider (id, name, user_id) + VALUES ('wahoo', 'Wahoo', ${TEST_USER_ID}) + ON CONFLICT DO NOTHING`, + ); + await executeWithSchema( + testCtx.db, + z.object({}), + sql`INSERT INTO fitness.activity ( + id, provider_id, user_id, external_id, activity_type, started_at, ended_at + ) VALUES + ( + ${activityIds[0]}::uuid, 'wahoo', ${TEST_USER_ID}, 'overlap-plan-a', 'cycling', + TIMESTAMPTZ '2026-01-10 10:00:00+00', + TIMESTAMPTZ '2026-01-10 11:00:00+00' + ), + ( + ${activityIds[1]}::uuid, 'wahoo', ${TEST_USER_ID}, 'overlap-plan-contained', 'cycling', + TIMESTAMPTZ '2026-01-10 10:05:00+00', + TIMESTAMPTZ '2026-01-10 10:55:00+00' + ), + ( + ${activityIds[2]}::uuid, 'wahoo', ${TEST_USER_ID}, 'overlap-plan-touching', 'cycling', + TIMESTAMPTZ '2026-01-10 11:00:00+00', + TIMESTAMPTZ '2026-01-10 12:00:00+00' + )`, + ); + }); + + afterAll(async () => { + await testCtx?.cleanup(); + }); + + it("requires positive overlap for candidate pairs", async () => { + if (!testCtx) { + throw new Error("Test database setup did not complete"); + } + + const rows = await executeWithSchema( + testCtx.db, + z.object({ member_activity_ids: z.array(z.string()) }), + sql`SELECT member_activity_ids::text[] AS member_activity_ids + FROM fitness.v_activity + WHERE user_id = ${TEST_USER_ID} + AND member_activity_ids && ARRAY[ + ${activityIds[0]}::uuid, + ${activityIds[1]}::uuid, + ${activityIds[2]}::uuid + ] + ORDER BY started_at`, + ); + + expect( + rows.map((row) => row.member_activity_ids.length).sort((left, right) => left - right), + ).toEqual([1, 2]); + + const explainRows = await executeWithSchema( + testCtx.db, + z.object({ "QUERY PLAN": z.unknown() }), + sql`EXPLAIN (FORMAT JSON) + SELECT count(*) + FROM fitness.v_activity + WHERE user_id = ${TEST_USER_ID}`, + ); + const planText = collectPlanStrings(explainRows[0]?.["QUERY PLAN"]).join(" "); + + expect(hasBidirectionalOverlapGuards(planText)).toBe(true); + }); +});