Skip to content

Commit 36b4ed4

Browse files
authored
feat(security-contacts): resolve GitHub handles to emails using cdp identities [CM-1323] (#4342)
Signed-off-by: Mouad BANI <mouad-mb@outlook.com>
1 parent 4c5d3ee commit 36b4ed4

12 files changed

Lines changed: 170 additions & 22 deletions

File tree

scripts/services/security-contacts-worker.yaml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ x-env-args: &env-args
66
SERVICE: security-contacts-worker
77
SHELL: /bin/sh
88
SUPPRESS_NO_CONFIG_WARNING: 'true'
9-
CROWD_TEMPORAL_TASKQUEUE: packages-worker
9+
CROWD_TEMPORAL_TASKQUEUE: security-contacts-worker
1010
CROWD_TEMPORAL_NAMESPACE: ${CROWD_PACKAGES_TEMPORAL_NAMESPACE}
1111

1212
services:

services/apps/packages_worker/package.json

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -43,9 +43,9 @@
4343
"start:go-worker": "CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker tsx src/bin/go-worker.ts",
4444
"dev:go-worker": "CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9241 src/bin/go-worker.ts",
4545
"dev:go-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=go-worker SERVICE=go-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9241 src/bin/go-worker.ts",
46-
"start:security-contacts-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker SERVICE=security-contacts-worker tsx src/bin/security-contacts-worker.ts",
47-
"dev:security-contacts-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker SERVICE=security-contacts-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9243 src/bin/security-contacts-worker.ts",
48-
"dev:security-contacts-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=packages-worker SERVICE=security-contacts-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9243 src/bin/security-contacts-worker.ts",
46+
"start:security-contacts-worker": "CROWD_TEMPORAL_TASKQUEUE=security-contacts-worker SERVICE=security-contacts-worker tsx src/bin/security-contacts-worker.ts",
47+
"dev:security-contacts-worker": "CROWD_TEMPORAL_TASKQUEUE=security-contacts-worker SERVICE=security-contacts-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9243 src/bin/security-contacts-worker.ts",
48+
"dev:security-contacts-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=security-contacts-worker SERVICE=security-contacts-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9243 src/bin/security-contacts-worker.ts",
4949
"start:nuget-worker": "CROWD_TEMPORAL_TASKQUEUE=nuget-worker SERVICE=nuget-worker tsx src/bin/nuget-worker.ts",
5050
"dev:nuget-worker": "CROWD_TEMPORAL_TASKQUEUE=nuget-worker SERVICE=nuget-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9242 src/bin/nuget-worker.ts",
5151
"dev:nuget-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=nuget-worker SERVICE=nuget-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9242 src/bin/nuget-worker.ts",

services/apps/packages_worker/src/config.ts

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,16 @@ export function getPackagesDbConfig() {
1818
}
1919
}
2020

21+
export function getCdpDbConfig() {
22+
return {
23+
host: requireEnv('CROWD_DB_READ_HOST'),
24+
port: requireEnvInt('CROWD_DB_PORT'),
25+
database: requireEnv('CROWD_DB_DATABASE'),
26+
user: requireEnv('CROWD_DB_USERNAME'),
27+
password: requireEnv('CROWD_DB_PASSWORD'),
28+
}
29+
}
30+
2131
export function getGithubAppConfig() {
2232
const rawPrivateKey = requireEnv('CROWD_GITHUB_PRIVATE_KEY')
2333
const privateKeyPem = Buffer.from(rawPrivateKey, 'base64').toString('ascii')

services/apps/packages_worker/src/db.ts

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,18 @@
11
import { pgpQx } from '@crowd/data-access-layer/src/queryExecutor'
22
import { DbConnection, getDbConnection } from '@crowd/database'
33

4-
import { getPackagesDbConfig } from './config'
4+
import { getCdpDbConfig, getPackagesDbConfig } from './config'
55

66
export async function getPackagesDb() {
77
const conn = await getDbConnection(getPackagesDbConfig())
88
return pgpQx(conn)
99
}
1010

11+
export async function getCdpDb() {
12+
const conn = await getDbConnection(getCdpDbConfig())
13+
return pgpQx(conn)
14+
}
15+
1116
// Raw pg-promise connection for the same pool getPackagesDb() wraps
1217
// (getDbConnection caches per host:database). Needed for COPY FROM STDIN via
1318
// pg-copy-streams, which requires direct access to the underlying pg client.

services/apps/packages_worker/src/security-contacts/activities.ts

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
import { getServiceChildLogger } from '@crowd/logging'
22

33
import { getSecurityContactsConfig } from '../config'
4-
import { getPackagesDb } from '../db'
4+
import { getCdpDb, getPackagesDb } from '../db'
55

66
import { IngestSingleResult, ingestSecurityContactsForPurl } from './ingestSingle'
77
import { BatchResult, processBatch } from './processBatch'
@@ -11,8 +11,9 @@ const log = getServiceChildLogger('security-contacts-activity')
1111
export async function processSecurityContactsBatch(): Promise<BatchResult> {
1212
const config = getSecurityContactsConfig()
1313
const qx = await getPackagesDb()
14+
const cdpQx = await getCdpDb()
1415

15-
const result = await processBatch(qx, config)
16+
const result = await processBatch(qx, cdpQx, config)
1617
log.info({ ...result }, 'Security contacts batch activity complete')
1718
return result
1819
}
@@ -22,8 +23,9 @@ export async function ingestSecurityContactsForPurlActivity(
2223
): Promise<IngestSingleResult> {
2324
const config = getSecurityContactsConfig()
2425
const qx = await getPackagesDb()
26+
const cdpQx = await getCdpDb()
2527

26-
const result = await ingestSecurityContactsForPurl(qx, config, purl)
28+
const result = await ingestSecurityContactsForPurl(qx, cdpQx, config, purl)
2729
log.info({ purl, ...result }, 'On-demand security contacts ingest activity complete')
2830
return result
2931
}

services/apps/packages_worker/src/security-contacts/ingestSingle.ts

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,7 @@ function toTarget(row: SingleRepoRow): RepoTarget {
8080
*/
8181
export async function ingestSecurityContactsForPurl(
8282
qx: QueryExecutor,
83+
cdpQx: QueryExecutor,
8384
config: Config,
8485
purl: string,
8586
): Promise<IngestSingleResult> {
@@ -92,7 +93,7 @@ export async function ingestSecurityContactsForPurl(
9293
const target = toTarget(row)
9394
const baseDeps = buildBaseDeps(config)
9495
try {
95-
await processRepo(target, baseDeps, qx)
96+
await processRepo(target, baseDeps, qx, cdpQx)
9697
} catch (err) {
9798
log.error({ repoId: target.repoId, errMsg: (err as Error).message }, 'Repo processing failed')
9899
await markRepoAttempted(qx, target.repoId).catch(() => undefined)

services/apps/packages_worker/src/security-contacts/processBatch.ts

Lines changed: 20 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ import { extractSecurityMd } from './extractors/securityMd'
1616
import { extractSecurityTxt } from './extractors/securityTxt'
1717
import { githubApiGet } from './githubToken'
1818
import { reconcile } from './reconcile'
19+
import { resolveCdpEmails } from './resolveCdpEmails'
1920
import {
2021
Extractor,
2122
ExtractorDeps,
@@ -120,6 +121,7 @@ export async function processRepo(
120121
target: RepoTarget,
121122
baseDeps: Omit<ExtractorDeps, 'repoTree'>,
122123
qx: QueryExecutor,
124+
cdpQx: QueryExecutor,
123125
): Promise<void> {
124126
// One tree fetch per repo, shared by extractors that probe well-known paths.
125127
let repoTree: ExtractorDeps['repoTree'] = { paths: null }
@@ -161,11 +163,27 @@ export async function processRepo(
161163
contacts = contacts.filter((c) => c.channel !== 'github-pvr')
162164
}
163165

166+
const handleContacts = contacts.filter((c) => c.channel === 'github-handle')
167+
if (handleContacts.length > 0) {
168+
try {
169+
contacts.push(...(await resolveCdpEmails(cdpQx, handleContacts)))
170+
} catch (err) {
171+
log.warn(
172+
{ repoId: target.repoId, errMsg: (err as Error).message },
173+
'CDP email resolution failed — proceeding without resolved emails',
174+
)
175+
}
176+
}
177+
164178
const scored = reconcile(contacts)
165179
await writeContacts(qx, target.repoId, scored, policies)
166180
}
167181

168-
export async function processBatch(qx: QueryExecutor, config: Config): Promise<BatchResult> {
182+
export async function processBatch(
183+
qx: QueryExecutor,
184+
cdpQx: QueryExecutor,
185+
config: Config,
186+
): Promise<BatchResult> {
169187
const batch = await fetchBatch(qx)
170188
if (batch.length === 0) return { processed: 0 }
171189

@@ -191,7 +209,7 @@ export async function processBatch(qx: QueryExecutor, config: Config): Promise<B
191209
)
192210
}
193211
try {
194-
await processRepo(target, deps, qx)
212+
await processRepo(target, deps, qx, cdpQx)
195213
} catch (err) {
196214
log.error(
197215
{ repoId: target.repoId, errMsg: (err as Error).message },

services/apps/packages_worker/src/security-contacts/reconcile.ts

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,6 @@ import {
88
SourceTier,
99
} from './types'
1010

11-
const MAX_CONTACTS = 5
12-
1311
const ROLE_PRIORITY: Record<ContactRole, number> = {
1412
'security-team': 5,
1513
maintainer: 4,
@@ -143,5 +141,5 @@ export function reconcile(contacts: RawContact[], now: Date = new Date()): Score
143141
a.value.localeCompare(b.value),
144142
)
145143

146-
return scored.slice(0, MAX_CONTACTS)
144+
return scored
147145
}
Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,59 @@
1+
import {
2+
findMembersByGithubHandles,
3+
findResolvableEmailsForMembers,
4+
} from '@crowd/data-access-layer/src/members/identities'
5+
import { QueryExecutor } from '@crowd/data-access-layer/src/queryExecutor'
6+
7+
import { ProvenanceEntry, RawContact } from './types'
8+
9+
function latestTimestamp(provenance: ProvenanceEntry[]): string {
10+
const times = provenance.map((p) => p.declaredAt ?? p.fetchedAt)
11+
return times.length === 0
12+
? new Date().toISOString()
13+
: times.reduce((a, b) => (new Date(b).getTime() > new Date(a).getTime() ? b : a))
14+
}
15+
16+
export async function resolveCdpEmails(
17+
cdpQx: QueryExecutor,
18+
handleContacts: RawContact[],
19+
): Promise<RawContact[]> {
20+
if (handleContacts.length === 0) return []
21+
22+
const handles = [...new Set(handleContacts.map((c) => c.value.toLowerCase()))]
23+
const members = await findMembersByGithubHandles(cdpQx, handles)
24+
if (members.length === 0) return []
25+
26+
const memberIdsByHandle = new Map<string, string[]>()
27+
for (const m of members) {
28+
const key = m.githubHandle.toLowerCase()
29+
memberIdsByHandle.set(key, [...(memberIdsByHandle.get(key) ?? []), m.memberId])
30+
}
31+
32+
const emails = await findResolvableEmailsForMembers(cdpQx, [
33+
...new Set(members.map((m) => m.memberId)),
34+
])
35+
const emailsByMember = new Map<string, { verified: string[]; unverified: string[] }>()
36+
for (const e of emails) {
37+
const bucket = emailsByMember.get(e.memberId) ?? { verified: [], unverified: [] }
38+
;(e.verified ? bucket.verified : bucket.unverified).push(e.email)
39+
emailsByMember.set(e.memberId, bucket)
40+
}
41+
42+
return handleContacts.flatMap((contact) => {
43+
const fetchedAt = latestTimestamp(contact.provenance)
44+
const memberIds = memberIdsByHandle.get(contact.value.toLowerCase()) ?? []
45+
return memberIds.flatMap((memberId) => {
46+
const bucket = emailsByMember.get(memberId)
47+
if (!bucket) return []
48+
const useVerified = bucket.verified.length > 0
49+
const source = useVerified ? 'cdp-verified' : 'cdp-unverified'
50+
return (useVerified ? bucket.verified : bucket.unverified).map((email) => ({
51+
channel: 'email' as const,
52+
value: email,
53+
role: contact.role,
54+
tier: contact.tier,
55+
provenance: [{ source, sourceTier: contact.tier, fetchedAt }],
56+
}))
57+
})
58+
})
59+
}

services/apps/packages_worker/src/security-contacts/schedule.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ function scheduleAction() {
1313
type: 'startWorkflow' as const,
1414
workflowType: ingestSecurityContacts,
1515
workflowId: 'security-contacts-daily',
16-
taskQueue: 'packages-worker',
16+
taskQueue: 'security-contacts-worker',
1717
workflowExecutionTimeout: WORKFLOW_EXECUTION_TIMEOUT,
1818
retry: {
1919
initialInterval: '30 seconds',

0 commit comments

Comments
 (0)