From fab32e50102378676958ec8f8407f825e3b10ac0 Mon Sep 17 00:00:00 2001 From: Asher Cohen Date: Mon, 27 Jul 2026 12:34:40 -0700 Subject: [PATCH 1/3] perf(db): guard activity overlap ratios Deploy the canonical view update through migration 0059 so time-disjoint activity pairs are rejected by cheap interval predicates while preserving the existing overlap semantics. --- docs/production-incident-baseline.md | 50 +++ drizzle/0059_v_activity_positive_overlap.sql | 318 ++++++++++++++++++ drizzle/_views/01_v_activity.sql | 2 + drizzle/meta/_journal.json | 7 + .../activity-overlap-plan.integration.test.ts | 100 ++++++ 5 files changed, 477 insertions(+) create mode 100644 drizzle/0059_v_activity_positive_overlap.sql create mode 100644 src/db/activity-overlap-plan.integration.test.ts diff --git a/docs/production-incident-baseline.md b/docs/production-incident-baseline.md index f2fe69c4ff..5fb400aadd 100644 --- a/docs/production-incident-baseline.md +++ b/docs/production-incident-baseline.md @@ -18916,3 +18916,53 @@ 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. +- **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..ff987deef0 --- /dev/null +++ b/drizzle/0059_v_activity_positive_overlap.sql @@ -0,0 +1,318 @@ +-- Canonical definition of the fitness.v_activity view. +-- This file is the source definition for fresh databases, local test schemas, +-- and future forward migrations that need to update the deployed view. +-- +-- To change v_activity: edit THIS file and add a forward migration when the +-- deployed view definition must change. +-- Git merge conflicts here force developers to reconcile concurrent changes. + +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 a + LEFT JOIN fitness.provider_priority pp ON pp.provider_id = a.provider_id + LEFT JOIN LATERAL ( + SELECT dp2.priority + FROM fitness.device_priority 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 + ) 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 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 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 t + INNER JOIN fitness.activity a ON a.id = t.id + WHERE t.provider_id = 'apple_health' + AND NOT EXISTS ( + SELECT 1 + FROM fitness.activity 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 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 t +), +pairs AS ( + SELECT c1.id AS id1, c2.id AS id2 + FROM clusterable c1 + JOIN clusterable c2 + ON c1.user_id = c2.user_id + AND c1.id < c2.id + AND c1.started_at < c2.ended_at + AND c2.started_at < c1.ended_at + 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 + 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 + ) +), +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 e + JOIN clusters c ON c.activity_id = e.a + 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 fg + JOIN ranked r ON r.id = fg.activity_id + ORDER BY fg.group_id, 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 fg + JOIN ranked r ON r.id = fg.activity_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 fg + JOIN effective_tombstoned t ON t.id = fg.activity_id + GROUP BY fg.group_id +), +tombstoned_groups AS ( + SELECT DISTINCT fg.group_id + FROM final_groups fg + JOIN effective_tombstoned t ON t.id = fg.activity_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 fg2 JOIN ranked r ON r.id = fg2.activity_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 fg2 JOIN ranked r ON r.id = fg2.activity_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 fg2 JOIN ranked r ON r.id = fg2.activity_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(key, value) + FROM ( + SELECT key, value, ROW_NUMBER() OVER (PARTITION BY key ORDER BY r.prio ASC) AS rn + FROM final_groups fg2 + JOIN ranked r ON r.id = fg2.activity_id, + LATERAL jsonb_each(COALESCE(r.raw, '{}'::jsonb)) + WHERE fg2.group_id = b.group_id + ) sub WHERE rn = 1 + ) AS raw, + (SELECT array_agg(DISTINCT r.provider_id ORDER BY r.provider_id) + FROM final_groups fg2 JOIN ranked r ON r.id = fg2.activity_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 fg2 JOIN ranked r ON r.id = fg2.activity_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 fg2 + WHERE fg2.group_id = b.group_id) AS member_activity_ids, + absent_source_links.absent_source_external_ids + FROM best_per_group b + JOIN group_bounds bounds ON bounds.group_id = b.group_id + LEFT JOIN absent_source_links ON absent_source_links.group_id = b.group_id + WHERE NOT EXISTS ( + SELECT 1 FROM tombstoned_groups 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 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..ff987deef0 100644 --- a/drizzle/_views/01_v_activity.sql +++ b/drizzle/_views/01_v_activity.sql @@ -152,6 +152,8 @@ pairs AS ( JOIN clusterable c2 ON c1.user_id = c2.user_id AND c1.id < c2.id + AND c1.started_at < c2.ended_at + AND c2.started_at < c1.ended_at CROSS JOIN LATERAL ( SELECT EXTRACT(EPOCH FROM ( 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..c0378651f6 --- /dev/null +++ b/src/db/activity-overlap-plan.integration.test.ts @@ -0,0 +1,100 @@ +import { sql } from "drizzle-orm"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; +import { TEST_USER_ID } from "./schema/core.ts"; +import { setupTestDatabase, type TestContext } from "./test-helpers.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 findPairsPlan(value: unknown): Record | undefined { + if (isRecord(value) && value["Subplan Name"] === "CTE pairs") { + return value; + } + + if (Array.isArray(value)) { + for (const child of value) { + const match = findPairsPlan(child); + if (match) return match; + } + } else if (isRecord(value)) { + for (const child of Object.values(value)) { + const match = findPairsPlan(child); + if (match) return match; + } + } + + return undefined; +} + +describe("activity overlap query plan", () => { + let testCtx: TestContext; + + beforeAll(async () => { + testCtx = await setupTestDatabase(); + await testCtx.db.execute( + sql`INSERT INTO fitness.provider (id, name, user_id) + VALUES ('wahoo', 'Wahoo', ${TEST_USER_ID}) + ON CONFLICT DO NOTHING`, + ); + await testCtx.db.execute( + 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 () => { + const rows = await testCtx.db.execute<{ member_activity_ids: 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()).toEqual([1, 2]); + + const explainRows = await testCtx.db.execute<{ "QUERY PLAN": unknown }>( + sql`EXPLAIN (FORMAT JSON) + SELECT count(*) + FROM fitness.v_activity + WHERE user_id = ${TEST_USER_ID}`, + ); + const pairsPlan = findPairsPlan(explainRows[0]?.["QUERY PLAN"]); + const joinFilter = pairsPlan?.["Join Filter"]; + + expect(joinFilter).toEqual(expect.any(String)); + expect(joinFilter).toContain("(c1.started_at < c2.ended_at)"); + expect(joinFilter).toContain("(c2.started_at < c1.ended_at)"); + }); +}); From 5c23379daae153065242284839205c0301358a43 Mon Sep 17 00:00:00 2001 From: Asher Cohen Date: Mon, 27 Jul 2026 12:44:58 -0700 Subject: [PATCH 2/3] fix(db): harden overlap migration checks Format the immutable migration for SQLFluff and make the database regression schema-validated and planner-shape tolerant. --- docs/production-incident-baseline.md | 6 + drizzle/0059_v_activity_positive_overlap.sql | 326 +++++++++++------- .../activity-overlap-plan.integration.test.ts | 54 +-- 3 files changed, 228 insertions(+), 158 deletions(-) diff --git a/docs/production-incident-baseline.md b/docs/production-incident-baseline.md index 27759ca2d4..a8407c4b7f 100644 --- a/docs/production-incident-baseline.md +++ b/docs/production-incident-baseline.md @@ -18968,6 +18968,12 @@ Drizzle schema and runtime Zod schemas. Findings and remediations: 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. diff --git a/drizzle/0059_v_activity_positive_overlap.sql b/drizzle/0059_v_activity_positive_overlap.sql index ff987deef0..8ae5c97789 100644 --- a/drizzle/0059_v_activity_positive_overlap.sql +++ b/drizzle/0059_v_activity_positive_overlap.sql @@ -1,29 +1,28 @@ --- Canonical definition of the fitness.v_activity view. --- This file is the source definition for fresh databases, local test schemas, --- and future forward migrations that need to update the deployed view. --- --- To change v_activity: edit THIS file and add a forward migration when the --- deployed view definition must change. --- Git merge conflicts here force developers to reconcile concurrent changes. +-- 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 a - LEFT JOIN fitness.provider_priority pp ON pp.provider_id = a.provider_id + 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 dp2 - WHERE dp2.provider_id = a.provider_id + 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 + ORDER BY LENGTH(dp2.source_name_pattern) DESC LIMIT 1 - ) dp ON true - WHERE a.provider_absent_at IS NULL - AND a.deleted_at IS NULL + ) AS dp ON true + WHERE + a.provider_absent_at IS null + AND a.deleted_at IS null ), + tombstoned AS ( SELECT a.id, @@ -35,15 +34,17 @@ tombstoned AS ( a.ended_at, a.provider_absent_at, COALESCE( - NULLIF(trim(a.raw->>'sourceName'), ''), - NULLIF(trim(a.source_name), '') + NULLIF(TRIM(a.raw ->> 'sourceName'), ''), + NULLIF(TRIM(a.source_name), '') ) AS subsource - FROM fitness.activity a - WHERE a.provider_absent_at IS NOT NULL - AND a.deleted_at IS NULL - AND a.external_id IS NOT NULL + 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, @@ -55,7 +56,7 @@ effective_tombstoned AS ( t.ended_at, t.provider_absent_at, t.subsource - FROM tombstoned t + FROM tombstoned AS t WHERE t.provider_id <> 'apple_health' UNION ALL SELECT @@ -68,57 +69,59 @@ effective_tombstoned AS ( t.ended_at, t.provider_absent_at, t.subsource - FROM tombstoned t - INNER JOIN fitness.activity a ON a.id = t.id - WHERE t.provider_id = 'apple_health' + 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 sib - WHERE sib.user_id = a.user_id + 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.deleted_at IS null AND sib.id <> a.id AND COALESCE( - NULLIF(trim(sib.raw->'metadata'->>'HKMetadataKeySyncIdentifier'), ''), + 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), ''), + NULLIF(TRIM(sib.raw ->> 'sourceName'), ''), + NULLIF(TRIM(sib.source_name), ''), '' ) ) = COALESCE( - NULLIF(trim(a.raw->'metadata'->>'HKMetadataKeySyncIdentifier'), ''), + 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), ''), + NULLIF(TRIM(a.raw ->> 'sourceName'), ''), + NULLIF(TRIM(a.source_name), ''), '' ) ) AND ( - sib.provider_absent_at IS NULL AND sib.deleted_at IS NULL + 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 + 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 + 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 + 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 + WHEN (a.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion') ~ '^[0-9]+$' + THEN (a.raw -> 'metadata' ->> 'HKMetadataKeySyncVersion')::bigint END, 0 ) @@ -127,6 +130,7 @@ effective_tombstoned AS ( ) ) ), + clusterable AS ( SELECT r.id, @@ -135,7 +139,7 @@ clusterable AS ( r.activity_type, r.started_at, COALESCE(r.ended_at, r.started_at + interval '1 hour') AS ended_at - FROM ranked r + FROM ranked AS r UNION ALL SELECT t.id, @@ -144,54 +148,86 @@ clusterable AS ( t.activity_type, t.started_at, COALESCE(t.ended_at, t.started_at + interval '1 hour') AS ended_at - FROM effective_tombstoned t + 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 c1.id AS id1, c2.id AS id2 - FROM clusterable c1 - JOIN clusterable c2 - ON c1.user_id = c2.user_id - AND c1.id < c2.id - AND c1.started_at < c2.ended_at - AND c2.started_at < c1.ended_at - 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 + 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 ( - SELECT id1 AS a, id2 AS b FROM pairs + SELECT + id1 AS a, + id2 AS b + FROM pairs UNION ALL - SELECT id2 AS a, id1 AS b FROM pairs + SELECT + id2 AS a, + id1 AS b + FROM pairs ), -clusters(activity_id, group_id, depth) AS ( - SELECT id, id::text, 0 FROM clusterable + +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 e - JOIN clusters c ON c.activity_id = e.a + 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 + 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, @@ -203,24 +239,26 @@ best_per_group AS ( r.ended_at, r.source_name, r.prio - FROM final_groups fg - JOIN ranked r ON r.id = fg.activity_id - ORDER BY fg.group_id, r.prio ASC, r.id ASC + 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 fg - JOIN ranked r ON r.id = fg.activity_id + 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( + JSONB_AGG( + JSONB_BUILD_OBJECT( 'providerId', t.provider_id, 'externalId', t.external_id, 'memberActivityId', t.id::text, @@ -229,15 +267,17 @@ absent_source_links AS ( ) ORDER BY t.provider_id, t.id ) AS absent_source_external_ids - FROM final_groups fg - JOIN effective_tombstoned t ON t.id = fg.activity_id + 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 fg - JOIN effective_tombstoned t ON t.id = fg.activity_id + FROM final_groups AS fg + INNER JOIN effective_tombstoned AS t ON fg.activity_id = t.id ), + merged AS ( SELECT b.canonical_id, @@ -247,56 +287,76 @@ merged AS ( bounds.started_at, bounds.ended_at, b.source_name, - (SELECT r.name FROM final_groups fg2 JOIN ranked r ON r.id = fg2.activity_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 fg2 JOIN ranked r ON r.id = fg2.activity_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 fg2 JOIN ranked r ON r.id = fg2.activity_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(key, value) - FROM ( - SELECT key, value, ROW_NUMBER() OVER (PARTITION BY key ORDER BY r.prio ASC) AS rn - FROM final_groups fg2 - JOIN ranked r ON r.id = fg2.activity_id, - LATERAL jsonb_each(COALESCE(r.raw, '{}'::jsonb)) - WHERE fg2.group_id = b.group_id - ) sub WHERE rn = 1 + ( + 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 fg2 JOIN ranked r ON r.id = fg2.activity_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 fg2 JOIN ranked r ON r.id = fg2.activity_id - WHERE fg2.group_id = b.group_id - AND r.external_id IS NOT NULL - AND r.external_id <> '' + ( + 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 fg2 - WHERE fg2.group_id = b.group_id) AS member_activity_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 b - JOIN group_bounds bounds ON bounds.group_id = b.group_id - LEFT JOIN absent_source_links ON absent_source_links.group_id = b.group_id + 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 tg WHERE tg.group_id = b.group_id + SELECT 1 FROM tombstoned_groups AS tg + WHERE tg.group_id = b.group_id ) ) + SELECT m.canonical_id AS id, m.provider_id, @@ -314,5 +374,5 @@ SELECT m.source_external_ids, m.member_activity_ids, m.absent_source_external_ids -FROM merged m +FROM merged AS m ORDER BY m.started_at DESC; diff --git a/src/db/activity-overlap-plan.integration.test.ts b/src/db/activity-overlap-plan.integration.test.ts index c0378651f6..3af6863e8b 100644 --- a/src/db/activity-overlap-plan.integration.test.ts +++ b/src/db/activity-overlap-plan.integration.test.ts @@ -1,7 +1,9 @@ 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", @@ -13,37 +15,33 @@ function isRecord(value: unknown): value is Record { return value !== null && typeof value === "object" && !Array.isArray(value); } -function findPairsPlan(value: unknown): Record | undefined { - if (isRecord(value) && value["Subplan Name"] === "CTE pairs") { - return value; +function collectPlanStrings(value: unknown): string[] { + if (Array.isArray(value)) { + return value.flatMap(collectPlanStrings); } - if (Array.isArray(value)) { - for (const child of value) { - const match = findPairsPlan(child); - if (match) return match; - } - } else if (isRecord(value)) { - for (const child of Object.values(value)) { - const match = findPairsPlan(child); - if (match) return match; - } + if (isRecord(value)) { + return Object.values(value).flatMap(collectPlanStrings); } - return undefined; + return typeof value === "string" ? [value] : []; } describe("activity overlap query plan", () => { - let testCtx: TestContext; + let testCtx: TestContext | undefined; beforeAll(async () => { testCtx = await setupTestDatabase(); - await testCtx.db.execute( + 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 testCtx.db.execute( + 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 @@ -66,11 +64,17 @@ describe("activity overlap query plan", () => { }); afterAll(async () => { - await testCtx.cleanup(); + await testCtx?.cleanup(); }); it("requires positive overlap for candidate pairs", async () => { - const rows = await testCtx.db.execute<{ member_activity_ids: string[] }>( + 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} @@ -84,17 +88,17 @@ describe("activity overlap query plan", () => { expect(rows.map((row) => row.member_activity_ids.length).sort()).toEqual([1, 2]); - const explainRows = await testCtx.db.execute<{ "QUERY PLAN": unknown }>( + 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 pairsPlan = findPairsPlan(explainRows[0]?.["QUERY PLAN"]); - const joinFilter = pairsPlan?.["Join Filter"]; + const planText = collectPlanStrings(explainRows[0]?.["QUERY PLAN"]).join(" "); - expect(joinFilter).toEqual(expect.any(String)); - expect(joinFilter).toContain("(c1.started_at < c2.ended_at)"); - expect(joinFilter).toContain("(c2.started_at < c1.ended_at)"); + expect(planText).toMatch(/c1\.started_at < c2\.ended_at|c2\.ended_at > c1\.started_at/); + expect(planText).toMatch(/c2\.started_at < c1\.ended_at|c1\.ended_at > c2\.started_at/); }); }); From dc8b1326b1117e2edb9bc19eda5f5d731c7bb260 Mon Sep 17 00:00:00 2001 From: Asher Cohen Date: Mon, 27 Jul 2026 13:19:40 -0700 Subject: [PATCH 3/3] fix(db): align overlap view definitions Keep the canonical pair implementation structurally aligned with the linted migration and make the plan regression alias-independent. --- drizzle/_views/01_v_activity.sql | 47 +++++++++++-------- .../activity-overlap-plan.integration.test.ts | 34 ++++++++++++-- 2 files changed, 58 insertions(+), 23 deletions(-) diff --git a/drizzle/_views/01_v_activity.sql b/drizzle/_views/01_v_activity.sql index ff987deef0..a59d54c4b0 100644 --- a/drizzle/_views/01_v_activity.sql +++ b/drizzle/_views/01_v_activity.sql @@ -146,32 +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 AND c1.started_at < c2.ended_at - AND c2.started_at < c1.ended_at - 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.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/src/db/activity-overlap-plan.integration.test.ts b/src/db/activity-overlap-plan.integration.test.ts index 3af6863e8b..e8f2f76d01 100644 --- a/src/db/activity-overlap-plan.integration.test.ts +++ b/src/db/activity-overlap-plan.integration.test.ts @@ -27,6 +27,33 @@ function collectPlanStrings(value: unknown): string[] { 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; @@ -86,7 +113,9 @@ describe("activity overlap query plan", () => { ORDER BY started_at`, ); - expect(rows.map((row) => row.member_activity_ids.length).sort()).toEqual([1, 2]); + expect( + rows.map((row) => row.member_activity_ids.length).sort((left, right) => left - right), + ).toEqual([1, 2]); const explainRows = await executeWithSchema( testCtx.db, @@ -98,7 +127,6 @@ describe("activity overlap query plan", () => { ); const planText = collectPlanStrings(explainRows[0]?.["QUERY PLAN"]).join(" "); - expect(planText).toMatch(/c1\.started_at < c2\.ended_at|c2\.ended_at > c1\.started_at/); - expect(planText).toMatch(/c2\.started_at < c1\.ended_at|c1\.ended_at > c2\.started_at/); + expect(hasBidirectionalOverlapGuards(planText)).toBe(true); }); });