Skip to content

Commit 1b9e96a

Browse files
committed
feat: add mailing list onboarding endpoint (CM-1318)
Signed-off-by: Uroš Marolt <uros@marolt.me>
1 parent 04136f0 commit 1b9e96a

5 files changed

Lines changed: 130 additions & 0 deletions

File tree

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,14 @@
1+
import Permissions from '../../../security/permissions'
2+
import IntegrationService from '../../../services/integrationService'
3+
import PermissionChecker from '../../../services/user/permissionChecker'
4+
5+
export default async (req, res) => {
6+
new PermissionChecker(req).validateHas(Permissions.values.tenantEdit)
7+
const integrationData = {
8+
...req.body,
9+
lists: req.body.lists || [],
10+
}
11+
12+
const payload = await new IntegrationService(req).mailingListConnectOrUpdate(integrationData)
13+
await req.responseHandler.success(req, res, payload)
14+
}

backend/src/api/integration/index.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -77,6 +77,8 @@ export default (app) => {
7777

7878
// Git
7979
app.put(`/git-connect`, safeWrap(require('./helpers/gitAuthenticate').default))
80+
81+
app.put(`/mailing-list-connect`, safeWrap(require('./helpers/mailingListAuthenticate').default))
8082
app.put(`/confluence-connect`, safeWrap(require('./helpers/confluenceAuthenticate').default))
8183
app.put(`/gerrit-connect`, safeWrap(require('./helpers/gerritAuthenticate').default))
8284
app.get('/devto-validate', safeWrap(require('./helpers/devtoValidators').default))

backend/src/services/integrationService.ts

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import {
1616
} from '@crowd/common'
1717
import { CommonIntegrationService, getGithubInstallationToken } from '@crowd/common_services'
1818
import { ICreateInsightsProject } from '@crowd/data-access-layer/src/collections'
19+
import { upsertMailingLists } from '@crowd/data-access-layer/src/mailinglist'
1920
import {
2021
ICreateRepository,
2122
IRepository,
@@ -1423,6 +1424,61 @@ export default class IntegrationService {
14231424
return integration
14241425
}
14251426

1427+
/**
1428+
* Adds/updates a mailing list (public-inbox/lore) integration and onboards
1429+
* its lists for processing by the mailing_list_integration worker.
1430+
*
1431+
* @param integrationData.lists - Mailing lists to onboard (name + sourceUrl)
1432+
* @param options - Optional repository options
1433+
*/
1434+
async mailingListConnectOrUpdate(
1435+
integrationData: {
1436+
lists: Array<{ name: string; sourceUrl: string }>
1437+
},
1438+
options?: IRepositoryOptions,
1439+
) {
1440+
const lists = integrationData.lists || []
1441+
1442+
if (lists.length === 0) {
1443+
this.options.log.warn('No lists provided - skipping mailing list integration update')
1444+
return null
1445+
}
1446+
1447+
const currentOptions = options || this.options
1448+
const existingTransaction =
1449+
currentOptions.transaction || SequelizeRepository.getTransaction(currentOptions)
1450+
const transaction =
1451+
existingTransaction || (await SequelizeRepository.createTransaction(options || this.options))
1452+
let integration
1453+
1454+
try {
1455+
integration = await this.createOrUpdate(
1456+
{
1457+
platform: PlatformType.MAILINGLIST,
1458+
settings: { lists },
1459+
status: 'done',
1460+
},
1461+
transaction,
1462+
options,
1463+
)
1464+
1465+
const currentSegmentId = (options || this.options).currentSegments[0].id
1466+
const qx = SequelizeRepository.getQueryExecutor({ ...(options || this.options), transaction })
1467+
await upsertMailingLists(qx, currentSegmentId, integration.id, lists)
1468+
1469+
if (!existingTransaction) {
1470+
await SequelizeRepository.commitTransaction(transaction)
1471+
}
1472+
} catch (err) {
1473+
if (!existingTransaction) {
1474+
await SequelizeRepository.rollbackTransaction(transaction)
1475+
}
1476+
this.options.log.error(`mailingListConnectOrUpdate failed with error: ${err}`)
1477+
throw err
1478+
}
1479+
return integration
1480+
}
1481+
14261482
async atlassianAdminConnect(adminApi: string, organizationId: string) {
14271483
const nangoPayload = {
14281484
params: {
Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,57 @@
1+
import { QueryExecutor } from '../queryExecutor'
2+
3+
export interface IMailingListToOnboard {
4+
name: string
5+
sourceUrl: string
6+
}
7+
8+
/**
9+
* Upsert mailing lists (public-inbox/lore) for a segment/integration and
10+
* seed their processing state so the mailing_list_integration worker picks
11+
* them up. Re-running with the same sourceUrl re-points the list at the
12+
* given segment/integration without resetting its processing progress.
13+
* @param qx - Query executor (should be transactional)
14+
* @param segmentId - Segment the lists belong to
15+
* @param integrationId - Integration these lists are onboarded under
16+
* @param lists - Mailing lists to onboard, keyed by sourceUrl
17+
*/
18+
export async function upsertMailingLists(
19+
qx: QueryExecutor,
20+
segmentId: string,
21+
integrationId: string,
22+
lists: IMailingListToOnboard[],
23+
): Promise<string[]> {
24+
if (lists.length === 0) {
25+
return []
26+
}
27+
28+
const rows = await qx.select(
29+
`
30+
WITH upserted_list AS (
31+
INSERT INTO mailinglist.lists (id, name, "sourceUrl", "segmentId", "integrationId", "createdAt", "updatedAt")
32+
SELECT uuid_generate_v4(), v.name, v."sourceUrl", $(segmentId)::uuid, $(integrationId)::uuid, NOW(), NOW()
33+
FROM json_to_recordset($(lists)::json) AS v(name text, "sourceUrl" text)
34+
ON CONFLICT ("sourceUrl") DO UPDATE SET
35+
name = EXCLUDED.name,
36+
"segmentId" = EXCLUDED."segmentId",
37+
"integrationId" = EXCLUDED."integrationId",
38+
"updatedAt" = NOW()
39+
RETURNING id
40+
),
41+
seed_processing AS (
42+
INSERT INTO mailinglist."listProcessing" ("listId", state, priority, "lastProcessedHeads", "createdAt", "updatedAt")
43+
SELECT id, 'pending', 2, '{}'::jsonb, NOW(), NOW()
44+
FROM upserted_list
45+
ON CONFLICT ("listId") DO NOTHING
46+
)
47+
SELECT id FROM upserted_list
48+
`,
49+
{
50+
segmentId,
51+
integrationId,
52+
lists: JSON.stringify(lists),
53+
},
54+
)
55+
56+
return rows.map((row: { id: string }) => row.id)
57+
}

services/libs/types/src/enums/platforms.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -20,6 +20,7 @@ export enum PlatformType {
2020
GIT = 'git',
2121
CRUNCHBASE = 'crunchbase',
2222
GROUPSIO = 'groupsio',
23+
MAILINGLIST = 'mailinglist',
2324
COMMITTEES = 'committees',
2425
CONFLUENCE = 'confluence',
2526
GERRIT = 'gerrit',

0 commit comments

Comments
 (0)