Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
97 changes: 70 additions & 27 deletions analytics/models/read_models/activity_power_curve.sql
Original file line number Diff line number Diff line change
Expand Up @@ -79,43 +79,86 @@ power_samples AS (
AND sensor.is_deleted = 0
),

sample_rate AS (
power_segments AS (
SELECT
current_sample.activity_id AS activity_id,
current_sample.user_id AS user_id,
current_sample.started_at AS started_at,
current_sample.recorded_at AS recorded_at,
next_sample.recorded_at AS next_recorded_at,
current_sample.power AS power,
current_sample.row_number AS row_number,
dateDiff('millisecond', current_sample.recorded_at, next_sample.recorded_at) / 1000.0 AS segment_seconds
FROM power_samples AS current_sample
INNER JOIN power_samples AS next_sample
ON next_sample.activity_id = current_sample.activity_id
AND toInt64(next_sample.row_number) = toInt64(current_sample.row_number) + 1
WHERE next_sample.recorded_at > current_sample.recorded_at
),

sample_gap_stats AS (
SELECT
activity_id,
greatest(
toInt32(round(
dateDiff('second', min(recorded_at), max(recorded_at))
/ nullIf(count() - 1, 0)
)),
1
) AS interval_s
FROM power_samples
greatest(5.0, quantileExact(0.5) (segment_seconds) * 2.0) AS max_continuous_gap_seconds
FROM power_segments
GROUP BY activity_id
HAVING count() > 1
),

duration_values AS (
SELECT arrayJoin([5, 15, 30, 60, 120, 180, 300, 420, 600, 1200, 1800, 3600, 5400, 7200]) AS duration_seconds
SELECT
duration_seconds,
toInt64(duration_seconds) * 1000 AS duration_milliseconds
FROM (
SELECT arrayJoin([5, 15, 30, 60, 120, 180, 300, 420, 600, 1200, 1800, 3600, 5400, 7200]) AS duration_seconds
)
),

duration_windows AS (
candidate_duration_windows AS (
SELECT
ps.activity_id AS activity_id,
ps.user_id AS user_id,
ps.started_at AS started_at,
start_sample.activity_id AS activity_id,
start_sample.user_id AS user_id,
start_sample.started_at AS started_at,
duration_values.duration_seconds AS duration_seconds,
greatest(1, toInt32(round(duration_values.duration_seconds / sr.interval_s))) AS window_samples,
(
ps.cumulative_sum - coalesce(prev_sample.cumulative_sum, 0)
) / toFloat64(greatest(1, toInt32(round(duration_values.duration_seconds / sr.interval_s)))) AS avg_power
FROM duration_values
CROSS JOIN power_samples AS ps
INNER JOIN sample_rate AS sr
ON sr.activity_id = ps.activity_id
LEFT JOIN power_samples AS prev_sample
ON prev_sample.activity_id = ps.activity_id
AND toInt64(prev_sample.row_number) = toInt64(ps.row_number) - toInt64(greatest(1, toInt32(round(duration_values.duration_seconds / sr.interval_s))))
WHERE toInt64(ps.row_number) >= greatest(1, toInt32(round(duration_values.duration_seconds / sr.interval_s)))
start_sample.recorded_at AS window_started_at,
max(window_sample.recorded_at) AS window_ended_at,
max(dateDiff('millisecond', start_sample.recorded_at, window_sample.recorded_at)) / 1000.0 AS elapsed_seconds,
max(segment.segment_seconds) AS max_gap_seconds,
sum(segment.power * segment.segment_seconds) / sum(segment.segment_seconds) AS avg_power
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
FROM power_samples AS start_sample
CROSS JOIN duration_values
Comment thread
Asherlc marked this conversation as resolved.
INNER JOIN power_samples AS window_sample
ON window_sample.activity_id = start_sample.activity_id
AND toInt64(window_sample.row_number) >= toInt64(start_sample.row_number)
INNER JOIN power_segments AS segment
ON segment.activity_id = start_sample.activity_id
AND toInt64(segment.row_number) >= toInt64(start_sample.row_number)
AND toInt64(segment.row_number) < toInt64(window_sample.row_number)
WHERE dateDiff('millisecond', start_sample.recorded_at, window_sample.recorded_at)
<= duration_values.duration_milliseconds
Comment thread
Asherlc marked this conversation as resolved.
GROUP BY
start_sample.activity_id,
start_sample.user_id,
start_sample.started_at,
start_sample.recorded_at,
window_sample.row_number,
duration_values.duration_seconds
),
Comment thread
Asherlc marked this conversation as resolved.

duration_windows AS (
SELECT
candidate_duration_windows.activity_id AS activity_id,
candidate_duration_windows.user_id AS user_id,
candidate_duration_windows.started_at AS started_at,
candidate_duration_windows.duration_seconds AS duration_seconds,
candidate_duration_windows.window_started_at AS window_started_at,
candidate_duration_windows.window_ended_at AS window_ended_at,
candidate_duration_windows.elapsed_seconds AS covered_seconds,
candidate_duration_windows.avg_power AS avg_power
FROM candidate_duration_windows
INNER JOIN sample_gap_stats AS gap_stats
ON gap_stats.activity_id = candidate_duration_windows.activity_id
WHERE candidate_duration_windows.elapsed_seconds >= toFloat64(candidate_duration_windows.duration_seconds)
AND candidate_duration_windows.max_gap_seconds <= gap_stats.max_continuous_gap_seconds
),
Comment thread
Asherlc marked this conversation as resolved.

best_powers AS (
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,202 @@
import { randomUUID } from "node:crypto";
import { sql } from "drizzle-orm";
import { afterAll, beforeAll, describe, expect, it } from "vitest";
import { z } from "zod";
import { readModelSql } from "../../../../analytics/models/read_models/read-model-sql-test-helpers.ts";
import { setupTestDatabase, type TestContext } from "../../../../src/db/test-helpers.ts";
import type { ActivitySensorStore } from "../repositories/activity-repository.ts";
import {
type ClickHouseMetricStreamSeedRow,
createClickHouseTestActivitySensorStore,
seedClickHouseMetricStreamRows,
syncClickHouseTestActivitySensorStore,
} from "./clickhouse-integration-test-helpers.ts";

const testUserId = "00000000-0000-0000-0000-000000000001";
const regularActivityStartedAt = "2026-07-01T12:00:00.000Z";
const gappedActivityStartedAt = "2026-07-01T13:00:00.000Z";
const varyingPowerStartedAt = "2026-07-01T14:00:00.000Z";
const readModelRowSchema = z.object({
activity_id: z.string(),
duration_seconds: z.coerce.number(),
best_power: z.coerce.number().nullable(),
is_deleted: z.coerce.number(),
});

function renderNonIncrementalActivityPowerCurveSql(): string {
return readModelSql("activity_power_curve.sql")
.replace(/^\{\{ config\([\s\S]*?\n\) \}\}\s*/, "")
.replace(/\{\{\s*ref\('activity_summary_rows'\)\s*\}\}/g, "analytics.activity_summary")
.replace(/\{\{\s*ref\('([^']+)'\)\s*\}\}/g, "analytics.$1")
.replace(/FROM analytics\.activity_summary FINAL/g, "FROM analytics.activity_summary")
.replace(
/\n {4}WHERE is_deleted = 0\n {8}AND ended_at IS NOT NULL/,
"\n WHERE ended_at IS NOT NULL",
)
.replace(/\{%\s*if is_incremental\(\)\s*%\}[\s\S]*?\{%\s*endif\s*%\}/g, "");
}

function powerSampleRows(
activityId: string,
startedAt: string,
samples: readonly { offsetSeconds: number; power: number }[],
): ClickHouseMetricStreamSeedRow[] {
const startedAtMs = Date.parse(startedAt);

return samples.map((sample) => ({
activityId,
userId: testUserId,
recordedAt: new Date(startedAtMs + sample.offsetSeconds * 1000).toISOString(),
providerId: "test_provider",
sourceType: "api",
channel: "power",
scalar: sample.power,
}));
}

async function insertActivity(
testContext: TestContext,
activityId: string,
name: string,
startedAt: string,
endedAt: string,
): Promise<void> {
await testContext.db.execute(sql`
INSERT INTO fitness.activity (
id, provider_id, user_id, external_id, activity_type, started_at, ended_at, name
) VALUES (
${activityId}, 'test_provider', ${testUserId}, ${`${name}-${activityId}`}, 'cycling',
${startedAt}, ${endedAt}, ${name}
)
`);
}

describe("activity_power_curve read model", () => {
let testContext: TestContext;
let sensorStore: ActivitySensorStore;

beforeAll(async () => {
testContext = await setupTestDatabase();
await testContext.db.execute(sql`
INSERT INTO fitness.provider (id, name, user_id)
VALUES ('test_provider', 'Test Provider', ${testUserId})
ON CONFLICT DO NOTHING
`);
sensorStore = await createClickHouseTestActivitySensorStore(testContext);
});

afterAll(async () => {
await testContext?.cleanup();
});

it("uses elapsed timestamp duration instead of sample count for power windows", async () => {
const regularActivityId = randomUUID();
const gappedActivityId = randomUUID();
const renderedSql = renderNonIncrementalActivityPowerCurveSql();

await insertActivity(
testContext,
regularActivityId,
"regular-power",
regularActivityStartedAt,
"2026-07-01T12:00:30.000Z",
);
await insertActivity(
testContext,
gappedActivityId,
"gapped-power",
gappedActivityStartedAt,
"2026-07-01T13:00:30.000Z",
);
await syncClickHouseTestActivitySensorStore(testContext);
await seedClickHouseMetricStreamRows(testContext, [
...powerSampleRows(regularActivityId, regularActivityStartedAt, [
{ offsetSeconds: 0, power: 200 },
{ offsetSeconds: 1, power: 200 },
{ offsetSeconds: 2, power: 200 },
{ offsetSeconds: 3, power: 200 },
{ offsetSeconds: 4, power: 200 },
{ offsetSeconds: 5, power: 200 },
]),
...powerSampleRows(gappedActivityId, gappedActivityStartedAt, [
{ offsetSeconds: 0, power: 100 },
{ offsetSeconds: 1, power: 100 },
{ offsetSeconds: 2, power: 100 },
{ offsetSeconds: 20, power: 500 },
{ offsetSeconds: 21, power: 500 },
{ offsetSeconds: 22, power: 500 },
]),
]);

const rows = await sensorStore.query(
readModelRowSchema,
`
SELECT
toString(activity_id) AS activity_id,
duration_seconds,
best_power,
is_deleted
FROM (${renderedSql}) AS power_curve
WHERE duration_seconds = 5
ORDER BY activity_id
`,
);

expect(rows).toEqual([
{
activity_id: regularActivityId,
best_power: 200,
duration_seconds: 5,
is_deleted: 0,
},
]);
});

it("computes average power correctly for varying-power windows", async () => {
const varyingActivityId = randomUUID();
const renderedSql = renderNonIncrementalActivityPowerCurveSql();

await insertActivity(
testContext,
varyingActivityId,
"varying-power",
varyingPowerStartedAt,
"2026-07-01T14:00:06.000Z",
);
await syncClickHouseTestActivitySensorStore(testContext);
await seedClickHouseMetricStreamRows(testContext, [
...powerSampleRows(varyingActivityId, varyingPowerStartedAt, [
{ offsetSeconds: 0, power: 100 },
{ offsetSeconds: 1, power: 200 },
{ offsetSeconds: 2, power: 300 },
{ offsetSeconds: 3, power: 400 },
{ offsetSeconds: 4, power: 500 },
Comment thread
cubic-dev-ai[bot] marked this conversation as resolved.
{ offsetSeconds: 5, power: 300 },
]),
]);

const rows = await sensorStore.query(
readModelRowSchema,
`
SELECT
toString(activity_id) AS activity_id,
duration_seconds,
best_power,
is_deleted
FROM (${renderedSql}) AS power_curve
WHERE activity_id = '${varyingActivityId}'
AND duration_seconds = 5
ORDER BY activity_id
`,
);

expect(rows).toEqual([
{
activity_id: varyingActivityId,
best_power: 300,
duration_seconds: 5,
is_deleted: 0,
},
]);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
});
});
Comment thread
Asherlc marked this conversation as resolved.
Loading