Skip to content

Commit 2228dc6

Browse files
authored
Merge branch 'main' into chore/update-pvtr-version
2 parents b884d89 + 6fb17af commit 2228dc6

6 files changed

Lines changed: 143 additions & 16 deletions

File tree

services/apps/packages_worker/src/activities.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@ export { processMavenCriticalBatch } from './maven/activities'
1616
export { criticalityComputePageRank, rankPackages } from './criticality/activities'
1717
export {
1818
cargoDownloadAndLoad,
19+
cargoNormalizeRepos,
1920
cargoEnrichPackages,
2021
cargoEnrichVersions,
2122
cargoEnrichRepos,

services/apps/packages_worker/src/cargo/activities.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,13 +15,15 @@ import {
1515
flushAudit,
1616
} from './enrich'
1717
import { STAGING_SCHEMA, loadDump } from './loadDump'
18+
import { normalizeRepos } from './normalizeRepos'
1819
import {
1920
EnrichDownloadsDailyResult,
2021
EnrichMaintainersResult,
2122
EnrichPackagesResult,
2223
EnrichReposResult,
2324
EnrichVersionsResult,
2425
LoadResult,
26+
NormalizeReposResult,
2527
} from './types'
2628

2729
const log = getServiceChildLogger('cargo-activity')
@@ -37,6 +39,10 @@ export async function cargoDownloadAndLoad(): Promise<LoadResult> {
3739
return result
3840
}
3941

42+
export async function cargoNormalizeRepos(): Promise<NormalizeReposResult> {
43+
return normalizeRepos(await getPackagesDb())
44+
}
45+
4046
export async function cargoEnrichPackages(): Promise<EnrichPackagesResult> {
4147
return enrichPackages(await getPackagesDb())
4248
}

services/apps/packages_worker/src/cargo/enrich.ts

Lines changed: 54 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ export async function enrichPackages(qx: QueryExecutor): Promise<EnrichPackagesR
4848
tx.selectOne(
4949
`WITH snap AS (
5050
SELECT p.id, p.status, p.description, p.homepage, p.declared_repository_url,
51+
p.repository_url,
5152
p.licenses, p.licenses_raw, p.keywords, p.versions_count, p.latest_version,
5253
p.first_release_at, p.latest_release_at, p.dependent_count,
5354
p.dependent_repos_count, p.downloads_last_30d
@@ -60,6 +61,8 @@ export async function enrichPackages(qx: QueryExecutor): Promise<EnrichPackagesR
6061
description = COALESCE(e.description, p.description),
6162
homepage = COALESCE(e.homepage, p.homepage),
6263
declared_repository_url = COALESCE(e.declared_repository_url, p.declared_repository_url),
64+
repository_url = CASE WHEN e.declared_repository_url IS NOT NULL
65+
THEN rn.repository_url ELSE p.repository_url END,
6366
licenses = COALESCE(e.licenses, p.licenses),
6467
licenses_raw = COALESCE(e.licenses_raw, p.licenses_raw),
6568
keywords = COALESCE(e.keywords, p.keywords),
@@ -73,18 +76,22 @@ export async function enrichPackages(qx: QueryExecutor): Promise<EnrichPackagesR
7376
ingestion_source = $(ingestionSource),
7477
last_synced_at = NOW()
7578
FROM ${STAGING_SCHEMA}.enrich_packages e
79+
LEFT JOIN ${STAGING_SCHEMA}.repo_norm rn ON rn.declared = e.declared_repository_url
7680
WHERE p.id = e.package_id
7781
RETURNING p.id
7882
),
7983
diff AS (
8084
SELECT s.id AS package_id, f.field
8185
FROM snap s
8286
JOIN ${STAGING_SCHEMA}.enrich_packages e ON e.package_id = s.id
87+
LEFT JOIN ${STAGING_SCHEMA}.repo_norm rn ON rn.declared = e.declared_repository_url
8388
CROSS JOIN LATERAL (VALUES
8489
('packages.status', s.status IS DISTINCT FROM e.status),
8590
('packages.description', s.description IS DISTINCT FROM COALESCE(e.description, s.description)),
8691
('packages.homepage', s.homepage IS DISTINCT FROM COALESCE(e.homepage, s.homepage)),
8792
('packages.declared_repository_url', s.declared_repository_url IS DISTINCT FROM COALESCE(e.declared_repository_url, s.declared_repository_url)),
93+
('packages.repository_url', s.repository_url IS DISTINCT FROM
94+
CASE WHEN e.declared_repository_url IS NOT NULL THEN rn.repository_url ELSE s.repository_url END),
8895
('packages.licenses', s.licenses IS DISTINCT FROM COALESCE(e.licenses, s.licenses)),
8996
('packages.licenses_raw', s.licenses_raw IS DISTINCT FROM COALESCE(e.licenses_raw, s.licenses_raw)),
9097
('packages.keywords', s.keywords IS DISTINCT FROM COALESCE(e.keywords, s.keywords)),
@@ -159,36 +166,65 @@ export async function enrichVersions(qx: QueryExecutor): Promise<EnrichVersionsR
159166
return { upserted: row.upserted }
160167
}
161168

162-
// Writes only url + host — other repo fields belong to the GitHub enricher.
169+
// Writes only url + host — other repo fields belong to the GitHub enricher. Uses
170+
// repo_norm (built by normalizeRepos) so repos.url/package_repos always agree with
171+
// the canonical packages.repository_url — never the raw declared_repository_url.
163172
export async function enrichRepos(qx: QueryExecutor): Promise<EnrichReposResult> {
164173
return withTunedSession(qx, 'repos', async (tx) => {
165174
const repoRow = await tx.selectOne(
166175
`WITH new_repos AS (
167176
INSERT INTO repos (url, host, updated_at)
168-
SELECT DISTINCT e.declared_repository_url,
169-
CASE
170-
WHEN e.declared_repository_url ~* '://([^/]+\\.)?github\\.com(/|$)' THEN 'github'
171-
WHEN e.declared_repository_url ~* '://[^/]*gitlab' THEN 'gitlab'
172-
WHEN e.declared_repository_url ~* '://([^/]+\\.)?bitbucket\\.org(/|$)' THEN 'bitbucket'
173-
ELSE 'other'
174-
END,
175-
NOW()
177+
SELECT DISTINCT rn.repository_url, rn.host, NOW()
176178
FROM ${STAGING_SCHEMA}.enrich_packages e
177-
WHERE e.declared_repository_url IS NOT NULL AND e.declared_repository_url LIKE 'http%'
179+
JOIN ${STAGING_SCHEMA}.repo_norm rn ON rn.declared = e.declared_repository_url
178180
ON CONFLICT (url) DO NOTHING
179181
RETURNING url
180182
),
181183
ins_audit AS (
182184
INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field)
183185
SELECT e.package_id, f.field
184186
FROM ${STAGING_SCHEMA}.enrich_packages e
185-
JOIN new_repos nr ON nr.url = e.declared_repository_url
187+
JOIN ${STAGING_SCHEMA}.repo_norm rn ON rn.declared = e.declared_repository_url
188+
JOIN new_repos nr ON nr.url = rn.repository_url
186189
CROSS JOIN LATERAL (VALUES ('repos.url'), ('repos.host')) AS f(field)
187190
RETURNING 1
188191
)
189192
SELECT (SELECT COUNT(*) FROM new_repos)::int AS repos`,
190193
)
191194

195+
// Prunes stale 'declared' links before relinking — covers junk/unparseable declared
196+
// values, rewrites (declared URL now maps elsewhere), and removals (this dump's
197+
// declared_repository_url is NULL, meaning the crate no longer declares a repo at
198+
// all — loadDump.ts stages every matched crate every run, so NULL here is
199+
// authoritative, not "no data this run"). Safe to always prune on that signal because
200+
// the DELETE is scoped to source = 'declared' — cargo only ever removes links it owns.
201+
// Without this, package_repos would accumulate a link to a repo no crate declares
202+
// anymore, and consumers such as security-contacts (which join through
203+
// repos ⋈ package_repos, not packages.repository_url) would keep reading it.
204+
const pruneRow = await tx.selectOne(
205+
`WITH targets AS (
206+
SELECT e.package_id, r.id AS repo_id
207+
FROM ${STAGING_SCHEMA}.enrich_packages e
208+
LEFT JOIN ${STAGING_SCHEMA}.repo_norm rn ON rn.declared = e.declared_repository_url
209+
LEFT JOIN repos r ON r.url = rn.repository_url
210+
),
211+
del AS (
212+
DELETE FROM package_repos pr
213+
USING targets t
214+
WHERE pr.package_id = t.package_id
215+
AND pr.source = $(source)
216+
AND (t.repo_id IS NULL OR pr.repo_id IS DISTINCT FROM t.repo_id)
217+
RETURNING pr.package_id
218+
),
219+
ins_audit AS (
220+
INSERT INTO ${STAGING_SCHEMA}.audit_changes (package_id, field)
221+
SELECT package_id, 'package_repos.repo_id' FROM del
222+
RETURNING 1
223+
)
224+
SELECT (SELECT COUNT(*) FROM del)::int AS pruned`,
225+
{ source: REPO_LINK_SOURCE },
226+
)
227+
192228
const linkRow = await tx.selectOne(
193229
`WITH old AS (
194230
SELECT pr.package_id, pr.repo_id, pr.source, pr.confidence
@@ -201,11 +237,13 @@ export async function enrichRepos(qx: QueryExecutor): Promise<EnrichReposResult>
201237
INSERT INTO package_repos (package_id, repo_id, source, confidence, created_at, verified_at)
202238
SELECT e.package_id, r.id, $(source), $(confidence), NOW(), NOW()
203239
FROM ${STAGING_SCHEMA}.enrich_packages e
204-
JOIN repos r ON r.url = e.declared_repository_url
205-
WHERE e.declared_repository_url IS NOT NULL
240+
JOIN ${STAGING_SCHEMA}.repo_norm rn ON rn.declared = e.declared_repository_url
241+
JOIN repos r ON r.url = rn.repository_url
242+
-- Leaves source untouched on conflict — a link another enricher already owns for
243+
-- this (package_id, repo_id) keeps its provenance instead of being reassigned to
244+
-- 'declared', matching upsertMavenPackageRepo's confidence-only merge.
206245
ON CONFLICT (package_id, repo_id) DO UPDATE SET
207-
source = EXCLUDED.source,
208-
confidence = EXCLUDED.confidence,
246+
confidence = GREATEST(EXCLUDED.confidence, package_repos.confidence),
209247
verified_at = NOW()
210248
RETURNING package_id, repo_id, source, confidence
211249
),
@@ -228,7 +266,7 @@ export async function enrichRepos(qx: QueryExecutor): Promise<EnrichReposResult>
228266
{ source: REPO_LINK_SOURCE, confidence: REPO_LINK_CONFIDENCE },
229267
)
230268

231-
return { repos: repoRow.repos, links: linkRow.links }
269+
return { repos: repoRow.repos, links: linkRow.links, pruned: pruneRow.pruned }
232270
})
233271
}
234272

Lines changed: 72 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,72 @@
1+
import { QueryExecutor } from '@crowd/data-access-layer'
2+
import { getServiceChildLogger } from '@crowd/logging'
3+
4+
import { CanonicalRepo, canonicalizeRepoUrl } from '../utils/canonicalizeRepoUrl'
5+
6+
import { STAGING_SCHEMA } from './loadDump'
7+
import { NormalizeReposResult } from './types'
8+
9+
const log = getServiceChildLogger('cargo-normalize')
10+
11+
// Batch size for the mapping upsert — keeps parameter arrays well under Postgres limits.
12+
const INSERT_BATCH = 5000
13+
14+
/**
15+
* Builds `cargo_sync.repo_norm(declared, repository_url, host)`, mapping every
16+
* distinct `declared_repository_url` staged in `enrich_packages` to its canonical
17+
* `https://<host>/<owner>/<name>` form via the shared `canonicalizeRepoUrl`
18+
* (same normalizer npm/pypi use — no per-ecosystem fork).
19+
*
20+
* Only inputs `canonicalizeRepoUrl` can reduce to an owner/name pair are stored;
21+
* unparseable URLs or those with fewer than two path segments are omitted, so a
22+
* LEFT JOIN from `enrich_packages` yields NULL — that both recovers derivable
23+
* URLs (Gap A) and clears that class of junk from `packages.repository_url`
24+
* (Gap C). Note this does not validate that the result is an actual repository
25+
* host — a URL-shaped homepage with ≥2 path segments (e.g.
26+
* `https://example.com/owner/repo`) still canonicalizes and is stored.
27+
*
28+
* Normalization runs in TypeScript because the bulk set-based cargo pipeline
29+
* cannot call the parser per row; the resulting mapping table lets the SQL
30+
* enrich phases join back to it. Idempotent: the table is rebuilt each run.
31+
*/
32+
export async function normalizeRepos(qx: QueryExecutor): Promise<NormalizeReposResult> {
33+
const rows: Array<{ declared_repository_url: string }> = await qx.select(
34+
`SELECT DISTINCT declared_repository_url
35+
FROM ${STAGING_SCHEMA}.enrich_packages
36+
WHERE declared_repository_url IS NOT NULL`,
37+
)
38+
39+
await qx.result(
40+
`DROP TABLE IF EXISTS ${STAGING_SCHEMA}.repo_norm CASCADE;
41+
CREATE TABLE ${STAGING_SCHEMA}.repo_norm (
42+
declared text PRIMARY KEY,
43+
repository_url text NOT NULL,
44+
host text NOT NULL
45+
)`,
46+
)
47+
48+
const mapped = rows
49+
.map(({ declared_repository_url }) => ({
50+
declared: declared_repository_url,
51+
canonical: canonicalizeRepoUrl(declared_repository_url),
52+
}))
53+
.filter((r): r is { declared: string; canonical: CanonicalRepo } => r.canonical !== null)
54+
55+
for (let i = 0; i < mapped.length; i += INSERT_BATCH) {
56+
const batch = mapped.slice(i, i + INSERT_BATCH)
57+
await qx.result(
58+
`INSERT INTO ${STAGING_SCHEMA}.repo_norm (declared, repository_url, host)
59+
SELECT * FROM unnest($(declared)::text[], $(urls)::text[], $(hosts)::text[])
60+
ON CONFLICT (declared) DO NOTHING`,
61+
{
62+
declared: batch.map((r) => r.declared),
63+
urls: batch.map((r) => r.canonical.url),
64+
hosts: batch.map((r) => r.canonical.host),
65+
},
66+
)
67+
}
68+
69+
const result: NormalizeReposResult = { scanned: rows.length, normalized: mapped.length }
70+
log.info(result, 'cargo repo normalization complete')
71+
return result
72+
}

services/apps/packages_worker/src/cargo/types.ts

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,11 @@ export interface LoadResult {
1212
durationMs: number
1313
}
1414

15+
export interface NormalizeReposResult {
16+
scanned: number
17+
normalized: number
18+
}
19+
1520
export interface EnrichPackagesResult {
1621
updated: number
1722
}
@@ -23,6 +28,7 @@ export interface EnrichVersionsResult {
2328
export interface EnrichReposResult {
2429
repos: number
2530
links: number
31+
pruned: number
2632
}
2733

2834
export interface EnrichMaintainersResult {

services/apps/packages_worker/src/cargo/workflows.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ const { cargoDownloadAndLoad } = proxyActivities<typeof activities>({
1414
})
1515

1616
const {
17+
cargoNormalizeRepos,
1718
cargoEnrichPackages,
1819
cargoEnrichVersions,
1920
cargoEnrichRepos,
@@ -31,6 +32,8 @@ export async function cargoSyncWorkflow(): Promise<void> {
3132
const load = await cargoDownloadAndLoad()
3233
log.info('cargoSync loaded dump', { ...load })
3334

35+
const normalizedRepos = await cargoNormalizeRepos()
36+
log.info('cargoSync normalized repos', { ...normalizedRepos })
3437
const packages = await cargoEnrichPackages()
3538
const versions = await cargoEnrichVersions()
3639
const repos = await cargoEnrichRepos()
@@ -41,6 +44,7 @@ export async function cargoSyncWorkflow(): Promise<void> {
4144

4245
log.info('cargoSync complete', {
4346
matched: load.matched,
47+
normalizedRepos,
4448
packages,
4549
versions,
4650
repos,

0 commit comments

Comments
 (0)