Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
7715e88
feat(backend): queue onboarding refreshes with direct Cloudflare tele…
WcaleNieWolny Sep 17, 2026
add79e1
fix(backend): ignore null dates in latest onboarding bundle lookup
WcaleNieWolny Sep 17, 2026
cc40f68
fix(backend): clamp onboarding telemetry windows at month ends
WcaleNieWolny Sep 17, 2026
14ed146
refactor(backend): use Drizzle for onboarding refresh transactions
WcaleNieWolny Sep 18, 2026
adb043f
refactor(onboarding): remove obsolete batch refresh function
WcaleNieWolny Sep 19, 2026
cb91970
chore(onboarding): restamp refresh migration after main
WcaleNieWolny Sep 19, 2026
4482949
Merge remote-tracking branch 'origin/main' into wolny/backend-onboard…
WcaleNieWolny Sep 19, 2026
e3bdd9c
fix(onboarding): preserve refresh parity for legacy apps
WcaleNieWolny Sep 19, 2026
927fff2
Merge remote-tracking branch 'origin/main' into wolny/backend-onboard…
WcaleNieWolny Sep 19, 2026
79893fb
test(cloudflare): retry synthetic event after worker restart
WcaleNieWolny Sep 19, 2026
0d6da9d
Merge remote-tracking branch 'origin/main' into wolny/backend-onboard…
WcaleNieWolny Sep 19, 2026
291da6e
fix(onboarding): reject missing analytics timestamps
WcaleNieWolny Sep 19, 2026
2c4bc64
refactor(onboarding): run refresh producer in SQL scheduler
WcaleNieWolny Sep 19, 2026
5eb45d0
refactor(onboarding): use Postgres signals for queued refreshes
WcaleNieWolny Sep 20, 2026
238c645
test(onboarding): use valid billing plan and short app ids
WcaleNieWolny Sep 20, 2026
986d27c
test(onboarding): give fixture orgs unique names
WcaleNieWolny Sep 20, 2026
0f1aac6
test(onboarding): account for automatically created org trials
WcaleNieWolny Sep 20, 2026
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
2 changes: 2 additions & 0 deletions cloudflare_workers/api/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ import { app as cron_app_fame } from '../../supabase/functions/_backend/triggers
import { app as cron_clean_orphan_images } from '../../supabase/functions/_backend/triggers/cron_clean_orphan_images.ts'
import { app as cron_clear_versions } from '../../supabase/functions/_backend/triggers/cron_clear_versions.ts'
import { app as cron_email } from '../../supabase/functions/_backend/triggers/cron_email.ts'
import { app as cron_onboarding_refresh_apps } from '../../supabase/functions/_backend/triggers/cron_onboarding_refresh_apps.ts'
import { app as cron_reconcile_build_status } from '../../supabase/functions/_backend/triggers/cron_reconcile_build_status.ts'
import { app as cron_rollout_auto_pause } from '../../supabase/functions/_backend/triggers/cron_rollout_auto_pause.ts'
import { app as cron_stat_app } from '../../supabase/functions/_backend/triggers/cron_stat_app.ts'
Expand Down Expand Up @@ -224,6 +225,7 @@ appTriggers.route('/cron_stat_app', cron_stat_app)
appTriggers.route('/cron_stat_org', cron_stat_org)
appTriggers.route('/cron_sync_sub', cron_sync_sub)
appTriggers.route('/cron_rollout_auto_pause', cron_rollout_auto_pause)
appTriggers.route('/cron_onboarding_refresh_apps', cron_onboarding_refresh_apps)
appTriggers.route('/queue_consumer', queue_consumer)
appTriggers.route('/send_email', send_email)
appTriggers.route('/webhook_delivery', webhook_delivery)
Expand Down
7 changes: 7 additions & 0 deletions read_replicate/schema_replicate.catalog.json
Original file line number Diff line number Diff line change
Expand Up @@ -2265,6 +2265,13 @@
"table": "apps",
"valid": true
},
{
"constraintOwned": false,
"definition": "CREATE INDEX idx_apps_onboarding_queued_refresh_at ON public.apps USING btree (COALESCE((onboarding ->> 'queued_refresh_at'::text), ''::text), COALESCE((onboarding ->> 'refreshed_at'::text), ''::text), app_id)",
"name": "idx_apps_onboarding_queued_refresh_at",
"table": "apps",
"valid": true
},
{
"constraintOwned": false,
"definition": "CREATE INDEX idx_apps_onboarding_refreshed_at ON public.apps USING btree (COALESCE((onboarding ->> 'refreshed_at'::text), ''::text), app_id)",
Expand Down
7 changes: 7 additions & 0 deletions read_replicate/schema_replicate.sql
Original file line number Diff line number Diff line change
Expand Up @@ -891,6 +891,13 @@ CREATE INDEX idx_apps_onboarding_login_creator ON public.apps USING btree (((onb
CREATE INDEX idx_apps_onboarding_ota_stage ON public.apps USING btree (((((onboarding -> 'features'::text) -> 'ota'::text) ->> 'stage'::text)));


--
-- Name: idx_apps_onboarding_queued_refresh_at; Type: INDEX; Schema: public; Owner: -
--

CREATE INDEX idx_apps_onboarding_queued_refresh_at ON public.apps USING btree (COALESCE((onboarding ->> 'queued_refresh_at'::text), ''::text), COALESCE((onboarding ->> 'refreshed_at'::text), ''::text), app_id);


--
-- Name: idx_apps_onboarding_refreshed_at; Type: INDEX; Schema: public; Owner: -
--
Expand Down
4 changes: 0 additions & 4 deletions src/types/supabase.types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5380,10 +5380,6 @@ export type Database = {
Args: { p_user_id: string }
Returns: string
}
refresh_app_onboarding_progress: {
Args: { p_batch_size?: number }
Returns: number
}
refresh_app_rollout_channel_count_for_app: {
Args: { p_app_id: string }
Returns: undefined
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
import type { MiddlewareKeyVariables } from '../utils/hono.ts'
import { Hono } from 'hono/tiny'
import { onboardingRefreshBody, refreshAppOnboardingBatch } from '../utils/app_onboarding_refresh.ts'
import { BRES, middlewareAPISecret, parseBody, quickError } from '../utils/hono.ts'
import { cloudlog } from '../utils/logging.ts'
import { closeClient, getDrizzleClient, getPgClient } from '../utils/pg.ts'

export const app = new Hono<MiddlewareKeyVariables>()
app.post('/', middlewareAPISecret, async (c) => {
const parsed = onboardingRefreshBody.safeParse(await parseBody(c))
if (!parsed.success || new Set(parsed.data.appIds).size !== parsed.data.appIds.length)
throw quickError(400, 'invalid_body', 'Invalid onboarding refresh batch')
const pool = getPgClient(c)
try {
const refreshed = await refreshAppOnboardingBatch(getDrizzleClient(pool), parsed.data)
cloudlog({ requestId: c.get('requestId'), message: 'onboarding refresh batch finished', requested: parsed.data.appIds.length, refreshed })
return c.json(BRES)
}
finally {
await closeClient(c, pool)
}
})
20 changes: 15 additions & 5 deletions supabase/functions/_backend/triggers/queue_consumer.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,15 @@
import type { Context } from 'hono'
import type { MiddlewareKeyVariables } from '../utils/hono.ts'
import type { Database } from '../utils/supabase.types.ts'
import { z } from 'zod'
import { Hono } from 'hono/tiny'
// --- Worker logic imports ---
import { integerLikeSchema, safeParseSchema } from '../utils/schema_validation.ts'
import { z } from 'zod'
import { ONBOARDING_MESSAGES_PER_MINUTE } from '../utils/app_onboarding_refresh.ts'
import { sendDiscordAlert } from '../utils/discord.ts'
import { BRES, middlewareAPISecret, parseBody, simpleError } from '../utils/hono.ts'
import { cloudlog, cloudlogErr, serializeError } from '../utils/logging.ts'
import { closeClient, getPgClient } from '../utils/pg.ts'
// --- Worker logic imports ---
import { integerLikeSchema, safeParseSchema } from '../utils/schema_validation.ts'
import { backgroundTask, getEnv, WAIT_FOR_COMPLETION_HEADER } from '../utils/utils.ts'
import { updateManifestSize } from './on_manifest_create.ts'

Expand Down Expand Up @@ -277,6 +278,8 @@ function prepareQueueHttpBody(functionName: string, body: Record<string, unknown
}

function getQueueHttpTimeoutMs(functionName: string): number {
if (isOnboardingQueue(functionName))
return 45_000
if (isVersionQueueFunction(functionName))
return VERSION_QUEUE_HTTP_TIMEOUT_MS
return QUEUE_HTTP_TIMEOUT_MS
Expand Down Expand Up @@ -530,12 +533,16 @@ async function processQueueMessage(c: Context, queueName: string, message: Messa
}

function getQueueBatchSize(queueName: string, requestedBatchSize: number): number {
if (queueName === 'cron_onboarding_refresh_apps')
return Math.min(requestedBatchSize, ONBOARDING_MESSAGES_PER_MINUTE)
if (isVersionQueueFunction(queueName))
return Math.min(requestedBatchSize, VERSION_QUEUE_BATCH_SIZE)
return requestedBatchSize
}

function getQueueHttpConcurrency(queueName: string): number {
if (isOnboardingQueue(queueName))
return ONBOARDING_MESSAGES_PER_MINUTE
if (queueName === 'on_manifest_create')
return MANIFEST_QUEUE_HTTP_CONCURRENCY
if (isVersionQueueFunction(queueName))
Expand Down Expand Up @@ -1019,7 +1026,6 @@ export async function http_post_helper(
}
}


// Helper function to delete multiple messages from the queue in a single batch
async function delete_queue_message_batch(c: Context, db: ReturnType<typeof getPgClient>, queueName: string, msgIds: number[]) {
try {
Expand Down Expand Up @@ -1096,7 +1102,11 @@ async function mass_edit_queue_messages_cf_ids(

// --- Hono app setup ---
function shouldRunQueueSyncInBackground(queueName: string): boolean {
return queueName !== 'on_manifest_create'
return queueName !== 'on_manifest_create' && !isOnboardingQueue(queueName)
}

function isOnboardingQueue(queueName: string): boolean {
return queueName === 'cron_onboarding_refresh_apps'
}

async function runQueueSync(
Expand Down
144 changes: 144 additions & 0 deletions supabase/functions/_backend/utils/app_onboarding_refresh.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
import type { getDrizzleClient } from './pg.ts'
import { sql } from 'drizzle-orm'
import { z } from 'zod'

export const ONBOARDING_APPS_PER_MESSAGE = 25
export const ONBOARDING_MESSAGES_PER_MINUTE = 4
export const onboardingRefreshBody = z.object({
appIds: z.array(z.string().min(1)).min(1).max(ONBOARDING_APPS_PER_MESSAGE),
queuedAt: z.iso.datetime(),
})

interface OnboardingSignals extends Record<string, unknown> {
app_id: string
last_device_at: Date | null
has_app_store: boolean | null
has_testflight: boolean | null
has_play_unknown: boolean | null
has_native: boolean | null
has_install_source: boolean | null
first_bundle_at: Date | null
last_bundle_at: Date | null
first_install_at: Date | null
last_install_at: Date | null
first_build_at: Date | null
first_success_at: Date | null
last_build_at: Date | null
}

export async function refreshAppOnboardingBatch(
database: Pick<ReturnType<typeof getDrizzleClient>, 'transaction'>,
body: z.infer<typeof onboardingRefreshBody>,
) {
return database.transaction(async (tx) => {
await tx.execute(sql`SELECT
pg_catalog.set_config('statement_timeout', '35s', true),
pg_catalog.set_config('lock_timeout', '5s', true),
pg_catalog.set_config('TimeZone', 'UTC', true)
`)

// Avoid reading large signal tables for messages already covered by a
// later refresh. Recheck this after locking because another worker may win.
const { rows: dueApps } = await tx.execute<{ app_id: string }>(sql`
SELECT app_id FROM public.apps
WHERE app_id = ANY(${sql.param(body.appIds)}::varchar[])
AND COALESCE(onboarding->>'refreshed_at', '') < ${body.queuedAt}
ORDER BY app_id
`)
if (!dueApps.length)
return 0

const dueIds = dueApps.map(app => app.app_id)
// Preserve the old cron's PostgreSQL evidence and date precision. Each
// lateral aggregate starts from an indexed app_id and returns one row.
const { rows: signals } = await tx.execute<OnboardingSignals>(sql`
SELECT batch.app_id,
d.last_device_at, d.has_app_store, d.has_testflight,
d.has_play_unknown, d.has_native, d.has_install_source,
v.first_bundle_at, v.last_bundle_at,
dv.first_install_at, dv.last_install_at,
br.first_build_at, br.first_success_at, br.last_build_at
FROM pg_catalog.unnest(${sql.param(dueIds)}::varchar[]) AS batch(app_id)
LEFT JOIN LATERAL (
SELECT
bool_or(install_source = 'app_store') AS has_app_store,
bool_or(install_source = 'testflight') AS has_testflight,
bool_or(install_source IN ('google_play', 'amazon_appstore', 'samsung_galaxy_store', 'huawei_appgallery')) AS has_play_unknown,
bool_or(is_prod IS TRUE AND is_emulator IS NOT TRUE) AS has_native,
bool_or(install_source IS NOT NULL) AS has_install_source,
max(updated_at) AS last_device_at
FROM public.devices
WHERE app_id = batch.app_id
AND (install_source IS NOT NULL OR (is_prod IS TRUE AND is_emulator IS NOT TRUE))
) d ON true
LEFT JOIN LATERAL (
SELECT min(created_at) AS first_bundle_at, max(created_at) AS last_bundle_at
FROM public.app_versions
WHERE app_id = batch.app_id AND deleted IS NOT TRUE
AND name IS DISTINCT FROM 'builtin' AND name IS DISTINCT FROM 'unknown'
) v ON true
LEFT JOIN LATERAL (
SELECT min(date)::timestamptz AS first_install_at,
max(date)::timestamptz AS last_install_at
FROM public.daily_version
WHERE app_id = batch.app_id AND COALESCE(install, 0) > 0
) dv ON true
LEFT JOIN LATERAL (
SELECT min(created_at) AS first_build_at,
min(completed_at) FILTER (WHERE status IN ('succeeded', 'released')) AS first_success_at,
max(COALESCE(completed_at, created_at)) AS last_build_at
FROM public.build_requests WHERE app_id = batch.app_id
) br ON true
ORDER BY batch.app_id
`)

await tx.execute(sql`
SELECT app_id FROM public.apps
WHERE app_id = ANY(${sql.param(dueIds)}::varchar[])
ORDER BY app_id FOR UPDATE
`)
const result = await tx.execute(sql`
WITH signals AS (
SELECT * FROM pg_catalog.jsonb_to_recordset(${JSON.stringify(signals)}::jsonb) AS s(
app_id varchar, last_device_at timestamptz,
has_app_store boolean, has_testflight boolean, has_play_unknown boolean,
has_native boolean, has_install_source boolean,
first_bundle_at timestamptz, last_bundle_at timestamptz,
first_install_at timestamptz, last_install_at timestamptz,
first_build_at timestamptz, first_success_at timestamptz, last_build_at timestamptz
)
)
UPDATE public.apps a SET onboarding = pg_catalog.jsonb_set(
pg_catalog.jsonb_set(
a.onboarding, '{features}',
COALESCE(a.onboarding->'features', '{}'::jsonb) || pg_catalog.jsonb_build_object(
'cli_install', public.merge_app_onboarding_feature(
a.onboarding->'features'->'cli_install', s.last_device_at,
s.last_device_at, s.last_device_at, NULL),
'ota', public.merge_app_onboarding_feature(
a.onboarding->'features'->'ota', s.first_bundle_at, s.first_install_at,
GREATEST(s.last_install_at, s.last_bundle_at),
CASE
WHEN s.has_app_store THEN 'store_live'
WHEN s.has_testflight THEN 'testflight'
WHEN s.has_play_unknown THEN 'play_unknown'
WHEN s.has_native THEN 'native_unknown'
WHEN s.has_install_source THEN 'local_only'
ELSE 'no_device'
END),
'builder', public.merge_app_onboarding_feature(
a.onboarding->'features'->'builder', s.first_build_at,
s.first_success_at, s.last_build_at, NULL)
), true
), '{refreshed_at}',
pg_catalog.to_jsonb(pg_catalog.to_char((now() AT TIME ZONE 'UTC'), 'YYYY-MM-DD"T"HH24:MI:SS.MS"Z"')),
true
)
FROM signals s
WHERE a.app_id = s.app_id
AND COALESCE(a.onboarding->>'refreshed_at', '') < ${body.queuedAt}
RETURNING a.app_id
`)
return result.rowCount ?? 0
})
}
4 changes: 0 additions & 4 deletions supabase/functions/_backend/utils/supabase.types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5380,10 +5380,6 @@ export type Database = {
Args: { p_user_id: string }
Returns: string
}
refresh_app_onboarding_progress: {
Args: { p_batch_size?: number }
Returns: number
}
refresh_app_rollout_channel_count_for_app: {
Args: { p_app_id: string }
Returns: undefined
Expand Down
2 changes: 2 additions & 0 deletions supabase/functions/triggers/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import { app as cron_app_fame } from '../_backend/triggers/cron_app_fame.ts'
import { app as cron_clean_orphan_images } from '../_backend/triggers/cron_clean_orphan_images.ts'
import { app as cron_clear_versions } from '../_backend/triggers/cron_clear_versions.ts'
import { app as cron_email } from '../_backend/triggers/cron_email.ts'
import { app as cron_onboarding_refresh_apps } from '../_backend/triggers/cron_onboarding_refresh_apps.ts'
import { app as cron_reconcile_build_status } from '../_backend/triggers/cron_reconcile_build_status.ts'
import { app as cron_rollout_auto_pause } from '../_backend/triggers/cron_rollout_auto_pause.ts'
import { app as cron_stat_app } from '../_backend/triggers/cron_stat_app.ts'
Expand Down Expand Up @@ -81,6 +82,7 @@ appGlobal.route('/cron_clear_versions', cron_clear_versions)
appGlobal.route('/cron_clean_orphan_images', cron_clean_orphan_images)
appGlobal.route('/cron_reconcile_build_status', cron_reconcile_build_status)
appGlobal.route('/cron_rollout_auto_pause', cron_rollout_auto_pause)
appGlobal.route('/cron_onboarding_refresh_apps', cron_onboarding_refresh_apps)
appGlobal.route('/canceled_org_retention_alerts', canceled_org_retention_alerts)
appGlobal.route('/credit_usage_alerts', credit_usage_alerts)
appGlobal.route('/credit_usage_posthog', credit_usage_posthog)
Expand Down
Loading
Loading