wip
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
import { getPreviousMetric, groupByLabels } from '@openpanel/common';
|
||||
import type { ISerieDataItem } from '@openpanel/common';
|
||||
import { groupByLabels } from '@openpanel/common';
|
||||
import { alphabetIds } from '@openpanel/constants';
|
||||
import type {
|
||||
FinalChart,
|
||||
@@ -33,7 +33,7 @@ export async function executeChart(input: IReportInput): Promise<FinalChart> {
|
||||
// Handle subscription end date limit
|
||||
const endDate = await getOrganizationSubscriptionChartEndDate(
|
||||
input.projectId,
|
||||
normalized.endDate,
|
||||
normalized.endDate
|
||||
);
|
||||
if (endDate) {
|
||||
normalized.endDate = endDate;
|
||||
@@ -73,6 +73,7 @@ export async function executeChart(input: IReportInput): Promise<FinalChart> {
|
||||
executionPlan.definitions,
|
||||
includeAlphaIds,
|
||||
previousSeries,
|
||||
normalized.limit
|
||||
);
|
||||
|
||||
return response;
|
||||
@@ -83,7 +84,7 @@ export async function executeChart(input: IReportInput): Promise<FinalChart> {
|
||||
* Executes a simplified pipeline: normalize -> fetch aggregate -> format
|
||||
*/
|
||||
export async function executeAggregateChart(
|
||||
input: IReportInput,
|
||||
input: IReportInput
|
||||
): Promise<FinalChart> {
|
||||
// Stage 1: Normalize input
|
||||
const normalized = await normalize(input);
|
||||
@@ -91,7 +92,7 @@ export async function executeAggregateChart(
|
||||
// Handle subscription end date limit
|
||||
const endDate = await getOrganizationSubscriptionChartEndDate(
|
||||
input.projectId,
|
||||
normalized.endDate,
|
||||
normalized.endDate
|
||||
);
|
||||
if (endDate) {
|
||||
normalized.endDate = endDate;
|
||||
@@ -137,7 +138,7 @@ export async function executeAggregateChart(
|
||||
getAggregateChartSql(queryInput),
|
||||
{
|
||||
session_timezone: timezone,
|
||||
},
|
||||
}
|
||||
);
|
||||
|
||||
// Fallback: if no results with breakdowns, try without breakdowns
|
||||
@@ -149,7 +150,7 @@ export async function executeAggregateChart(
|
||||
}),
|
||||
{
|
||||
session_timezone: timezone,
|
||||
},
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
@@ -262,7 +263,7 @@ export async function executeAggregateChart(
|
||||
getAggregateChartSql(queryInput),
|
||||
{
|
||||
session_timezone: timezone,
|
||||
},
|
||||
}
|
||||
);
|
||||
|
||||
if (queryResult.length === 0 && normalized.breakdowns.length > 0) {
|
||||
@@ -273,7 +274,7 @@ export async function executeAggregateChart(
|
||||
}),
|
||||
{
|
||||
session_timezone: timezone,
|
||||
},
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
@@ -344,7 +345,7 @@ export async function executeAggregateChart(
|
||||
normalized.series,
|
||||
includeAlphaIds,
|
||||
previousSeries,
|
||||
normalized.limit,
|
||||
normalized.limit
|
||||
);
|
||||
|
||||
return response;
|
||||
|
||||
@@ -181,16 +181,8 @@ export async function getGroupTypes(projectId: string): Promise<string[]> {
|
||||
}
|
||||
|
||||
export async function createGroup(input: IServiceUpsertGroup) {
|
||||
const { id, projectId, type, name, properties = {} } = input;
|
||||
await writeGroupToCh({
|
||||
id,
|
||||
projectId,
|
||||
type,
|
||||
name,
|
||||
properties: properties as Record<string, string>,
|
||||
createdAt: new Date(),
|
||||
});
|
||||
return getGroupById(id, projectId);
|
||||
await upsertGroup(input);
|
||||
return getGroupById(input.id, input.projectId);
|
||||
}
|
||||
|
||||
export async function updateGroup(
|
||||
@@ -248,6 +240,49 @@ export async function getGroupPropertyKeys(
|
||||
return rows.map((r) => r.key).sort();
|
||||
}
|
||||
|
||||
export type IServiceGroupStats = {
|
||||
groupId: string;
|
||||
memberCount: number;
|
||||
lastActiveAt: Date | null;
|
||||
};
|
||||
|
||||
export async function getGroupStats(
|
||||
projectId: string,
|
||||
groupIds: string[]
|
||||
): Promise<Map<string, IServiceGroupStats>> {
|
||||
if (groupIds.length === 0) {
|
||||
return new Map();
|
||||
}
|
||||
|
||||
const rows = await chQuery<{
|
||||
group_id: string;
|
||||
member_count: number;
|
||||
last_active_at: string;
|
||||
}>(`
|
||||
SELECT
|
||||
g AS group_id,
|
||||
uniqExact(profile_id) AS member_count,
|
||||
max(created_at) AS last_active_at
|
||||
FROM ${TABLE_NAMES.events}
|
||||
ARRAY JOIN groups AS g
|
||||
WHERE project_id = ${sqlstring.escape(projectId)}
|
||||
AND g IN (${groupIds.map((id) => sqlstring.escape(id)).join(',')})
|
||||
AND profile_id != device_id
|
||||
GROUP BY g
|
||||
`);
|
||||
|
||||
return new Map(
|
||||
rows.map((r) => [
|
||||
r.group_id,
|
||||
{
|
||||
groupId: r.group_id,
|
||||
memberCount: r.member_count,
|
||||
lastActiveAt: r.last_active_at ? new Date(r.last_active_at) : null,
|
||||
},
|
||||
])
|
||||
);
|
||||
}
|
||||
|
||||
export async function getGroupsByIds(
|
||||
projectId: string,
|
||||
ids: string[]
|
||||
@@ -284,26 +319,12 @@ export async function getGroupMemberProfiles({
|
||||
? `AND (email ILIKE ${sqlstring.escape(`%${search.trim()}%`)} OR first_name ILIKE ${sqlstring.escape(`%${search.trim()}%`)} OR last_name ILIKE ${sqlstring.escape(`%${search.trim()}%`)})`
|
||||
: '';
|
||||
|
||||
const countResult = await chQuery<{ count: number }>(`
|
||||
SELECT count() AS count
|
||||
FROM (
|
||||
SELECT profile_id
|
||||
FROM ${TABLE_NAMES.events}
|
||||
WHERE project_id = ${sqlstring.escape(projectId)}
|
||||
AND has(groups, ${sqlstring.escape(groupId)})
|
||||
AND profile_id != device_id
|
||||
GROUP BY profile_id
|
||||
) gm
|
||||
INNER JOIN (
|
||||
SELECT id FROM ${TABLE_NAMES.profiles} FINAL
|
||||
WHERE project_id = ${sqlstring.escape(projectId)}
|
||||
${searchCondition}
|
||||
) p ON p.id = gm.profile_id
|
||||
`);
|
||||
const count = countResult[0]?.count ?? 0;
|
||||
|
||||
const idRows = await chQuery<{ profile_id: string }>(`
|
||||
SELECT gm.profile_id
|
||||
// count() OVER () is evaluated after JOINs/WHERE but before LIMIT,
|
||||
// so we get the total match count and the paginated IDs in one query.
|
||||
const rows = await chQuery<{ profile_id: string; total_count: number }>(`
|
||||
SELECT
|
||||
gm.profile_id,
|
||||
count() OVER () AS total_count
|
||||
FROM (
|
||||
SELECT profile_id, max(created_at) AS last_seen
|
||||
FROM ${TABLE_NAMES.events}
|
||||
@@ -311,7 +332,6 @@ export async function getGroupMemberProfiles({
|
||||
AND has(groups, ${sqlstring.escape(groupId)})
|
||||
AND profile_id != device_id
|
||||
GROUP BY profile_id
|
||||
ORDER BY last_seen DESC
|
||||
) gm
|
||||
INNER JOIN (
|
||||
SELECT id FROM ${TABLE_NAMES.profiles} FINAL
|
||||
@@ -322,10 +342,14 @@ export async function getGroupMemberProfiles({
|
||||
LIMIT ${take}
|
||||
OFFSET ${offset}
|
||||
`);
|
||||
const profileIds = idRows.map((r) => r.profile_id);
|
||||
|
||||
const count = rows[0]?.total_count ?? 0;
|
||||
const profileIds = rows.map((r) => r.profile_id);
|
||||
|
||||
if (profileIds.length === 0) {
|
||||
return { data: [], count };
|
||||
}
|
||||
|
||||
const profiles = await getProfiles(profileIds, projectId);
|
||||
const byId = new Map(profiles.map((p) => [p.id, p]));
|
||||
const data = profileIds
|
||||
|
||||
@@ -7,6 +7,7 @@ import {
|
||||
getGroupListCount,
|
||||
getGroupMemberProfiles,
|
||||
getGroupPropertyKeys,
|
||||
getGroupStats,
|
||||
getGroupsByIds,
|
||||
getGroupTypes,
|
||||
TABLE_NAMES,
|
||||
@@ -32,7 +33,18 @@ export const groupRouter = createTRPCRouter({
|
||||
getGroupList(input),
|
||||
getGroupListCount(input),
|
||||
]);
|
||||
return { data, meta: { count, take: input.take } };
|
||||
const stats = await getGroupStats(
|
||||
input.projectId,
|
||||
data.map((g) => g.id)
|
||||
);
|
||||
return {
|
||||
data: data.map((g) => ({
|
||||
...g,
|
||||
memberCount: stats.get(g.id)?.memberCount ?? 0,
|
||||
lastActiveAt: stats.get(g.id)?.lastActiveAt ?? null,
|
||||
})),
|
||||
meta: { count, take: input.take },
|
||||
};
|
||||
}),
|
||||
|
||||
byId: protectedProcedure
|
||||
@@ -160,6 +172,60 @@ export const groupRouter = createTRPCRouter({
|
||||
};
|
||||
}),
|
||||
|
||||
mostEvents: protectedProcedure
|
||||
.input(z.object({ id: z.string(), projectId: z.string() }))
|
||||
.query(async ({ input: { id, projectId } }) => {
|
||||
return chQuery<{ count: number; name: string }>(`
|
||||
SELECT count() as count, name
|
||||
FROM ${TABLE_NAMES.events}
|
||||
WHERE project_id = ${sqlstring.escape(projectId)}
|
||||
AND has(groups, ${sqlstring.escape(id)})
|
||||
AND name NOT IN ('screen_view', 'session_start', 'session_end')
|
||||
GROUP BY name
|
||||
ORDER BY count DESC
|
||||
LIMIT 10
|
||||
`);
|
||||
}),
|
||||
|
||||
popularRoutes: protectedProcedure
|
||||
.input(z.object({ id: z.string(), projectId: z.string() }))
|
||||
.query(async ({ input: { id, projectId } }) => {
|
||||
return chQuery<{ count: number; path: string }>(`
|
||||
SELECT count() as count, path
|
||||
FROM ${TABLE_NAMES.events}
|
||||
WHERE project_id = ${sqlstring.escape(projectId)}
|
||||
AND has(groups, ${sqlstring.escape(id)})
|
||||
AND name = 'screen_view'
|
||||
GROUP BY path
|
||||
ORDER BY count DESC
|
||||
LIMIT 10
|
||||
`);
|
||||
}),
|
||||
|
||||
memberGrowth: protectedProcedure
|
||||
.input(z.object({ id: z.string(), projectId: z.string() }))
|
||||
.query(async ({ input: { id, projectId } }) => {
|
||||
return chQuery<{ date: string; count: number }>(`
|
||||
SELECT
|
||||
toDate(toStartOfDay(min_date)) AS date,
|
||||
count() AS count
|
||||
FROM (
|
||||
SELECT profile_id, min(created_at) AS min_date
|
||||
FROM ${TABLE_NAMES.events}
|
||||
WHERE project_id = ${sqlstring.escape(projectId)}
|
||||
AND has(groups, ${sqlstring.escape(id)})
|
||||
AND profile_id != device_id
|
||||
AND created_at >= now() - INTERVAL 30 DAY
|
||||
GROUP BY profile_id
|
||||
)
|
||||
GROUP BY date
|
||||
ORDER BY date ASC WITH FILL
|
||||
FROM toDate(now() - INTERVAL 29 DAY)
|
||||
TO toDate(now() + INTERVAL 1 DAY)
|
||||
STEP 1
|
||||
`);
|
||||
}),
|
||||
|
||||
properties: protectedProcedure
|
||||
.input(z.object({ projectId: z.string() }))
|
||||
.query(async ({ input: { projectId } }) => {
|
||||
|
||||
Reference in New Issue
Block a user