feat: create workspaces and user workspaces
This commit is contained in:
@@ -28,7 +28,7 @@
|
||||
},
|
||||
"yarnLockChecksum": "31c83d74119b810f2c956d41c2ee3576",
|
||||
"packageJsonChecksum": "1807c41d2a51ee36385ece5d35b3ef4f",
|
||||
"apiClientChecksum": "3016815faabfbe0aa66c15e37330545f"
|
||||
"apiClientChecksum": null
|
||||
},
|
||||
"objects": [
|
||||
{
|
||||
@@ -1579,8 +1579,8 @@
|
||||
"logicFunctions": [
|
||||
{
|
||||
"universalIdentifier": "3897e059-715e-4a4b-b165-c44f17d2e30a",
|
||||
"name": "product-data-create-users",
|
||||
"description": "A simple logic function",
|
||||
"name": "sync-product-data",
|
||||
"description": "Syncs cloud users, cloud workspaces, and cloud user workspaces from ClickHouse",
|
||||
"timeoutSeconds": 120,
|
||||
"cronTriggerSettings": {
|
||||
"pattern": "*/10 * * * *"
|
||||
@@ -1590,9 +1590,9 @@
|
||||
"properties": {}
|
||||
},
|
||||
"handlerName": "default.config.handler",
|
||||
"sourceHandlerPath": "src/logic-functions/create-users.ts",
|
||||
"builtHandlerPath": "src/logic-functions/create-users.mjs",
|
||||
"builtHandlerChecksum": "dc59169411b3508d7a008e213b31f3e2"
|
||||
"sourceHandlerPath": "src/logic-functions/sync-product-data.ts",
|
||||
"builtHandlerPath": "src/logic-functions/sync-product-data.mjs",
|
||||
"builtHandlerChecksum": "9e7a645ec39adf0585692bf9273efec6"
|
||||
}
|
||||
],
|
||||
"frontComponents": [],
|
||||
|
||||
-7
File diff suppressed because one or more lines are too long
+14203
-13929
File diff suppressed because it is too large
Load Diff
+7
File diff suppressed because one or more lines are too long
@@ -1,350 +0,0 @@
|
||||
import { defineLogicFunction } from 'twenty-sdk';
|
||||
import Twenty, { enumCloudUser2ActivityStatusEnum } from 'twenty-sdk/generated';
|
||||
import { z } from 'zod';
|
||||
|
||||
const client = new Twenty();
|
||||
|
||||
const applicationConfigSchema = z.object({
|
||||
CLICKHOUSE_DATABASE: z.string().nonempty(),
|
||||
CLICKHOUSE_URL: z.url(),
|
||||
CLICKHOUSE_USERNAME: z.string().nonempty(),
|
||||
CLICKHOUSE_PASSWORD: z.string().nonempty(),
|
||||
})
|
||||
|
||||
const getApplicationConfig = () => {
|
||||
const env = applicationConfigSchema.parse({
|
||||
CLICKHOUSE_DATABASE: process.env.CLICKHOUSE_DATABASE,
|
||||
CLICKHOUSE_URL: process.env.CLICKHOUSE_URL,
|
||||
CLICKHOUSE_USERNAME: process.env.CLICKHOUSE_USERNAME,
|
||||
CLICKHOUSE_PASSWORD: process.env.CLICKHOUSE_PASSWORD,
|
||||
});
|
||||
|
||||
return {
|
||||
clickHouseDatabase: env.CLICKHOUSE_DATABASE,
|
||||
clickHouseUrl: env.CLICKHOUSE_URL,
|
||||
clickHouseUsername: env.CLICKHOUSE_USERNAME,
|
||||
clickHousePassword: env.CLICKHOUSE_PASSWORD,
|
||||
};
|
||||
};
|
||||
|
||||
const clickHouseUserSchema = z.object({
|
||||
userId: z.uuid(),
|
||||
firstName: z.string(),
|
||||
lastName: z.string(),
|
||||
email: z.email(),
|
||||
fullName: z.string(),
|
||||
isEmailVerified: z.boolean(),
|
||||
disabled: z.boolean(),
|
||||
canImpersonate: z.boolean(),
|
||||
canAccessFullAdminPanel: z.boolean(),
|
||||
createdAt: z.string(),
|
||||
updatedAt: z.string(),
|
||||
deletedAt: z.string(),
|
||||
locale: z.string(),
|
||||
createdDate: z.string(),
|
||||
workspaceCount: z.coerce.number(),
|
||||
workspaceIds: z.string(),
|
||||
workspaceDomains: z.string(),
|
||||
firstActivityDate: z.string(),
|
||||
lastActivityDate: z.string(),
|
||||
lastWorkspaceId: z.string(),
|
||||
totalPageviews: z.coerce.number(),
|
||||
pageviewsLast30d: z.coerce.number(),
|
||||
pageviewsLast7d: z.coerce.number(),
|
||||
pageviewsLast24h: z.coerce.number(),
|
||||
userAgeDays: z.coerce.number(),
|
||||
daysSinceLastActivity: z.coerce.number(),
|
||||
isActiveLast30d: z.boolean(),
|
||||
isActiveLast7d: z.boolean(),
|
||||
isActiveLast24h: z.boolean(),
|
||||
activityStatus: z
|
||||
.string()
|
||||
.transform((val) => val.toUpperCase())
|
||||
.pipe(z.enum(enumCloudUser2ActivityStatusEnum)),
|
||||
avgDailyPageviewsLast30d: z.coerce.number(),
|
||||
isTwenty: z.coerce.boolean(),
|
||||
maxWorkspaceMembers: z.coerce.number(),
|
||||
inTrial: z.boolean(),
|
||||
});
|
||||
|
||||
type ClickHouseUser = z.infer<typeof clickHouseUserSchema>;
|
||||
|
||||
const fetchUsersFromClickHouse = async (): Promise<{
|
||||
users: ClickHouseUser[];
|
||||
}> => {
|
||||
const {
|
||||
clickHouseDatabase,
|
||||
clickHouseUrl,
|
||||
clickHouseUsername,
|
||||
clickHousePassword,
|
||||
} = getApplicationConfig();
|
||||
|
||||
// const nowDate = 'now()'
|
||||
const nowDate = "'2026-02-04 14:28:52.000'";
|
||||
|
||||
const findUsersWithRecentActivity = `
|
||||
SELECT
|
||||
*
|
||||
FROM
|
||||
${clickHouseDatabase}.user
|
||||
WHERE
|
||||
lastActivityDate >= ${nowDate} - INTERVAL 500 MINUTE
|
||||
AND
|
||||
lastActivityDate <= ${nowDate}
|
||||
FORMAT
|
||||
JSONEachRow;
|
||||
`;
|
||||
|
||||
const res = await fetch(clickHouseUrl, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization:
|
||||
'Basic ' +
|
||||
Buffer.from(`${clickHouseUsername}:${clickHousePassword}`).toString(
|
||||
'base64',
|
||||
),
|
||||
'Content-Type': 'text/plain',
|
||||
},
|
||||
body: findUsersWithRecentActivity,
|
||||
});
|
||||
|
||||
if (!res.ok) {
|
||||
const errText = await res.text();
|
||||
throw new Error(`ClickHouse error: ${res.status} - ${errText}`);
|
||||
}
|
||||
|
||||
/**
|
||||
* Format is a list of JSON objects, one per line, so we need to split by line and parse each line as JSON.
|
||||
*/
|
||||
const text = await res.text();
|
||||
|
||||
const rows = text
|
||||
.trim()
|
||||
.split('\n')
|
||||
.filter((line) => line.trim()) // Filter out empty lines
|
||||
.map((line) => JSON.parse(line));
|
||||
|
||||
const users = z.array(clickHouseUserSchema).parse(rows);
|
||||
|
||||
return { users };
|
||||
};
|
||||
|
||||
const fetchAllPeopleFromTwentyByEmail = async (emails: string[]) => {
|
||||
if (emails.length === 0) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const allPeople = await client.query({
|
||||
people: {
|
||||
edges: {
|
||||
node: {
|
||||
id: true,
|
||||
emails: {
|
||||
primaryEmail: true,
|
||||
},
|
||||
},
|
||||
},
|
||||
__args: {
|
||||
filter: {
|
||||
emails: {
|
||||
primaryEmail: {
|
||||
in: emails,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
return allPeople.people?.edges.map((edge) => edge.node) ?? [];
|
||||
};
|
||||
|
||||
const buildCloudUserInput = ({
|
||||
user,
|
||||
personId,
|
||||
}: {
|
||||
user: ClickHouseUser;
|
||||
personId: string;
|
||||
}) => ({
|
||||
id: user.userId,
|
||||
name: user.fullName,
|
||||
email: {
|
||||
primaryEmail: user.email,
|
||||
},
|
||||
personId,
|
||||
fullName: {
|
||||
lastName: user.lastName,
|
||||
firstName: user.firstName,
|
||||
},
|
||||
isTwenty: user.isTwenty,
|
||||
userTenure: user.userAgeDays,
|
||||
isActiveL7d: user.isActiveLast7d,
|
||||
isActiveL24h: user.isActiveLast24h,
|
||||
isActiveL30d: user.isActiveLast30d,
|
||||
pageViewsL7d: user.pageviewsLast7d,
|
||||
pageViewsL24h: user.pageviewsLast24h,
|
||||
pageViewsL30d: user.pageviewsLast30d,
|
||||
activityStatus: user.activityStatus,
|
||||
workspaceCount: user.workspaceCount,
|
||||
lastActivityDate: user.lastActivityDate,
|
||||
dataLastUpdatedAt: new Date().toISOString(),
|
||||
avgDailyPageviewsLast30d: user.avgDailyPageviewsLast30d,
|
||||
daysSinceLastActivity: user.daysSinceLastActivity,
|
||||
});
|
||||
|
||||
const handler = async (): Promise<{ message: string }> => {
|
||||
try {
|
||||
const { users } = await fetchUsersFromClickHouse();
|
||||
|
||||
console.log('fetch users from clickhouse', users);
|
||||
|
||||
const emails = users.map((user) => user.email);
|
||||
|
||||
console.log('fetch people from twenty with emails', emails);
|
||||
|
||||
const people = await fetchAllPeopleFromTwentyByEmail(emails);
|
||||
|
||||
console.log('fetched people from twenty', people);
|
||||
|
||||
// Build a map of email -> personId from existing people
|
||||
const emailToPersonId = new Map<string, string>();
|
||||
|
||||
for (const person of people) {
|
||||
if (person.emails?.primaryEmail !== undefined) {
|
||||
emailToPersonId.set(
|
||||
person.emails.primaryEmail.toLowerCase(),
|
||||
person.id,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Partition users into those with/without an existing person
|
||||
const usersWithoutPerson = users.filter(
|
||||
(user) => !emailToPersonId.has(user.email.toLowerCase()),
|
||||
);
|
||||
|
||||
// Step 1: Batch-create missing people
|
||||
const newlyCreatedPersonIds: string[] = [];
|
||||
|
||||
if (usersWithoutPerson.length > 0) {
|
||||
console.log(`Batch-creating ${usersWithoutPerson.length} people records`);
|
||||
|
||||
const createPeopleResult = await client.mutation({
|
||||
createPeople: {
|
||||
__args: {
|
||||
data: usersWithoutPerson.map((user) => ({
|
||||
name: {
|
||||
firstName: user.firstName,
|
||||
lastName: user.lastName,
|
||||
},
|
||||
emails: {
|
||||
primaryEmail: user.email,
|
||||
},
|
||||
})),
|
||||
},
|
||||
id: true,
|
||||
emails: {
|
||||
primaryEmail: true,
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
const createdPeople = createPeopleResult.createPeople ?? [];
|
||||
|
||||
for (const person of createdPeople) {
|
||||
if (
|
||||
person.emails?.primaryEmail === undefined ||
|
||||
person.id === undefined
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
|
||||
newlyCreatedPersonIds.push(person.id);
|
||||
emailToPersonId.set(
|
||||
person.emails.primaryEmail.toLowerCase(),
|
||||
person.id,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
// Step 2: Batch-upsert all cloud users, with rollback on failure
|
||||
const cloudUserInputs = users.map((user) => {
|
||||
const personId = emailToPersonId.get(user.email.toLowerCase());
|
||||
|
||||
if (personId === undefined) {
|
||||
throw new Error(
|
||||
`No personId found for user ${user.email} — this should not happen`,
|
||||
);
|
||||
}
|
||||
|
||||
return buildCloudUserInput({ user, personId });
|
||||
});
|
||||
|
||||
try {
|
||||
console.log(`Batch-upserting ${cloudUserInputs.length} cloud users`);
|
||||
|
||||
await client.mutation({
|
||||
createCloudUsers2: {
|
||||
__args: {
|
||||
data: cloudUserInputs,
|
||||
upsert: true,
|
||||
},
|
||||
__scalar: true,
|
||||
},
|
||||
});
|
||||
} catch (cloudUserError) {
|
||||
console.log(
|
||||
'Cloud user upsert failed, rolling back newly created people',
|
||||
cloudUserError,
|
||||
);
|
||||
|
||||
// Rollback: hard-delete people that were created in step 1
|
||||
if (newlyCreatedPersonIds.length > 0) {
|
||||
try {
|
||||
await client.mutation({
|
||||
destroyPeople: {
|
||||
__args: {
|
||||
filter: {
|
||||
id: {
|
||||
in: newlyCreatedPersonIds,
|
||||
},
|
||||
},
|
||||
},
|
||||
id: true,
|
||||
},
|
||||
});
|
||||
|
||||
console.log(
|
||||
`Rolled back ${newlyCreatedPersonIds.length} newly created people`,
|
||||
);
|
||||
} catch (rollbackError) {
|
||||
console.log(
|
||||
'Rollback of newly created people also failed',
|
||||
rollbackError,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
throw cloudUserError;
|
||||
}
|
||||
|
||||
return {
|
||||
message: `Successfully processed ${users.length} users from ClickHouse`,
|
||||
};
|
||||
} catch (err) {
|
||||
console.log(err);
|
||||
|
||||
throw err;
|
||||
}
|
||||
};
|
||||
|
||||
export default defineLogicFunction({
|
||||
universalIdentifier: '3897e059-715e-4a4b-b165-c44f17d2e30a',
|
||||
name: 'product-data-create-users',
|
||||
description: 'A simple logic function',
|
||||
timeoutSeconds: 120,
|
||||
handler,
|
||||
cronTriggerSettings: {
|
||||
pattern: '*/10 * * * *',
|
||||
},
|
||||
});
|
||||
+49
@@ -0,0 +1,49 @@
|
||||
import { defineLogicFunction } from 'twenty-sdk';
|
||||
|
||||
import { syncCloudUserWorkspaces } from 'src/logic-functions/sync-product-data/sync-cloud-user-workspaces';
|
||||
import { syncCloudUsers } from 'src/logic-functions/sync-product-data/sync-cloud-users';
|
||||
import { syncCloudWorkspaces } from 'src/logic-functions/sync-product-data/sync-cloud-workspaces';
|
||||
|
||||
const handler = async (): Promise<{ message: string }> => {
|
||||
try {
|
||||
console.log('Starting product data sync');
|
||||
|
||||
const cloudUsersResult = await syncCloudUsers();
|
||||
|
||||
console.log(
|
||||
`Cloud users sync complete: ${cloudUsersResult.syncedCount} users`,
|
||||
);
|
||||
|
||||
const cloudWorkspacesResult = await syncCloudWorkspaces();
|
||||
|
||||
console.log(
|
||||
`Cloud workspaces sync complete: ${cloudWorkspacesResult.syncedCount} workspaces`,
|
||||
);
|
||||
|
||||
const cloudUserWorkspacesResult = await syncCloudUserWorkspaces();
|
||||
|
||||
console.log(
|
||||
`Cloud user workspaces sync complete: ${cloudUserWorkspacesResult.syncedCount} user workspaces`,
|
||||
);
|
||||
|
||||
return {
|
||||
message: `Product data sync complete — ${cloudUsersResult.syncedCount} users, ${cloudWorkspacesResult.syncedCount} workspaces, ${cloudUserWorkspacesResult.syncedCount} user workspaces`,
|
||||
};
|
||||
} catch (err) {
|
||||
console.log(err);
|
||||
|
||||
throw err;
|
||||
}
|
||||
};
|
||||
|
||||
export default defineLogicFunction({
|
||||
universalIdentifier: '3897e059-715e-4a4b-b165-c44f17d2e30a',
|
||||
name: 'sync-product-data',
|
||||
description:
|
||||
'Syncs cloud users, cloud workspaces, and cloud user workspaces from ClickHouse',
|
||||
timeoutSeconds: 120,
|
||||
handler,
|
||||
cronTriggerSettings: {
|
||||
pattern: '*/10 * * * *',
|
||||
},
|
||||
});
|
||||
+93
@@ -0,0 +1,93 @@
|
||||
import { z } from 'zod';
|
||||
|
||||
import { getApplicationConfig } from 'src/shared/application-config';
|
||||
import { fetchFromClickHouse } from 'src/shared/clickhouse-client';
|
||||
import { twentyClient } from 'src/shared/twenty-client';
|
||||
|
||||
const clickHouseUserWorkspaceSchema = z.object({
|
||||
id: z.string(),
|
||||
workspaceId: z.string(),
|
||||
userId: z.string(),
|
||||
createdAt: z.string(),
|
||||
updatedAt: z.string(),
|
||||
deletedAt: z.string(),
|
||||
is_active_membership: z.coerce.boolean(),
|
||||
});
|
||||
|
||||
type ClickHouseUserWorkspace = z.infer<typeof clickHouseUserWorkspaceSchema>;
|
||||
|
||||
const fetchUserWorkspacesFromClickHouse = async (): Promise<
|
||||
ClickHouseUserWorkspace[]
|
||||
> => {
|
||||
const { clickHouseDatabase } = getApplicationConfig();
|
||||
|
||||
// const nowDate = 'now()'
|
||||
const nowDate = "'2026-02-04 14:28:52.000'";
|
||||
|
||||
const query = `
|
||||
SELECT
|
||||
*
|
||||
FROM (
|
||||
SELECT
|
||||
*,
|
||||
row_number() OVER (PARTITION BY id ORDER BY updatedAt DESC) AS rn
|
||||
FROM
|
||||
${clickHouseDatabase}.user_workspace
|
||||
WHERE
|
||||
updatedAt >= ${nowDate} - INTERVAL 500 MINUTE
|
||||
AND
|
||||
updatedAt <= ${nowDate}
|
||||
)
|
||||
WHERE
|
||||
rn = 1
|
||||
FORMAT
|
||||
JSONEachRow;
|
||||
`;
|
||||
|
||||
return fetchFromClickHouse(query, clickHouseUserWorkspaceSchema);
|
||||
};
|
||||
|
||||
const buildCloudUserWorkspaceInput = (
|
||||
userWorkspace: ClickHouseUserWorkspace,
|
||||
) => ({
|
||||
id: userWorkspace.id,
|
||||
twentyUserIdentifier: userWorkspace.userId,
|
||||
twentyWorkspaceIdentifier: userWorkspace.workspaceId,
|
||||
idOfTheUserWorkspace: userWorkspace.id,
|
||||
cloudUser2Id: userWorkspace.userId,
|
||||
cloudWorkspace2Id: userWorkspace.workspaceId,
|
||||
});
|
||||
|
||||
export const syncCloudUserWorkspaces = async (): Promise<{
|
||||
syncedCount: number;
|
||||
}> => {
|
||||
const userWorkspaces = await fetchUserWorkspacesFromClickHouse();
|
||||
|
||||
console.log(
|
||||
`Fetched ${userWorkspaces.length} user workspaces from ClickHouse`,
|
||||
);
|
||||
|
||||
if (userWorkspaces.length === 0) {
|
||||
return { syncedCount: 0 };
|
||||
}
|
||||
|
||||
const cloudUserWorkspaceInputs = userWorkspaces.map(
|
||||
buildCloudUserWorkspaceInput,
|
||||
);
|
||||
|
||||
console.log(
|
||||
`Batch-upserting ${cloudUserWorkspaceInputs.length} cloud user workspaces`,
|
||||
);
|
||||
|
||||
await twentyClient.mutation({
|
||||
createCloudUserWorkspaces2: {
|
||||
__args: {
|
||||
data: cloudUserWorkspaceInputs,
|
||||
upsert: true,
|
||||
},
|
||||
__scalar: true,
|
||||
},
|
||||
});
|
||||
|
||||
return { syncedCount: userWorkspaces.length };
|
||||
};
|
||||
+268
@@ -0,0 +1,268 @@
|
||||
import { enumCloudUser2ActivityStatusEnum } from 'twenty-sdk/generated';
|
||||
import { z } from 'zod';
|
||||
|
||||
import { getApplicationConfig } from 'src/shared/application-config';
|
||||
import { fetchFromClickHouse } from 'src/shared/clickhouse-client';
|
||||
import { clickHouseDateToIso } from 'src/shared/clickhouse-date-to-iso';
|
||||
import { twentyClient } from 'src/shared/twenty-client';
|
||||
|
||||
const clickHouseUserSchema = z.object({
|
||||
userId: z.uuid(),
|
||||
firstName: z.string(),
|
||||
lastName: z.string(),
|
||||
email: z.email(),
|
||||
fullName: z.string(),
|
||||
isEmailVerified: z.boolean(),
|
||||
disabled: z.boolean(),
|
||||
canImpersonate: z.boolean(),
|
||||
canAccessFullAdminPanel: z.boolean(),
|
||||
createdAt: z.string(),
|
||||
updatedAt: z.string(),
|
||||
deletedAt: z.string(),
|
||||
locale: z.string(),
|
||||
createdDate: z.string(),
|
||||
workspaceCount: z.coerce.number(),
|
||||
workspaceIds: z.string(),
|
||||
workspaceDomains: z.string(),
|
||||
firstActivityDate: z.string(),
|
||||
lastActivityDate: z.string(),
|
||||
lastWorkspaceId: z.string(),
|
||||
totalPageviews: z.coerce.number(),
|
||||
pageviewsLast30d: z.coerce.number(),
|
||||
pageviewsLast7d: z.coerce.number(),
|
||||
pageviewsLast24h: z.coerce.number(),
|
||||
userAgeDays: z.coerce.number(),
|
||||
daysSinceLastActivity: z.coerce.number(),
|
||||
isActiveLast30d: z.boolean(),
|
||||
isActiveLast7d: z.boolean(),
|
||||
isActiveLast24h: z.boolean(),
|
||||
activityStatus: z
|
||||
.string()
|
||||
.transform((val) => val.toUpperCase())
|
||||
.pipe(z.enum(enumCloudUser2ActivityStatusEnum)),
|
||||
avgDailyPageviewsLast30d: z.coerce.number(),
|
||||
isTwenty: z.coerce.boolean(),
|
||||
maxWorkspaceMembers: z.coerce.number(),
|
||||
inTrial: z.boolean(),
|
||||
});
|
||||
|
||||
type ClickHouseUser = z.infer<typeof clickHouseUserSchema>;
|
||||
|
||||
const fetchUsersFromClickHouse = async (): Promise<ClickHouseUser[]> => {
|
||||
const { clickHouseDatabase } = getApplicationConfig();
|
||||
|
||||
// const nowDate = 'now()'
|
||||
const nowDate = "'2026-02-04 14:28:52.000'";
|
||||
|
||||
const query = `
|
||||
SELECT
|
||||
*
|
||||
FROM
|
||||
${clickHouseDatabase}.user
|
||||
WHERE
|
||||
lastActivityDate >= ${nowDate} - INTERVAL 500 MINUTE
|
||||
AND
|
||||
lastActivityDate <= ${nowDate}
|
||||
FORMAT
|
||||
JSONEachRow;
|
||||
`;
|
||||
|
||||
return fetchFromClickHouse(query, clickHouseUserSchema);
|
||||
};
|
||||
|
||||
const fetchAllPeopleFromTwentyByEmail = async (emails: string[]) => {
|
||||
if (emails.length === 0) {
|
||||
return [];
|
||||
}
|
||||
|
||||
const allPeople = await twentyClient.query({
|
||||
people: {
|
||||
edges: {
|
||||
node: {
|
||||
id: true,
|
||||
emails: {
|
||||
primaryEmail: true,
|
||||
},
|
||||
},
|
||||
},
|
||||
__args: {
|
||||
filter: {
|
||||
emails: {
|
||||
primaryEmail: {
|
||||
in: emails,
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
return allPeople.people?.edges.map((edge) => edge.node) ?? [];
|
||||
};
|
||||
|
||||
const buildCloudUserInput = ({
|
||||
user,
|
||||
personId,
|
||||
}: {
|
||||
user: ClickHouseUser;
|
||||
personId: string;
|
||||
}) => ({
|
||||
id: user.userId,
|
||||
name: user.fullName,
|
||||
email: {
|
||||
primaryEmail: user.email,
|
||||
},
|
||||
personId,
|
||||
fullName: {
|
||||
lastName: user.lastName,
|
||||
firstName: user.firstName,
|
||||
},
|
||||
isTwenty: user.isTwenty,
|
||||
userTenure: user.userAgeDays,
|
||||
isActiveL7d: user.isActiveLast7d,
|
||||
isActiveL24h: user.isActiveLast24h,
|
||||
isActiveL30d: user.isActiveLast30d,
|
||||
pageViewsL7d: user.pageviewsLast7d,
|
||||
pageViewsL24h: user.pageviewsLast24h,
|
||||
pageViewsL30d: user.pageviewsLast30d,
|
||||
activityStatus: user.activityStatus,
|
||||
workspaceCount: user.workspaceCount,
|
||||
lastActivityDate: clickHouseDateToIso(user.lastActivityDate),
|
||||
dataLastUpdatedAt: new Date().toISOString(),
|
||||
avgDailyPageviewsLast30d: user.avgDailyPageviewsLast30d,
|
||||
daysSinceLastActivity: user.daysSinceLastActivity,
|
||||
});
|
||||
|
||||
export const syncCloudUsers = async (): Promise<{
|
||||
syncedCount: number;
|
||||
}> => {
|
||||
const users = await fetchUsersFromClickHouse();
|
||||
|
||||
console.log('Fetched users from ClickHouse', users);
|
||||
|
||||
const emails = users.map((user) => user.email);
|
||||
|
||||
console.log('Fetching people from Twenty with emails', emails);
|
||||
|
||||
const people = await fetchAllPeopleFromTwentyByEmail(emails);
|
||||
|
||||
console.log('Fetched people from Twenty', people);
|
||||
|
||||
// Build a map of email -> personId from existing people
|
||||
const emailToPersonId = new Map<string, string>();
|
||||
|
||||
for (const person of people) {
|
||||
if (person.emails?.primaryEmail !== undefined) {
|
||||
emailToPersonId.set(person.emails.primaryEmail.toLowerCase(), person.id);
|
||||
}
|
||||
}
|
||||
|
||||
// Partition users into those with/without an existing person
|
||||
const usersWithoutPerson = users.filter(
|
||||
(user) => !emailToPersonId.has(user.email.toLowerCase()),
|
||||
);
|
||||
|
||||
// Step 1: Batch-create missing people
|
||||
const newlyCreatedPersonIds: string[] = [];
|
||||
|
||||
if (usersWithoutPerson.length > 0) {
|
||||
console.log(`Batch-creating ${usersWithoutPerson.length} people records`);
|
||||
|
||||
const createPeopleResult = await twentyClient.mutation({
|
||||
createPeople: {
|
||||
__args: {
|
||||
data: usersWithoutPerson.map((user) => ({
|
||||
name: {
|
||||
firstName: user.firstName,
|
||||
lastName: user.lastName,
|
||||
},
|
||||
emails: {
|
||||
primaryEmail: user.email,
|
||||
},
|
||||
})),
|
||||
},
|
||||
id: true,
|
||||
emails: {
|
||||
primaryEmail: true,
|
||||
},
|
||||
},
|
||||
});
|
||||
|
||||
const createdPeople = createPeopleResult.createPeople ?? [];
|
||||
|
||||
for (const person of createdPeople) {
|
||||
if (
|
||||
person.emails?.primaryEmail === undefined ||
|
||||
person.id === undefined
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
|
||||
newlyCreatedPersonIds.push(person.id);
|
||||
emailToPersonId.set(person.emails.primaryEmail.toLowerCase(), person.id);
|
||||
}
|
||||
}
|
||||
|
||||
// Step 2: Batch-upsert all cloud users, with rollback on failure
|
||||
const cloudUserInputs = users.map((user) => {
|
||||
const personId = emailToPersonId.get(user.email.toLowerCase());
|
||||
|
||||
if (personId === undefined) {
|
||||
throw new Error(
|
||||
`No personId found for user ${user.email} — this should not happen`,
|
||||
);
|
||||
}
|
||||
|
||||
return buildCloudUserInput({ user, personId });
|
||||
});
|
||||
|
||||
try {
|
||||
console.log(`Batch-upserting ${cloudUserInputs.length} cloud users`);
|
||||
|
||||
await twentyClient.mutation({
|
||||
createCloudUsers2: {
|
||||
__args: {
|
||||
data: cloudUserInputs,
|
||||
upsert: true,
|
||||
},
|
||||
__scalar: true,
|
||||
},
|
||||
});
|
||||
} catch (cloudUserError) {
|
||||
console.log(
|
||||
'Cloud user upsert failed, rolling back newly created people',
|
||||
cloudUserError,
|
||||
);
|
||||
|
||||
// Rollback: hard-delete people that were created in step 1
|
||||
if (newlyCreatedPersonIds.length > 0) {
|
||||
try {
|
||||
await twentyClient.mutation({
|
||||
destroyPeople: {
|
||||
__args: {
|
||||
filter: {
|
||||
id: {
|
||||
in: newlyCreatedPersonIds,
|
||||
},
|
||||
},
|
||||
},
|
||||
id: true,
|
||||
},
|
||||
});
|
||||
|
||||
console.log(
|
||||
`Rolled back ${newlyCreatedPersonIds.length} newly created people`,
|
||||
);
|
||||
} catch (rollbackError) {
|
||||
console.log(
|
||||
'Rollback of newly created people also failed',
|
||||
rollbackError,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
throw cloudUserError;
|
||||
}
|
||||
|
||||
return { syncedCount: users.length };
|
||||
};
|
||||
+201
@@ -0,0 +1,201 @@
|
||||
import {
|
||||
enumCloudWorkspace2ActivationStatusEnum,
|
||||
enumCloudWorkspace2PaymentFrequencyEnum,
|
||||
enumCloudWorkspace2SubscriptionStatusEnum,
|
||||
} from 'twenty-sdk/generated';
|
||||
import { z } from 'zod';
|
||||
|
||||
import { getApplicationConfig } from 'src/shared/application-config';
|
||||
import { fetchFromClickHouse } from 'src/shared/clickhouse-client';
|
||||
import { clickHouseDateToIso } from 'src/shared/clickhouse-date-to-iso';
|
||||
import { twentyClient } from 'src/shared/twenty-client';
|
||||
|
||||
const clickHouseWorkspaceSchema = z.object({
|
||||
workspaceId: z.string(),
|
||||
workspaceName: z.string(),
|
||||
subdomain: z.string(),
|
||||
customDomain: z.string(),
|
||||
createdAt: z.string(),
|
||||
updatedAt: z.string(),
|
||||
deletedAt: z.string(),
|
||||
activationStatus: z
|
||||
.string()
|
||||
.transform((val) => {
|
||||
const upper = val.toUpperCase();
|
||||
if (upper === '') return 'EMPTY';
|
||||
const valid: string[] = Object.values(
|
||||
enumCloudWorkspace2ActivationStatusEnum,
|
||||
);
|
||||
return valid.includes(upper) ? upper : 'EMPTY';
|
||||
})
|
||||
.pipe(z.enum(enumCloudWorkspace2ActivationStatusEnum)),
|
||||
createdDate: z.string(),
|
||||
lastPageviewDate: z.string(),
|
||||
pageviewsLast30d: z.coerce.number(),
|
||||
pageviewsLast7d: z.coerce.number(),
|
||||
pageviewsLast24h: z.coerce.number(),
|
||||
totalEverActiveUsers: z.coerce.number(),
|
||||
totalWorkspaceUsers: z.coerce.number(),
|
||||
activeUsersLast30d: z.coerce.number(),
|
||||
activeUsersLast7d: z.coerce.number(),
|
||||
activeUsersLast24h: z.coerce.number(),
|
||||
isActiveLast30d: z.coerce.boolean(),
|
||||
isActiveLast7d: z.coerce.boolean(),
|
||||
isActiveLast24h: z.coerce.boolean(),
|
||||
activeUserRatioLast30d: z.coerce.number(),
|
||||
activeUserRatioLast7d: z.coerce.number(),
|
||||
workspaceAgeDays: z.coerce.number(),
|
||||
totalEvents: z.coerce.number(),
|
||||
eventsLast30d: z.coerce.number(),
|
||||
eventsPerUser: z.coerce.number(),
|
||||
subscription_status: z
|
||||
.string()
|
||||
.transform((val) => {
|
||||
const upper = val.toUpperCase();
|
||||
if (upper === '') return 'EMPTY';
|
||||
const valid: string[] = Object.values(
|
||||
enumCloudWorkspace2SubscriptionStatusEnum,
|
||||
);
|
||||
return valid.includes(upper) ? upper : 'OTHER';
|
||||
})
|
||||
.pipe(z.enum(enumCloudWorkspace2SubscriptionStatusEnum)),
|
||||
payment_frequency: z
|
||||
.string()
|
||||
.transform((val) => {
|
||||
const upper = val.toUpperCase();
|
||||
if (upper === '') return 'EMPTY';
|
||||
const valid: string[] = Object.values(
|
||||
enumCloudWorkspace2PaymentFrequencyEnum,
|
||||
);
|
||||
return valid.includes(upper) ? upper : 'OTHER';
|
||||
})
|
||||
.pipe(z.enum(enumCloudWorkspace2PaymentFrequencyEnum)),
|
||||
trial_status: z.string(),
|
||||
mrr: z.coerce.number(),
|
||||
potential_arr: z.coerce.number(),
|
||||
arr: z.coerce.number().nullable(),
|
||||
next_renewal_date: z.string(),
|
||||
workspace_domain: z.string(),
|
||||
domain_source: z.string(),
|
||||
creatorUserId: z.string(),
|
||||
creatorEmail: z.string(),
|
||||
creator_domain_type: z.string(),
|
||||
primary_business_domain: z.string(),
|
||||
business_domain_user_count: z.coerce.number(),
|
||||
isTwenty: z.coerce.boolean(),
|
||||
});
|
||||
|
||||
type ClickHouseWorkspace = z.infer<typeof clickHouseWorkspaceSchema>;
|
||||
|
||||
const fetchWorkspacesFromClickHouse =
|
||||
async (): Promise<ClickHouseWorkspace[]> => {
|
||||
const { clickHouseDatabase } = getApplicationConfig();
|
||||
|
||||
// const nowDate = 'now()'
|
||||
const nowDate = "'2026-02-04 14:28:52.000'";
|
||||
|
||||
const query = `
|
||||
SELECT
|
||||
*
|
||||
FROM (
|
||||
SELECT
|
||||
*,
|
||||
row_number() OVER (PARTITION BY workspaceId ORDER BY updatedAt DESC) AS rn
|
||||
FROM
|
||||
${clickHouseDatabase}.workspace
|
||||
WHERE
|
||||
updatedAt >= ${nowDate} - INTERVAL 500 MINUTE
|
||||
AND
|
||||
updatedAt <= ${nowDate}
|
||||
)
|
||||
WHERE
|
||||
rn = 1
|
||||
FORMAT
|
||||
JSONEachRow;
|
||||
`;
|
||||
|
||||
return fetchFromClickHouse(query, clickHouseWorkspaceSchema);
|
||||
};
|
||||
|
||||
// Convert a dollar amount to micros (1 USD = 1_000_000 micros).
|
||||
const dollarsToAmountMicros = (dollars: number) =>
|
||||
Math.round(dollars * 1_000_000);
|
||||
|
||||
const buildCloudWorkspaceInput = (workspace: ClickHouseWorkspace) => ({
|
||||
id: workspace.workspaceId,
|
||||
name: workspace.workspaceName,
|
||||
subDomain: workspace.subdomain,
|
||||
customDomain: {
|
||||
primaryLinkUrl: workspace.customDomain,
|
||||
primaryLinkLabel: '',
|
||||
},
|
||||
activationStatus: workspace.activationStatus,
|
||||
subscriptionStatus: workspace.subscription_status,
|
||||
paymentFrequency: workspace.payment_frequency,
|
||||
lastPageViewDate: workspace.lastPageviewDate,
|
||||
pageViewsL30D: workspace.pageviewsLast30d,
|
||||
pageViewsL7D: workspace.pageviewsLast7d,
|
||||
pageViewsL24H: workspace.pageviewsLast24h,
|
||||
totalEverActiveWorkspaceUsers: workspace.totalEverActiveUsers,
|
||||
totalWorkspaceUsers: workspace.totalWorkspaceUsers,
|
||||
activeUsersL30D: workspace.activeUsersLast30d,
|
||||
activeUsersL7D: workspace.activeUsersLast7d,
|
||||
activeUsersL24H: workspace.activeUsersLast24h,
|
||||
isActiveL30D: workspace.isActiveLast30d,
|
||||
isActiveL7D: workspace.isActiveLast7d,
|
||||
isActiveL24H: workspace.isActiveLast24h,
|
||||
workspaceTenure: workspace.workspaceAgeDays,
|
||||
numberOfEventsTotal: workspace.totalEvents,
|
||||
numberOfEventsL30D: workspace.eventsLast30d,
|
||||
mrr: {
|
||||
amountMicros: dollarsToAmountMicros(workspace.mrr),
|
||||
currencyCode: 'USD',
|
||||
},
|
||||
potentialArr: {
|
||||
amountMicros: dollarsToAmountMicros(workspace.potential_arr),
|
||||
currencyCode: 'USD',
|
||||
},
|
||||
arr: {
|
||||
amountMicros: dollarsToAmountMicros(workspace.arr ?? 0),
|
||||
currencyCode: 'USD',
|
||||
},
|
||||
nextRenewalDate: clickHouseDateToIso(workspace.next_renewal_date),
|
||||
creatorEmail: {
|
||||
primaryEmail: workspace.creatorEmail,
|
||||
},
|
||||
workspaceBusinessDomain: {
|
||||
primaryLinkUrl: workspace.primary_business_domain,
|
||||
primaryLinkLabel: '',
|
||||
},
|
||||
dataLastUpdatedAt: new Date().toISOString(),
|
||||
});
|
||||
|
||||
export const syncCloudWorkspaces = async (): Promise<{
|
||||
syncedCount: number;
|
||||
}> => {
|
||||
const workspaces = await fetchWorkspacesFromClickHouse();
|
||||
|
||||
console.log(`Fetched ${workspaces.length} workspaces from ClickHouse`);
|
||||
|
||||
if (workspaces.length === 0) {
|
||||
return { syncedCount: 0 };
|
||||
}
|
||||
|
||||
const cloudWorkspaceInputs = workspaces.map(buildCloudWorkspaceInput);
|
||||
|
||||
console.log(
|
||||
`Batch-upserting ${cloudWorkspaceInputs.length} cloud workspaces`,
|
||||
);
|
||||
|
||||
await twentyClient.mutation({
|
||||
createCloudWorkspaces2: {
|
||||
__args: {
|
||||
data: cloudWorkspaceInputs,
|
||||
upsert: true,
|
||||
},
|
||||
__scalar: true,
|
||||
},
|
||||
});
|
||||
|
||||
return { syncedCount: workspaces.length };
|
||||
};
|
||||
@@ -0,0 +1,25 @@
|
||||
import { z } from 'zod';
|
||||
|
||||
const applicationConfigSchema = z
|
||||
.object({
|
||||
CLICKHOUSE_DATABASE: z.string().nonempty(),
|
||||
CLICKHOUSE_URL: z.url(),
|
||||
CLICKHOUSE_USERNAME: z.string().nonempty(),
|
||||
CLICKHOUSE_PASSWORD: z.string().nonempty(),
|
||||
});
|
||||
|
||||
export const getApplicationConfig = () => {
|
||||
const env = applicationConfigSchema.parse({
|
||||
CLICKHOUSE_DATABASE: process.env.CLICKHOUSE_DATABASE,
|
||||
CLICKHOUSE_URL: process.env.CLICKHOUSE_URL,
|
||||
CLICKHOUSE_USERNAME: process.env.CLICKHOUSE_USERNAME,
|
||||
CLICKHOUSE_PASSWORD: process.env.CLICKHOUSE_PASSWORD,
|
||||
});
|
||||
|
||||
return {
|
||||
clickHouseDatabase: env.CLICKHOUSE_DATABASE,
|
||||
clickHouseUrl: env.CLICKHOUSE_URL,
|
||||
clickHouseUsername: env.CLICKHOUSE_USERNAME,
|
||||
clickHousePassword: env.CLICKHOUSE_PASSWORD,
|
||||
};
|
||||
};
|
||||
@@ -0,0 +1,43 @@
|
||||
import { z } from 'zod';
|
||||
|
||||
import { getApplicationConfig } from 'src/shared/application-config';
|
||||
|
||||
// Fetches rows from ClickHouse using a raw SQL query,
|
||||
// parses the JSONEachRow response, and validates each row against the provided Zod schema.
|
||||
export const fetchFromClickHouse = async <T>(
|
||||
query: string,
|
||||
schema: z.ZodType<T>,
|
||||
): Promise<T[]> => {
|
||||
const { clickHouseUrl, clickHouseUsername, clickHousePassword } =
|
||||
getApplicationConfig();
|
||||
|
||||
const response = await fetch(clickHouseUrl, {
|
||||
method: 'POST',
|
||||
headers: {
|
||||
Authorization:
|
||||
'Basic ' +
|
||||
Buffer.from(`${clickHouseUsername}:${clickHousePassword}`).toString(
|
||||
'base64',
|
||||
),
|
||||
'Content-Type': 'text/plain',
|
||||
},
|
||||
body: query,
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const errorText = await response.text();
|
||||
|
||||
throw new Error(`ClickHouse error: ${response.status} - ${errorText}`);
|
||||
}
|
||||
|
||||
// Format is a list of JSON objects, one per line (JSONEachRow format).
|
||||
const text = await response.text();
|
||||
|
||||
const rows = text
|
||||
.trim()
|
||||
.split('\n')
|
||||
.filter((line) => line.trim())
|
||||
.map((line) => JSON.parse(line));
|
||||
|
||||
return z.array(schema).parse(rows);
|
||||
};
|
||||
@@ -0,0 +1,16 @@
|
||||
// ClickHouse returns DateTime64 values in the format 'YYYY-MM-DD HH:mm:ss.SSSSSS'
|
||||
// (space-separated, up to 6 fractional digits). Twenty's DATE_TIME fields expect
|
||||
// ISO 8601 format ('YYYY-MM-DDTHH:mm:ss.SSSZ', at most 3 fractional digits).
|
||||
//
|
||||
// Additionally, ClickHouse uses '1970-01-01 00:00:00.000000' as a sentinel for
|
||||
// null/empty dates. We map that to null.
|
||||
|
||||
const CLICKHOUSE_EPOCH_SENTINEL = '1970-01-01 00:00:00.000000';
|
||||
|
||||
export const clickHouseDateToIso = (value: string): string | null => {
|
||||
if (!value || value === CLICKHOUSE_EPOCH_SENTINEL) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return new Date(value.replace(' ', 'T') + 'Z').toISOString();
|
||||
};
|
||||
@@ -0,0 +1,3 @@
|
||||
import Twenty from 'twenty-sdk/generated';
|
||||
|
||||
export const twentyClient = new Twenty();
|
||||
Reference in New Issue
Block a user