Correct product analytics cohorts and reporting
CI / verify (push) Has been cancelled
CI / publish (push) Has been cancelled

This commit is contained in:
Rocky
2026-08-21 17:50:35 +08:00
parent b25f5ae6e9
commit edd0d9feca
31 changed files with 1134 additions and 262 deletions
@@ -604,7 +604,7 @@ private suspend fun AdminStatsService.getRange(
)
}
private fun parseAdminStatsRange(range: String?, clock: Clock): Pair<java.time.Instant, java.time.Instant>? {
internal fun parseAdminStatsRange(range: String?, clock: Clock): Pair<java.time.Instant, java.time.Instant>? {
val days = when (range) {
null, "30d" -> 30L
"7d" -> 7L
@@ -612,7 +612,9 @@ private fun parseAdminStatsRange(range: String?, clock: Clock): Pair<java.time.I
else -> return null
}
val until = clock.instant()
return until.minus(Duration.ofDays(days)) to until
val firstIncludedDate = until.atZone(ZoneOffset.UTC).toLocalDate().minusDays(days - 1)
val from = firstIncludedDate.atStartOfDay(ZoneOffset.UTC).toInstant()
return from to until
}
private data class AdminReferralQueryOptions(
@@ -953,6 +955,7 @@ private fun adminCookie(
private fun AdminStatsDto.toOverviewResponse(): AdminOverviewResponse {
val consumedByDate = creditFlow.associateBy { it.date }
return AdminOverviewResponse(
period = period,
totalUsers = overview.totalUsers,
activeUsers = overview.activeUsers,
newUsers = overview.registrations,
@@ -972,12 +975,13 @@ private fun AdminStatsDto.toOverviewResponse(): AdminOverviewResponse {
private fun AdminStatsDto.toReferralResponse(): AdminReferralResponse =
AdminReferralResponse(
period = period,
pendingBindings = referralFunnel.pendingBindings,
ineligibleBindings = referralFunnel.ineligibleBindings,
funnel = listOf(
AdminFunnelResponse("邀请码创建", referralFunnel.codesCreated),
AdminFunnelResponse("成功绑定", referralFunnel.bindings),
AdminFunnelResponse("有效使用并奖励", referralFunnel.rewardedBindings),
AdminFunnelResponse("绑定后首次 AI 成功", referralFunnel.activatedBindings),
AdminFunnelResponse("完成奖励", referralFunnel.rewardedBindings),
),
ranking = referralRanking.map {
AdminReferralRankResponse(
@@ -1137,6 +1141,7 @@ private data class PageResponse<T>(val items: List<T>, val nextCursor: String? =
@Serializable
private data class AdminOverviewResponse(
val period: com.osglab.account.features.admin.stats.models.AdminStatsPeriodDto,
val totalUsers: Long,
val activeUsers: Long,
val newUsers: Long,
@@ -1156,6 +1161,7 @@ private data class AdminTrendResponse(
@Serializable
private data class AdminReferralResponse(
val period: com.osglab.account.features.admin.stats.models.AdminStatsPeriodDto,
val pendingBindings: Long,
val ineligibleBindings: Long,
val funnel: List<AdminFunnelResponse>,
@@ -34,6 +34,19 @@ data class AdminAnalyticsFeatureUsageDto(
val successes: Long,
)
@Serializable
data class AdminAnalyticsReferralSignalsDto(
val shared: Long,
val opened: Long,
)
@Serializable
data class AdminAnalyticsLatencyBucketDto(
val bucket: String,
val successful: Long,
val failed: Long,
)
@Serializable
data class AdminAnalyticsFunnelStepDto(
val label: String,
@@ -88,6 +101,8 @@ data class AdminAnalyticsMonetizationDto(
val conversion7d: AdminAnalyticsRateDto,
val conversion30d: AdminAnalyticsRateDto,
val repeatPurchaseRate: AdminAnalyticsRateDto,
val purchaseFunnel: List<AdminAnalyticsFunnelStepDto>,
val cancelledUsers: Long,
)
@Serializable
@@ -95,6 +110,7 @@ data class AdminAnalyticsGuardrailsDto(
val clientAiSuccessRate: AdminAnalyticsRateDto,
val managedSuccessRate: AdminAnalyticsRateDto,
val creditBlockedUsers: Long,
val latencyBuckets: List<AdminAnalyticsLatencyBucketDto>,
)
@Serializable
@@ -130,6 +146,7 @@ data class AdminProductAnalyticsDto(
val retention: List<AdminAnalyticsCohortDto>,
val aiFeatures: List<AdminAnalyticsFeatureUsageDto>,
val keyboardUsage: AdminAnalyticsKeyboardUsageDto,
val referralSignals: AdminAnalyticsReferralSignalsDto,
val referralFunnel: List<AdminAnalyticsFunnelStepDto>,
val guardrails: AdminAnalyticsGuardrailsDto,
)
@@ -35,6 +35,7 @@ data class AdminCreditFlowPointDto(
data class AdminReferralFunnelDto(
val codesCreated: Long,
val bindings: Long,
val activatedBindings: Long,
val rewardedBindings: Long,
val pendingBindings: Long,
val ineligibleBindings: Long,
@@ -0,0 +1,29 @@
package com.osglab.account.features.admin.stats.repositories
/**
* Canonical AI value events used by product analytics. Managed usage is sourced
* from immutable billing records; LOCAL and BYOK usage comes from terminal
* client events. No user content is selected.
*/
internal fun identityValueEventsCte(): String =
"""
WITH value_events AS (
SELECT
CONCAT('a:', user_id) AS identity_key,
created_at AS occurred_at
FROM credit_usage_records
UNION ALL
SELECT
COALESCE(
CONCAT('a:', i.account_id),
CONCAT('i:', e.installation_hash)
) AS identity_key,
e.occurred_at
FROM product_analytics_events e
JOIN product_analytics_installations i
ON i.installation_hash = e.installation_hash
WHERE e.event_name = 'AI_FEATURE_SUCCEEDED'
AND e.execution_mode IN ('LOCAL', 'BYOK')
)
""".trimIndent()
@@ -74,6 +74,19 @@ data class AdminAnalyticsGuardrailRow(
val creditBlockedUsers: Long,
)
data class AdminAnalyticsLatencyRow(
val bucket: String,
val successful: Long,
val failed: Long,
)
data class AdminAnalyticsPurchaseFunnelRow(
val viewed: Long,
val started: Long,
val verified: Long,
val cancelled: Long,
)
data class AdminAnalyticsKeyboardUsageRow(
val activeUsers: Long,
val keyboardUsers: Long,
@@ -95,7 +108,6 @@ data class AdminAnalyticsGrowthFunnelRow(
val opened: Long,
val registered: Long,
val activated: Long,
val retainedD7: Long,
val purchased: Long,
)
@@ -120,6 +132,8 @@ data class AdminProductAnalyticsSnapshot(
val keyboardUsage: AdminAnalyticsKeyboardUsageRow,
val referrals: AdminAnalyticsReferralRow,
val guardrails: AdminAnalyticsGuardrailRow,
val latencyDistribution: List<AdminAnalyticsLatencyRow>,
val purchaseFunnel: AdminAnalyticsPurchaseFunnelRow,
)
interface AdminProductAnalyticsRepository {
@@ -142,7 +156,7 @@ class ExposedAdminProductAnalyticsRepository(
AdminProductAnalyticsSnapshot(
currentWeeklyUsers = loadValueActiveUsers(currentWeek),
previousWeeklyUsers = loadValueActiveUsers(previousWeek),
newInstallations = activation.denominator,
newInstallations = loadNewInstallations(range),
newAccounts = loadNewAccounts(range),
activation24h = AdminAnalyticsCountRow(activation.activated, activation.denominator),
medianTimeToValueMinutes = activation.medianMinutes,
@@ -166,12 +180,14 @@ class ExposedAdminProductAnalyticsRepository(
keyboardUsage = loadKeyboardUsage(range),
referrals = loadReferrals(range),
guardrails = loadGuardrails(range),
latencyDistribution = loadLatencyDistribution(range),
purchaseFunnel = loadPurchaseFunnel(range),
)
}
private fun loadValueActiveUsers(window: AdminAnalyticsWindow): Long =
querySingle(
valueEventsCte() +
identityValueEventsCte() +
"""
SELECT COUNT(DISTINCT identity_key) AS aggregate_value
FROM value_events
@@ -190,6 +206,21 @@ class ExposedAdminProductAnalyticsRepository(
range.arguments(),
) { it.exactLong("aggregate_value") }
private fun loadNewInstallations(range: AdminAnalyticsWindow): Long =
querySingle(
"""
SELECT COUNT(*) AS aggregate_value
FROM (
SELECT installation_hash, MIN(occurred_at) AS opened_at
FROM product_analytics_events
WHERE event_name = 'FIRST_OPEN'
GROUP BY installation_hash
) first_open
WHERE opened_at >= ? AND opened_at < ?
""",
range.arguments(),
) { it.exactLong("aggregate_value") }
private fun loadActivation(range: AdminAnalyticsWindow): ActivationRow =
querySingle(
"""
@@ -197,21 +228,30 @@ class ExposedAdminProductAnalyticsRepository(
SELECT installation_hash, MIN(occurred_at) AS opened_at
FROM product_analytics_events
WHERE event_name = 'FIRST_OPEN'
AND occurred_at >= ? AND occurred_at < ?
GROUP BY installation_hash
HAVING opened_at >= ?
AND opened_at <= DATE_SUB(?, INTERVAL 24 HOUR)
),
client_value AS (
SELECT installation_hash, MIN(occurred_at) AS value_at
FROM product_analytics_events
WHERE event_name = 'AI_FEATURE_SUCCEEDED'
AND execution_mode IN ('LOCAL', 'BYOK')
GROUP BY installation_hash
SELECT e.installation_hash, MIN(e.occurred_at) AS value_at
FROM first_open o
JOIN product_analytics_events e
ON e.installation_hash = o.installation_hash
WHERE e.event_name = 'AI_FEATURE_SUCCEEDED'
AND e.execution_mode IN ('LOCAL', 'BYOK')
AND e.occurred_at >= o.opened_at
AND e.occurred_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
GROUP BY e.installation_hash
),
managed_value AS (
SELECT i.installation_hash, MIN(u.created_at) AS value_at
FROM product_analytics_installations i
SELECT o.installation_hash, MIN(u.created_at) AS value_at
FROM first_open o
JOIN product_analytics_installations i
ON i.installation_hash = o.installation_hash
JOIN credit_usage_records u ON u.user_id = i.account_id
GROUP BY i.installation_hash
WHERE u.created_at >= o.opened_at
AND u.created_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
GROUP BY o.installation_hash
),
first_value_by_install AS (
SELECT installation_hash, MIN(value_at) AS value_at
@@ -228,8 +268,6 @@ class ExposedAdminProductAnalyticsRepository(
TIMESTAMPDIFF(SECOND, o.opened_at, v.value_at) AS seconds_to_value
FROM first_open o
JOIN first_value_by_install v ON v.installation_hash = o.installation_hash
WHERE v.value_at >= o.opened_at
AND v.value_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
),
ranked AS (
SELECT
@@ -277,21 +315,30 @@ class ExposedAdminProductAnalyticsRepository(
) AS channel
FROM product_analytics_events e
WHERE e.event_name = 'FIRST_OPEN'
AND e.occurred_at >= ? AND e.occurred_at < ?
GROUP BY e.installation_hash
HAVING opened_at >= ?
AND opened_at <= DATE_SUB(?, INTERVAL 24 HOUR)
),
client_value AS (
SELECT installation_hash, MIN(occurred_at) AS value_at
FROM product_analytics_events
WHERE event_name = 'AI_FEATURE_SUCCEEDED'
AND execution_mode IN ('LOCAL', 'BYOK')
GROUP BY installation_hash
SELECT e.installation_hash, MIN(e.occurred_at) AS value_at
FROM first_open o
JOIN product_analytics_events e
ON e.installation_hash = o.installation_hash
WHERE e.event_name = 'AI_FEATURE_SUCCEEDED'
AND e.execution_mode IN ('LOCAL', 'BYOK')
AND e.occurred_at >= o.opened_at
AND e.occurred_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
GROUP BY e.installation_hash
),
managed_value AS (
SELECT i.installation_hash, MIN(u.created_at) AS value_at
FROM product_analytics_installations i
SELECT o.installation_hash, MIN(u.created_at) AS value_at
FROM first_open o
JOIN product_analytics_installations i
ON i.installation_hash = o.installation_hash
JOIN credit_usage_records u ON u.user_id = i.account_id
GROUP BY i.installation_hash
WHERE u.created_at >= o.opened_at
AND u.created_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
GROUP BY o.installation_hash
),
first_value_by_install AS (
SELECT installation_hash, MIN(value_at) AS value_at
@@ -307,9 +354,7 @@ class ExposedAdminProductAnalyticsRepository(
COUNT(*) AS installations,
SUM(
CASE
WHEN v.value_at >= o.opened_at
AND v.value_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
THEN 1 ELSE 0
WHEN v.value_at IS NOT NULL THEN 1 ELSE 0
END
) AS activated
FROM first_open o
@@ -397,10 +442,8 @@ class ExposedAdminProductAnalyticsRepository(
)
}
private fun loadMonetization(range: AdminAnalyticsWindow): AdminAnalyticsMonetizationRow {
val sevenDayMaturity = range.until.minusSeconds(7 * DAY_SECONDS)
val thirtyDayMaturity = range.until.minusSeconds(30 * DAY_SECONDS)
return querySingle(
private fun loadMonetization(range: AdminAnalyticsWindow): AdminAnalyticsMonetizationRow =
querySingle(
"""
WITH first_purchase AS (
SELECT user_id, MIN(purchased_at) AS first_purchased_at, COUNT(*) AS lifetime_purchases
@@ -459,10 +502,10 @@ class ExposedAdminProductAnalyticsRepository(
""",
buildList {
addAll(range.arguments(repetitions = 3))
addAll(maturedWindowArguments(range.from, sevenDayMaturity))
addAll(maturedWindowArguments(range.from, sevenDayMaturity))
addAll(maturedWindowArguments(range.from, thirtyDayMaturity))
addAll(maturedWindowArguments(range.from, thirtyDayMaturity))
addAll(maturedCohortArguments(range, 7))
addAll(maturedCohortArguments(range, 7))
addAll(maturedCohortArguments(range, 30))
addAll(maturedCohortArguments(range, 30))
},
) {
AdminAnalyticsMonetizationRow(
@@ -483,7 +526,6 @@ class ExposedAdminProductAnalyticsRepository(
),
)
}
}
private fun loadGrowthFunnel(range: AdminAnalyticsWindow): AdminAnalyticsGrowthFunnelRow =
querySingle(
@@ -492,8 +534,22 @@ class ExposedAdminProductAnalyticsRepository(
SELECT installation_hash, MIN(occurred_at) AS opened_at
FROM product_analytics_events
WHERE event_name = 'FIRST_OPEN'
AND occurred_at >= ? AND occurred_at < ?
GROUP BY installation_hash
HAVING opened_at >= ?
AND opened_at <= DATE_SUB(?, INTERVAL 24 HOUR)
),
registered AS (
SELECT
o.installation_hash,
o.opened_at,
i.account_id,
a.created_at AS registered_at
FROM first_open o
JOIN product_analytics_installations i
ON i.installation_hash = o.installation_hash
JOIN accounts a ON a.id = i.account_id
WHERE a.created_at >= o.opened_at
AND a.created_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
),
client_values AS (
SELECT installation_hash, occurred_at
@@ -513,63 +569,42 @@ class ExposedAdminProductAnalyticsRepository(
),
activated AS (
SELECT
o.installation_hash,
o.opened_at,
r.installation_hash,
r.opened_at,
r.account_id,
MIN(v.occurred_at) AS first_value_at
FROM first_open o
JOIN values_by_install v ON v.installation_hash = o.installation_hash
WHERE v.occurred_at >= o.opened_at
AND v.occurred_at <= DATE_ADD(o.opened_at, INTERVAL 24 HOUR)
GROUP BY o.installation_hash, o.opened_at
FROM registered r
JOIN values_by_install v ON v.installation_hash = r.installation_hash
WHERE v.occurred_at >= r.registered_at
AND v.occurred_at <= DATE_ADD(r.opened_at, INTERVAL 24 HOUR)
GROUP BY r.installation_hash, r.opened_at, r.account_id
),
purchased AS (
SELECT DISTINCT a.installation_hash
FROM activated a
JOIN storekit_credit_purchases p ON p.user_id = a.account_id
WHERE p.purchased_at >= a.first_value_at
AND p.purchased_at <= DATE_ADD(a.opened_at, INTERVAL 24 HOUR)
)
SELECT
(SELECT COUNT(*) FROM first_open) AS opened,
(
SELECT COUNT(*)
FROM first_open o
JOIN product_analytics_installations i
ON i.installation_hash = o.installation_hash
WHERE i.account_id IS NOT NULL
) AS registered,
(SELECT COUNT(*) FROM registered) AS registered,
(SELECT COUNT(*) FROM activated) AS activated,
(
SELECT COUNT(*)
FROM activated a
WHERE EXISTS (
SELECT 1
FROM values_by_install v
WHERE v.installation_hash = a.installation_hash
AND DATE(v.occurred_at) = DATE_ADD(DATE(a.first_value_at), INTERVAL 7 DAY)
)
) AS retained_d7,
(
SELECT COUNT(*)
FROM first_open o
JOIN product_analytics_installations i
ON i.installation_hash = o.installation_hash
WHERE EXISTS (
SELECT 1
FROM storekit_credit_purchases p
WHERE p.user_id = i.account_id
AND p.purchased_at >= o.opened_at
AND p.purchased_at < ?
)
) AS purchased
(SELECT COUNT(*) FROM purchased) AS purchased
""",
range.arguments() + listOf(INSTANT_COLUMN_TYPE to range.until),
range.arguments(),
) {
AdminAnalyticsGrowthFunnelRow(
opened = it.exactLong("opened"),
registered = it.exactLong("registered"),
activated = it.exactLong("activated"),
retainedD7 = it.exactLong("retained_d7"),
purchased = it.exactLong("purchased"),
)
}
private fun loadRetention(range: AdminAnalyticsWindow): List<AdminAnalyticsCohortRow> =
queryRows(
valueEventsCte() +
identityValueEventsCte() +
"""
, first_value_by_identity AS (
SELECT identity_key, MIN(occurred_at) AS first_value_at
@@ -737,7 +772,7 @@ class ExposedAdminProductAnalyticsRepository(
private fun loadReferrals(range: AdminAnalyticsWindow): AdminAnalyticsReferralRow =
querySingle(
valueEventsCte() +
identityValueEventsCte() +
"""
SELECT
(
@@ -755,7 +790,7 @@ class ExposedAdminProductAnalyticsRepository(
SELECT COALESCE(SUM(counter_value), 0)
FROM product_analytics_daily_counters
WHERE counter_name = 'INVITE_PAGE_OPENED'
AND counter_date >= DATE(?) AND counter_date < DATE(?)
AND counter_date >= DATE(?) AND counter_date <= DATE(?)
) AS opened,
(
SELECT COUNT(*)
@@ -779,12 +814,14 @@ class ExposedAdminProductAnalyticsRepository(
FROM referral_bindings
WHERE bound_at >= ? AND bound_at < ?
AND reward_status = 'REWARDED'
AND rewarded_at < ?
) AS rewarded
""",
buildList {
addAll(range.arguments(repetitions = 5))
add(INSTANT_COLUMN_TYPE to range.until)
addAll(range.arguments())
add(INSTANT_COLUMN_TYPE to range.until)
},
) {
AdminAnalyticsReferralRow(
@@ -796,6 +833,91 @@ class ExposedAdminProductAnalyticsRepository(
)
}
private fun loadLatencyDistribution(range: AdminAnalyticsWindow): List<AdminAnalyticsLatencyRow> =
queryRows(
"""
SELECT
duration_bucket,
SUM(CASE WHEN event_name = 'AI_FEATURE_SUCCEEDED' THEN 1 ELSE 0 END) AS successful,
SUM(CASE WHEN event_name = 'AI_FEATURE_FAILED' THEN 1 ELSE 0 END) AS failed
FROM product_analytics_events
WHERE occurred_at >= ? AND occurred_at < ?
AND event_name IN ('AI_FEATURE_SUCCEEDED', 'AI_FEATURE_FAILED')
AND duration_bucket IS NOT NULL
GROUP BY duration_bucket
ORDER BY FIELD(
duration_bucket,
'LT_1S',
'S1_TO_3',
'S3_TO_10',
'S10_TO_30',
'GTE_30S'
)
""",
range.arguments(),
) {
AdminAnalyticsLatencyRow(
bucket = it.getString("duration_bucket"),
successful = it.exactLong("successful"),
failed = it.exactLong("failed"),
)
}
private fun loadPurchaseFunnel(range: AdminAnalyticsWindow): AdminAnalyticsPurchaseFunnelRow =
querySingle(
"""
WITH viewed AS (
SELECT installation_hash, MIN(occurred_at) AS viewed_at
FROM product_analytics_events
WHERE event_name = 'PURCHASE_VIEWED'
AND occurred_at >= ? AND occurred_at < ?
GROUP BY installation_hash
),
started AS (
SELECT v.installation_hash, MIN(e.occurred_at) AS started_at
FROM viewed v
JOIN product_analytics_events e
ON e.installation_hash = v.installation_hash
AND e.event_name = 'PURCHASE_STARTED'
AND e.occurred_at >= v.viewed_at
AND e.occurred_at < ?
GROUP BY v.installation_hash
),
verified AS (
SELECT DISTINCT s.installation_hash
FROM started s
JOIN product_analytics_installations i
ON i.installation_hash = s.installation_hash
JOIN storekit_credit_purchases p ON p.user_id = i.account_id
WHERE p.purchased_at >= s.started_at
AND p.purchased_at < ?
)
SELECT
(SELECT COUNT(*) FROM viewed) AS viewed,
(SELECT COUNT(*) FROM started) AS started,
(SELECT COUNT(*) FROM verified) AS verified,
(
SELECT COUNT(DISTINCT installation_hash)
FROM product_analytics_events
WHERE event_name = 'PURCHASE_CANCELLED'
AND occurred_at >= ? AND occurred_at < ?
) AS cancelled
""",
buildList {
addAll(range.arguments())
add(INSTANT_COLUMN_TYPE to range.until)
add(INSTANT_COLUMN_TYPE to range.until)
addAll(range.arguments())
},
) {
AdminAnalyticsPurchaseFunnelRow(
viewed = it.exactLong("viewed"),
started = it.exactLong("started"),
verified = it.exactLong("verified"),
cancelled = it.exactLong("cancelled"),
)
}
private fun loadGuardrails(range: AdminAnalyticsWindow): AdminAnalyticsGuardrailRow =
querySingle(
"""
@@ -861,28 +983,6 @@ private data class ActivationRow(
val medianMinutes: Double?,
)
private fun valueEventsCte(): String =
"""
WITH value_events AS (
SELECT
CONCAT('a:', user_id) AS identity_key,
created_at AS occurred_at
FROM credit_usage_records
UNION ALL
SELECT
COALESCE(
CONCAT('a:', i.account_id),
CONCAT('i:', e.installation_hash)
) AS identity_key,
e.occurred_at
FROM product_analytics_events e
JOIN product_analytics_installations i
ON i.installation_hash = e.installation_hash
WHERE e.event_name = 'AI_FEATURE_SUCCEEDED'
AND e.execution_mode IN ('LOCAL', 'BYOK')
)
""".trimIndent()
private fun AdminAnalyticsWindow.arguments(
repetitions: Int = 1,
): List<Pair<IColumnType<*>, Any?>> = buildList {
@@ -892,13 +992,13 @@ private fun AdminAnalyticsWindow.arguments(
}
}
private fun maturedWindowArguments(
from: Instant,
maturityEnd: Instant,
private fun maturedCohortArguments(
range: AdminAnalyticsWindow,
observationDays: Long,
): List<Pair<IColumnType<*>, Any?>> =
listOf(
INSTANT_COLUMN_TYPE to from,
INSTANT_COLUMN_TYPE to maxOf(from, maturityEnd),
INSTANT_COLUMN_TYPE to range.from.minusSeconds(observationDays * DAY_SECONDS),
INSTANT_COLUMN_TYPE to range.until.minusSeconds(observationDays * DAY_SECONDS),
)
private fun <T> querySingle(
@@ -114,8 +114,28 @@ class ExposedAdminStatsRepository(
) AS registrations,
(
SELECT COUNT(DISTINCT user_id)
FROM credit_usage_records
WHERE created_at >= ? AND created_at < ?
FROM (
SELECT user_id
FROM credit_usage_records u
JOIN accounts a ON a.id = u.user_id
WHERE u.created_at >= ? AND u.created_at < ?
UNION
SELECT i.account_id AS user_id
FROM product_analytics_events e
JOIN product_analytics_installations i
ON i.installation_hash = e.installation_hash
WHERE e.occurred_at >= ? AND e.occurred_at < ?
AND e.event_name = 'AI_FEATURE_SUCCEEDED'
AND e.execution_mode IN ('LOCAL', 'BYOK')
AND i.account_id IS NOT NULL
UNION
SELECT i.account_id AS user_id
FROM keyboard_usage_daily_summaries s
JOIN product_analytics_installations i
ON i.installation_hash = s.installation_hash
WHERE s.summary_date >= DATE(?) AND s.summary_date < DATE(?)
AND i.account_id IS NOT NULL
) registered_activity
) AS active_users,
(
SELECT COALESCE(SUM(balance), 0)
@@ -134,7 +154,7 @@ class ExposedAdminStatsRepository(
WHERE created_at >= ? AND created_at < ?
) AS consumed_credits
""",
range.arguments(repetitions = 4),
range.arguments(repetitions = 6),
) { result ->
AdminOverviewDto(
totalUsers = result.exactLong("total_users"),
@@ -148,7 +168,13 @@ class ExposedAdminStatsRepository(
private fun loadReferralFunnel(range: AdminStatsRange): AdminReferralFunnelDto =
querySingle(
"""
identityValueEventsCte() +
"""
, binding_cohort AS (
SELECT *
FROM referral_bindings
WHERE bound_at >= ? AND bound_at < ?
)
SELECT
(
SELECT COUNT(*)
@@ -157,33 +183,47 @@ class ExposedAdminStatsRepository(
) AS codes_created,
(
SELECT COUNT(*)
FROM referral_bindings
WHERE bound_at >= ? AND bound_at < ?
FROM binding_cohort
) AS bindings,
(
SELECT COUNT(*)
FROM referral_bindings
FROM binding_cohort r
WHERE EXISTS (
SELECT 1
FROM value_events v
WHERE v.identity_key = CONCAT('a:', r.invitee_user_id)
AND v.occurred_at >= r.bound_at
AND v.occurred_at < ?
)
) AS activated_bindings,
(
SELECT COUNT(*)
FROM binding_cohort
WHERE reward_status = 'REWARDED'
AND rewarded_at >= ? AND rewarded_at < ?
AND rewarded_at < ?
) AS rewarded_bindings,
(
SELECT COUNT(*)
FROM referral_bindings
FROM binding_cohort
WHERE reward_status = 'PENDING'
AND bound_at >= ? AND bound_at < ?
) AS pending_bindings,
(
SELECT COUNT(*)
FROM referral_bindings
FROM binding_cohort
WHERE reward_status = 'INELIGIBLE_BUDGET'
AND bound_at >= ? AND bound_at < ?
) AS ineligible_bindings
""",
range.arguments(repetitions = 5),
buildList {
addAll(range.arguments())
addAll(range.arguments())
add(INSTANT_COLUMN_TYPE to range.until)
add(INSTANT_COLUMN_TYPE to range.until)
},
) { result ->
AdminReferralFunnelDto(
codesCreated = result.exactLong("codes_created"),
bindings = result.exactLong("bindings"),
activatedBindings = result.exactLong("activated_bindings"),
rewardedBindings = result.exactLong("rewarded_bindings"),
pendingBindings = result.exactLong("pending_bindings"),
ineligibleBindings = result.exactLong("ineligible_bindings"),
@@ -197,23 +237,19 @@ class ExposedAdminStatsRepository(
"""
SELECT
inviter_user_id,
SUM(CASE WHEN bound_at >= ? AND bound_at < ? THEN 1 ELSE 0 END) AS invited_users,
COUNT(*) AS invited_users,
SUM(
CASE
WHEN reward_status = 'REWARDED'
AND rewarded_at >= ? AND rewarded_at < ?
AND rewarded_at < ?
THEN 1 ELSE 0
END
) AS rewarded_users
FROM referral_bindings
WHERE (bound_at >= ? AND bound_at < ?)
OR (
reward_status = 'REWARDED'
AND rewarded_at >= ? AND rewarded_at < ?
)
WHERE bound_at >= ? AND bound_at < ?
GROUP BY inviter_user_id
""",
range.arguments(repetitions = 4),
listOf(INSTANT_COLUMN_TYPE to range.until) + range.arguments(),
) { result ->
ReferralBindingAggregateRow(
inviterUserId = result.getString("inviter_user_id"),
@@ -225,14 +261,19 @@ class ExposedAdminStatsRepository(
private fun loadReferralCreditsByInviter(range: AdminStatsRange): Map<String, Long> =
queryRows(
"""
SELECT user_id, COALESCE(SUM(amount_delta), 0) AS earned_credits
FROM credit_ledger
WHERE created_at >= ? AND created_at < ?
AND entry_type = 'REFERRAL_INVITER'
AND amount_delta > 0
GROUP BY user_id
SELECT
r.inviter_user_id AS user_id,
COALESCE(SUM(l.amount_delta), 0) AS earned_credits
FROM referral_bindings r
JOIN credit_ledger l
ON l.reference_id = r.id
AND l.entry_type = 'REFERRAL_INVITER'
AND l.amount_delta > 0
WHERE r.bound_at >= ? AND r.bound_at < ?
AND l.created_at < ?
GROUP BY r.inviter_user_id
""",
range.arguments(),
range.arguments() + listOf(INSTANT_COLUMN_TYPE to range.until),
) { result ->
result.getString("user_id") to result.exactLong("earned_credits")
}.toMap()
@@ -312,13 +353,21 @@ private fun <T> queryRows(
sql: String,
arguments: List<Pair<IColumnType<*>, Any?>>,
transform: (ResultSet) -> T,
): List<T> = TransactionManager.current().exec(sql.trimIndent(), arguments) { result ->
buildList {
while (result.next()) {
add(transform(result))
}
): List<T> {
val normalized = sql.trimIndent()
val executable = if (normalized.startsWith("WITH ", ignoreCase = true)) {
"SELECT * FROM (\n$normalized\n) AS admin_stats_result"
} else {
normalized
}
} ?: emptyList()
return TransactionManager.current().exec(executable, arguments) { result ->
buildList {
while (result.next()) {
add(transform(result))
}
}
} ?: emptyList()
}
private fun ResultSet.exactLong(column: String): Long =
requireNotNull(getBigDecimal(column)) { "Aggregate column $column must not be null" }
@@ -9,10 +9,12 @@ import com.osglab.account.features.admin.stats.models.AdminAnalyticsFunnelStepDt
import com.osglab.account.features.admin.stats.models.AdminAnalyticsGrowthDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsGuardrailsDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsKeyboardUsageDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsLatencyBucketDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsMonetizationDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsNorthStarDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsPeriodDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsRateDto
import com.osglab.account.features.admin.stats.models.AdminAnalyticsReferralSignalsDto
import com.osglab.account.features.admin.stats.models.AdminProductAnalyticsDto
import com.osglab.account.features.admin.stats.repositories.AdminAnalyticsCountRow
import com.osglab.account.features.admin.stats.repositories.AdminAnalyticsWindow
@@ -100,13 +102,18 @@ class AdminProductAnalyticsService(
conversion7d = snapshot.monetization.conversion7d.toRate(),
conversion30d = snapshot.monetization.conversion30d.toRate(),
repeatPurchaseRate = snapshot.monetization.repeatPurchase.toRate(),
purchaseFunnel = listOf(
AdminAnalyticsFunnelStepDto("浏览购买页", snapshot.purchaseFunnel.viewed),
AdminAnalyticsFunnelStepDto("发起购买", snapshot.purchaseFunnel.started),
AdminAnalyticsFunnelStepDto("StoreKit 验证完成", snapshot.purchaseFunnel.verified),
),
cancelledUsers = snapshot.purchaseFunnel.cancelled,
),
growthFunnel = listOf(
AdminAnalyticsFunnelStepDto("首次启动", growth.opened),
AdminAnalyticsFunnelStepDto("完成注册", growth.registered),
AdminAnalyticsFunnelStepDto("已完成 24h 观察的新安装", growth.opened),
AdminAnalyticsFunnelStepDto("24 小时内完成注册", growth.registered),
AdminAnalyticsFunnelStepDto("24 小时内首次 AI 成功", growth.activated),
AdminAnalyticsFunnelStepDto("D7 再次使用 AI", growth.retainedD7),
AdminAnalyticsFunnelStepDto("首次购买", growth.purchased),
AdminAnalyticsFunnelStepDto("24 小时内完成首购", growth.purchased),
),
retention = snapshot.retention.map { cohort ->
AdminAnalyticsCohortDto(
@@ -160,17 +167,26 @@ class AdminProductAnalyticsService(
mixedLanguageSessions = keyboard.mixedLanguageSessions,
otherOnlySessions = keyboard.otherOnlySessions,
),
referralSignals = AdminAnalyticsReferralSignalsDto(
shared = referrals.shared,
opened = referrals.opened,
),
referralFunnel = listOf(
AdminAnalyticsFunnelStepDto("发起分享", referrals.shared),
AdminAnalyticsFunnelStepDto("打开邀请", referrals.opened),
AdminAnalyticsFunnelStepDto("完成绑定", referrals.bound),
AdminAnalyticsFunnelStepDto("首次 AI 成功", referrals.activated),
AdminAnalyticsFunnelStepDto("绑定后首次 AI 成功", referrals.activated),
AdminAnalyticsFunnelStepDto("完成奖励", referrals.rewarded),
),
guardrails = AdminAnalyticsGuardrailsDto(
clientAiSuccessRate = snapshot.guardrails.clientSuccess.toRate(),
managedSuccessRate = snapshot.guardrails.managedSuccess.toRate(),
creditBlockedUsers = snapshot.guardrails.creditBlockedUsers,
latencyBuckets = snapshot.latencyDistribution.map {
AdminAnalyticsLatencyBucketDto(
bucket = it.bucket,
successful = it.successful,
failed = it.failed,
)
},
),
)
}
@@ -12,7 +12,7 @@ import kotlinx.serialization.Serializable
@Serializable
data class AnalyticsBatchRequest(
val installationId: String,
val installationId: String? = null,
val events: List<AnalyticsEventRequest>,
) {
override fun toString(): String =
@@ -32,6 +32,8 @@ data class AnalyticsEventRequest(
val durationBucket: AnalyticsDurationBucket? = null,
val appVersion: String? = null,
val osVersion: String? = null,
// Transitional compatibility for clients released before installationId moved to the batch.
val installationId: String? = null,
) {
override fun toString(): String = "AnalyticsEventRequest([REDACTED])"
}
@@ -49,7 +49,7 @@ class DefaultAnalyticsService(
if (request.events.size !in MIN_BATCH_SIZE..MAX_BATCH_SIZE) {
throw InvalidRequestException("events must contain between 1 and 50 items")
}
val installationId = parseUuid(request.installationId, "installationId")
val installationId = parseUuid(resolveInstallationId(request), "installationId")
val now = clock.instant()
val events = request.events.map { validateAndMap(it, now) }
return repository.ingest(
@@ -62,6 +62,25 @@ class DefaultAnalyticsService(
)
}
private fun resolveInstallationId(request: AnalyticsBatchRequest): String {
val batchInstallationId = request.installationId
val eventInstallationIds = request.events.map(AnalyticsEventRequest::installationId)
if (batchInstallationId == null) {
if (eventInstallationIds.any { it == null }) {
throw InvalidRequestException("installationId is required")
}
val distinctIds = eventInstallationIds.filterNotNull().toSet()
if (distinctIds.size != 1) {
throw InvalidRequestException("event installationId values must match")
}
return distinctIds.single()
}
if (eventInstallationIds.filterNotNull().any { it != batchInstallationId }) {
throw InvalidRequestException("event installationId must match the batch")
}
return batchInstallationId
}
override suspend fun ingestKeyboardUsage(
accountId: UUID?,
request: KeyboardUsageBatchRequest,