Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
-- Extend osspckgs_ingest_jobs.job_kind CHECK constraint to include scorecard kinds.
-- Required for ingestScorecard workflow (CM-1227).
ALTER TABLE osspckgs_ingest_jobs
DROP CONSTRAINT osspckgs_ingest_jobs_job_kind_check,
ADD CONSTRAINT osspckgs_ingest_jobs_job_kind_check CHECK (job_kind IN (
'packages', 'versions', 'package_dependencies',
'repos', 'package_repos',
'advisories', 'advisory_packages',
'dependent_counts',
'scorecard_repos', 'scorecard_checks'
));
Original file line number Diff line number Diff line change
Expand Up @@ -3,18 +3,18 @@ version: '3.1'
x-env-args: &env-args
DOCKER_BUILDKIT: 1
NODE_ENV: docker
SERVICE: deps-dev-ingest
CROWD_TEMPORAL_TASKQUEUE: deps-dev-ingest
SERVICE: bq-dataset-ingest
CROWD_TEMPORAL_TASKQUEUE: bq-dataset-ingest
CROWD_TEMPORAL_NAMESPACE: ${CROWD_PACKAGES_TEMPORAL_NAMESPACE}
SHELL: /bin/sh
SUPPRESS_NO_CONFIG_WARNING: 'true'

services:
deps-dev-ingest:
bq-dataset-ingest:
build:
context: ../../
dockerfile: ./scripts/services/docker/Dockerfile.packages
command: 'pnpm run start:deps-dev-ingest'
command: 'pnpm run start:bq-dataset-ingest'
working_dir: /usr/crowd/app/services/apps/packages_worker
env_file:
- ../../backend/.env.dist.local
Expand All @@ -27,11 +27,11 @@ services:
networks:
- crowd-bridge

deps-dev-ingest-dev:
bq-dataset-ingest-dev:
build:
context: ../../
dockerfile: ./scripts/services/docker/Dockerfile.packages
command: 'pnpm run dev:deps-dev-ingest'
command: 'pnpm run dev:bq-dataset-ingest'
working_dir: /usr/crowd/app/services/apps/packages_worker
# user: '${USER_ID}:${GROUP_ID}'
env_file:
Expand All @@ -41,7 +41,7 @@ services:
- ../../backend/.env.override.composed
environment:
<<: *env-args
hostname: deps-dev-ingest
hostname: bq-dataset-ingest
networks:
- crowd-bridge
volumes:
Expand Down
10 changes: 5 additions & 5 deletions services/apps/packages_worker/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
"scripts": {
"start:packages-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=packages-worker tsx src/bin/packages-worker.ts",
"start:criticality-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=criticality-worker tsx src/bin/criticality-worker.ts",
"start:deps-dev-ingest": "CROWD_TEMPORAL_TASKQUEUE=deps-dev-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=deps-dev-ingest tsx src/bin/deps-dev-ingest.ts",
"start:bq-dataset-ingest": "CROWD_TEMPORAL_TASKQUEUE=bq-dataset-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=bq-dataset-ingest tsx src/bin/bq-dataset-ingest.ts",
"start:npm-worker": "CROWD_TEMPORAL_TASKQUEUE=npm-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=npm-worker tsx src/bin/npm-worker.ts",
"start:osv-worker": "CROWD_TEMPORAL_TASKQUEUE=osv-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=osv-worker tsx src/bin/osv-worker.ts",
"start:github-repos-enricher": "SERVICE=github-repos-enricher tsx src/bin/github-repos-enricher.ts",
Expand All @@ -18,21 +18,21 @@
"start:maven-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=maven-worker tsx src/bin/maven-worker.ts",
"backfill:maven:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=maven LOG_LEVEL=info tsx src/bin/maven-backfill.ts",
"dev:maven-worker": "CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=maven-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/maven-worker.ts",
"dev:deps-dev-ingest": "CROWD_TEMPORAL_TASKQUEUE=deps-dev-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=deps-dev-ingest nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/deps-dev-ingest.ts",
"dev:bq-dataset-ingest": "CROWD_TEMPORAL_TASKQUEUE=bq-dataset-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=bq-dataset-ingest nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/bq-dataset-ingest.ts",
"dev:npm-worker": "CROWD_TEMPORAL_TASKQUEUE=npm-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=npm-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/npm-worker.ts",
"dev:osv-worker": "CROWD_TEMPORAL_TASKQUEUE=osv-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=osv-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9238 src/bin/osv-worker.ts",
"dev:github-repos-enricher": "SERVICE=github-repos-enricher LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9234 src/bin/github-repos-enricher.ts",
"dev:packages-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=packages-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9233 src/bin/packages-worker.ts",
"dev:criticality-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=criticality-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9237 src/bin/criticality-worker.ts",
"dev:maven-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=packages-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=maven-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/maven-worker.ts",
"dev:deps-dev-ingest:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=deps-dev-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=deps-dev-ingest nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/deps-dev-ingest.ts",
"dev:bq-dataset-ingest:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=bq-dataset-ingest CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=bq-dataset-ingest nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/bq-dataset-ingest.ts",
"dev:npm-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=npm-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=npm-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9236 src/bin/npm-worker.ts",
"dev:osv-worker:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && CROWD_TEMPORAL_TASKQUEUE=osv-worker CROWD_TEMPORAL_NAMESPACE=$CROWD_PACKAGES_TEMPORAL_NAMESPACE SERVICE=osv-worker LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9238 src/bin/osv-worker.ts",
"start:maven": "SERVICE=maven tsx src/bin/maven.ts",
"dev:maven": "SERVICE=maven LOG_LEVEL=info nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9235 src/bin/maven.ts",
"dev:github-repos-enricher:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=github-repos-enricher LOG_LEVEL=trace nodemon --watch src --watch ../../libs --ext ts --exec tsx --inspect=0.0.0.0:9234 src/bin/github-repos-enricher.ts",
"export-to-bucket": "SERVICE=deps-dev-ingest tsx src/scripts/exportToBucket.ts",
"export-to-bucket:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=deps-dev-ingest tsx src/scripts/exportToBucket.ts",
"export-to-bucket": "SERVICE=bq-dataset-ingest tsx src/scripts/exportToBucket.ts",
"export-to-bucket:local": "set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && SERVICE=bq-dataset-ingest tsx src/scripts/exportToBucket.ts",
"monitor:osspckgs:local": "bash -c 'set -a && . ../../../backend/.env.dist.local && . ../../../backend/.env.override.local && set +a && node ../../../scripts/monitor-osspckgs.mjs'",
"lint": "npx eslint --ext .ts src --max-warnings=0",
"format": "npx prettier --write \"src/**/*.ts\"",
Expand Down
1 change: 1 addition & 0 deletions services/apps/packages_worker/src/deps-dev/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ function requireEnv(name: string): string {
export const GCP_PROJECT = requireEnv('OSSPCKGS_GCP_PROJECT')
export const GCS_BUCKET = requireEnv('OSSPCKGS_GCS_BUCKET')
export const DEPS_DEV_DATASET = 'bigquery-public-data.deps_dev_v1'
export const SCORECARD_DATASET = 'openssf.scorecardcron'

// ADR-0003: Option A = DependencyGraphEdgesLatest (prod default, has version_constraint).
// Set OSSPCKGS_DEPS_TABLE=B locally to use DependenciesLatest (cheaper, no version_constraint).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,7 @@ export async function scheduleOsspckgsBootstrap(): Promise<void> {
await temporal.schedule.create({
scheduleId: 'osspckgs-bootstrap-weekly',
spec: {
cronExpressions: ['0 2 * * 0'],
cronExpressions: ['0 2 * * 1'],
},
policies: {
overlap: ScheduleOverlapPolicy.SKIP,
Expand All @@ -20,7 +20,7 @@ export async function scheduleOsspckgsBootstrap(): Promise<void> {
action: {
type: 'startWorkflow',
workflowType: bootstrapOsspckgs,
taskQueue: 'deps-dev-ingest',
taskQueue: 'bq-dataset-ingest',
Comment thread
themarolt marked this conversation as resolved.
workflowExecutionTimeout: '12 hours',
Comment thread
themarolt marked this conversation as resolved.
retry: {
initialInterval: '1 minute',
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,14 +5,15 @@
workflowInfo,
} from '@temporalio/workflow'

import type * as depsDevActivities from '../activities'

Check failure on line 8 in services/apps/packages_worker/src/deps-dev/workflows/bootstrapOsspckgs.ts

View workflow job for this annotation

GitHub Actions / lint-format-services

Insert `{·ingestScorecard·}·from·'../../scorecard/workflows'⏎import·`

import { ingestAdvisories } from './ingestAdvisories'
import { ingestDependencies } from './ingestDependencies'
import { ingestDependentCounts } from './ingestDependentCounts'
import { ingestPackages } from './ingestPackages'
import { ingestRepos } from './ingestRepos'
import { ingestVersions } from './ingestVersions'

Check failure on line 15 in services/apps/packages_worker/src/deps-dev/workflows/bootstrapOsspckgs.ts

View workflow job for this annotation

GitHub Actions / lint-format-services

Delete `s'⏎import·{·ingestScorecard·}·from·'../../scorecard/workflow`
import { ingestScorecard } from '../../scorecard/workflows'

const { getLastSnapshot, probePartitionExists, resolveSnapshotDate } = proxyActivities<
typeof depsDevActivities
Expand Down Expand Up @@ -224,4 +225,9 @@
],
})
}
if (runs('scorecard')) {
await executeChild(ingestScorecard, {
args: [{ runId, reuseExports: opts.reuseExports, exportName: opts.exportName }],
})
}
}
2 changes: 1 addition & 1 deletion services/apps/packages_worker/src/schedules/cleanup.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ export async function scheduleOsspckgsCleanup(): Promise<void> {
action: {
type: 'startWorkflow',
workflowType: cleanupOsspckgs,
taskQueue: 'deps-dev-ingest',
taskQueue: 'bq-dataset-ingest',
workflowExecutionTimeout: '1 hour',
Comment thread
themarolt marked this conversation as resolved.
retry: {
initialInterval: '1 minute',
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import { SCORECARD_DATASET } from '../../deps-dev/config'

export const SCORECARD_REPOS_SQL = `
SELECT
CONCAT('https://', repo.name) AS repo_url,
Comment thread
themarolt marked this conversation as resolved.
Outdated
score,
date AS scanned_at
FROM \`${SCORECARD_DATASET}.scorecard-v2_latest\`
WHERE repo.name IS NOT NULL
`
Comment thread
themarolt marked this conversation as resolved.

export const SCORECARD_CHECKS_SQL = `
SELECT
CONCAT('https://', r.repo.name) AS repo_url,
Comment thread
themarolt marked this conversation as resolved.
Outdated
c.name AS check_name,
c.score AS check_score,
c.reason AS check_reason
FROM \`${SCORECARD_DATASET}.scorecard-v2_latest\` r,
UNNEST(r.checks) AS c
WHERE r.repo.name IS NOT NULL
`
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
export { ingestScorecard } from './ingestScorecard'
Original file line number Diff line number Diff line change
@@ -0,0 +1,223 @@
import { proxyActivities } from '@temporalio/workflow'

import type * as depsDevActivities from '../../deps-dev/activities'
import { SCORECARD_CHECKS_SQL, SCORECARD_REPOS_SQL } from '../queries/scorecardSql'

const { bqExportToGcs } = proxyActivities<typeof depsDevActivities>({
startToCloseTimeout: '1 hour',
retry: { maximumAttempts: 3, initialInterval: '1 minute', backoffCoefficient: 2 },
})

const { listParquetFiles } = proxyActivities<typeof depsDevActivities>({
startToCloseTimeout: '5 minutes',
retry: { maximumAttempts: 3 },
})

const { gcsParquetToStaging } = proxyActivities<typeof depsDevActivities>({
startToCloseTimeout: '2 hours',
heartbeatTimeout: '2 minutes',
retry: { maximumAttempts: 2 },
})

const { mergeStagingToTable } = proxyActivities<typeof depsDevActivities>({
startToCloseTimeout: '30 minutes',
retry: { maximumAttempts: 1 },
})

const SCORECARD_REPOS_STAGING_TABLE = 'staging.osspckgs_scorecard_repos_raw'

const SCORECARD_REPOS_STAGING_DDL = `
CREATE UNLOGGED TABLE IF NOT EXISTS staging.osspckgs_scorecard_repos_raw (
repo_url text,
score float8,
scanned_at text
)
`

// scanned_at is text because BQ DATE → Parquet INT32 DATE → JS Date; pg serialises it as ISO string.
// Cast to timestamptz happens in merge SQL.
const SCORECARD_REPOS_PG_COLUMNS = ['repo_url', 'score', 'scanned_at']

const SCORECARD_REPOS_MERGE_SQL = `
UPDATE repos r
SET scorecard_score = CASE
WHEN s.score IS NULL
OR s.score = 'NaN'::float8
OR s.score = 'Infinity'::float8
OR s.score = '-Infinity'::float8
THEN NULL
ELSE s.score::numeric(3,1)
END,
scorecard_last_run_at = s.scanned_at::timestamptz,
updated_at = NOW()
FROM staging.osspckgs_scorecard_repos_raw s
WHERE r.url = s.repo_url
`

const SCORECARD_CHECKS_STAGING_TABLE = 'staging.osspckgs_scorecard_checks_raw'

const SCORECARD_CHECKS_STAGING_DDL = `
CREATE UNLOGGED TABLE IF NOT EXISTS staging.osspckgs_scorecard_checks_raw (
repo_url text,
check_name text,
check_score int,
check_reason text
)
`

const SCORECARD_CHECKS_PG_COLUMNS = ['repo_url', 'check_name', 'check_score', 'check_reason']

const SCORECARD_CHECKS_MERGE_SQL = `
INSERT INTO repo_scorecard_checks (repo_id, check_name, score, reason)
SELECT r.id,
s.check_name,
NULLIF(s.check_score, -1)::numeric(3,1),
s.check_reason
FROM staging.osspckgs_scorecard_checks_raw s
JOIN repos r ON r.url = s.repo_url
ON CONFLICT (repo_id, check_name) DO UPDATE SET
score = EXCLUDED.score,
reason = EXCLUDED.reason,
updated_at = NOW()
`

const ROWS_PER_CHUNK = 1_000_000

export async function ingestScorecard(opts: {
runId: string
reuseExports?: boolean
exportName?: string
}): Promise<void> {
// Step 1: repos aggregate scores (plain UPDATE — repos already exist from deps-dev ingest)
const reposExport = await bqExportToGcs({
jobKind: 'scorecard_repos',
sql: SCORECARD_REPOS_SQL,
runId: opts.runId,
syncMode: 'full',
snapshotAt: null,
maxBytesGb: 10,
reuseExports: opts.reuseExports,
exportName: opts.exportName,
})

const { fileNames: repoFileNames, rowCounts: repoRowCounts } = await listParquetFiles({
gcsPrefix: reposExport.gcsPrefix,
})
const repoTotalFiles = repoFileNames.length

if (repoTotalFiles === 0) {
await mergeStagingToTable({ jobId: reposExport.jobId, mergeSql: [], tableNames: [], isFinal: true })

Check failure on line 109 in services/apps/packages_worker/src/scorecard/workflows/ingestScorecard.ts

View workflow job for this annotation

GitHub Actions / lint-format-services

Replace `·jobId:·reposExport.jobId,·mergeSql:·[],·tableNames:·[],·isFinal:·true` with `⏎······jobId:·reposExport.jobId,⏎······mergeSql:·[],⏎······tableNames:·[],⏎······isFinal:·true,⏎···`
} else {
const repoTotalRows = repoRowCounts.reduce((a, b) => a + b, 0)
const repoFilesPerChunk =
repoTotalRows > 0
? Math.max(1, Math.round((ROWS_PER_CHUNK * repoFileNames.length) / repoTotalRows))
: Math.min(repoFileNames.length, 2)
const repoTotalChunks = Math.ceil(repoFileNames.length / repoFilesPerChunk)
let repoPriorRowsAffected = 0
let repoPriorStagingRows = 0
const repoPriorTableRowCounts: Record<string, number> = {}

for (let chunkIndex = 0; chunkIndex < repoTotalChunks; chunkIndex++) {
const start = chunkIndex * repoFilesPerChunk
const chunk = repoFileNames.slice(start, start + repoFilesPerChunk)
const isFinal = chunkIndex === repoTotalChunks - 1

const { rowsLoaded } = await gcsParquetToStaging({
jobId: reposExport.jobId,
stagingTable: SCORECARD_REPOS_STAGING_TABLE,
stagingDdl: SCORECARD_REPOS_STAGING_DDL,
pgColumns: SCORECARD_REPOS_PG_COLUMNS,
fileNames: chunk,
filesOffset: start,
totalFiles: repoTotalFiles,
priorStagingRows: repoPriorStagingRows,
})
repoPriorStagingRows += rowsLoaded

const { rowsAffected, tableRowCounts } = await mergeStagingToTable({
jobId: reposExport.jobId,
mergeSql: SCORECARD_REPOS_MERGE_SQL,
tableNames: 'repos',
isFinal,
priorRowsAffected: repoPriorRowsAffected,
priorTableRowCounts: repoPriorTableRowCounts,
chunkInfo: { index: chunkIndex, total: repoTotalChunks },
})

repoPriorRowsAffected += rowsAffected
if (!isFinal) {
for (const [k, v] of Object.entries(tableRowCounts)) {
repoPriorTableRowCounts[k] = (repoPriorTableRowCounts[k] ?? 0) + v
}
}
}
}

// Step 2: per-check detail (FK → repos, so must run after Step 1)
const checksExport = await bqExportToGcs({
jobKind: 'scorecard_checks',
sql: SCORECARD_CHECKS_SQL,
runId: opts.runId,
syncMode: 'full',
snapshotAt: null,
maxBytesGb: 200,
reuseExports: opts.reuseExports,
exportName: opts.exportName,
})

const { fileNames: checkFileNames, rowCounts: checkRowCounts } = await listParquetFiles({
gcsPrefix: checksExport.gcsPrefix,
})
const checkTotalFiles = checkFileNames.length

if (checkTotalFiles === 0) {
await mergeStagingToTable({ jobId: checksExport.jobId, mergeSql: [], tableNames: [], isFinal: true })

Check failure on line 175 in services/apps/packages_worker/src/scorecard/workflows/ingestScorecard.ts

View workflow job for this annotation

GitHub Actions / lint-format-services

Replace `·jobId:·checksExport.jobId,·mergeSql:·[],·tableNames:·[],·isFinal:·true` with `⏎······jobId:·checksExport.jobId,⏎······mergeSql:·[],⏎······tableNames:·[],⏎······isFinal:·true,⏎···`
return
}

const checkTotalRows = checkRowCounts.reduce((a, b) => a + b, 0)
const checkFilesPerChunk =
checkTotalRows > 0
? Math.max(1, Math.round((ROWS_PER_CHUNK * checkFileNames.length) / checkTotalRows))
: Math.min(checkFileNames.length, 2)
const checkTotalChunks = Math.ceil(checkFileNames.length / checkFilesPerChunk)
let checkPriorRowsAffected = 0
let checkPriorStagingRows = 0
const checkPriorTableRowCounts: Record<string, number> = {}

for (let chunkIndex = 0; chunkIndex < checkTotalChunks; chunkIndex++) {
const start = chunkIndex * checkFilesPerChunk
const chunk = checkFileNames.slice(start, start + checkFilesPerChunk)
const isFinal = chunkIndex === checkTotalChunks - 1

const { rowsLoaded } = await gcsParquetToStaging({
jobId: checksExport.jobId,
stagingTable: SCORECARD_CHECKS_STAGING_TABLE,
stagingDdl: SCORECARD_CHECKS_STAGING_DDL,
pgColumns: SCORECARD_CHECKS_PG_COLUMNS,
fileNames: chunk,
filesOffset: start,
totalFiles: checkTotalFiles,
priorStagingRows: checkPriorStagingRows,
})
checkPriorStagingRows += rowsLoaded

const { rowsAffected, tableRowCounts } = await mergeStagingToTable({
jobId: checksExport.jobId,
mergeSql: SCORECARD_CHECKS_MERGE_SQL,
tableNames: 'repo_scorecard_checks',
isFinal,
priorRowsAffected: checkPriorRowsAffected,
priorTableRowCounts: checkPriorTableRowCounts,
chunkInfo: { index: chunkIndex, total: checkTotalChunks },
})

checkPriorRowsAffected += rowsAffected
if (!isFinal) {
for (const [k, v] of Object.entries(tableRowCounts)) {
checkPriorTableRowCounts[k] = (checkPriorTableRowCounts[k] ?? 0) + v
}
}
}
}
Loading
Loading