feat(scoring): nightly grant-embedding workflow (Stage 2, step 1)
embedGrants (cron 04:15, after the ingest crons): open grants with a
synopsis and no vector → buildGrantEmbeddingText (pure, tested:
title+funder+program areas+synopsis, 8K cap) → gemini-embedding-001 @
1536 dims RETRIEVAL_DOCUMENT (org profiles will embed as
RETRIEVAL_QUERY on the other side) → grants.synopsis_embedding.
Chunked embed→store (100/chunk) so failures resume from the last
stored chunk; 500/run spend cap.
Core: serverListGrantsNeedingEmbedding (open+unembedded, closest
deadline first), serverSetGrantEmbeddings. Worker gains
@novelpad/outreach-ai dep; wired into main.ts and run-once (incl. the
missed run-once deps injection).
Live-verified: 199/199 open grants embedded in 16s; semantic probe
('after-school STEM education for youth') ranks NCI Youth Enjoy
Science R25 first at 0.639 cosine via the hnsw index.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -17,6 +17,7 @@
|
||||
"dependencies": {
|
||||
"@dbos-inc/dbos-sdk": "4.17.6",
|
||||
"@dbos-inc/drizzle-datasource": "4.17.6",
|
||||
"@novelpad/outreach-ai": "workspace:^",
|
||||
"@novelpad/outreach-core": "workspace:^",
|
||||
"drizzle-orm": "0.44.6",
|
||||
"fast-xml-parser": "^4.5.0",
|
||||
|
||||
@@ -30,6 +30,7 @@ import pg from 'pg';
|
||||
// scheduled-function args, so the `db` handle can't be passed through the
|
||||
// scheduler; it's threaded in via this module-scope registry instead (same
|
||||
// pattern as novelpad-desktop's `setStartDeps`).
|
||||
import { setEmbedGrantsDeps } from './workflows/embed-grants.js';
|
||||
import { setEnrichOrgsDeps } from './workflows/enrich-orgs.js';
|
||||
import { setExpireGrantsDeps } from './workflows/expire-grants.js';
|
||||
import { setIngestGrantsDeps } from './workflows/ingest-grants.js';
|
||||
@@ -65,6 +66,7 @@ async function main() {
|
||||
setIngestPndRssDeps({ db });
|
||||
setIngestNhdojOrgsDeps({ db });
|
||||
setEnrichOrgsDeps({ db });
|
||||
setEmbedGrantsDeps({ db });
|
||||
|
||||
DBOS.setConfig({
|
||||
name: 'helmdocs-outreach-worker',
|
||||
|
||||
@@ -19,6 +19,10 @@ import { schema } from '@novelpad/outreach-core';
|
||||
import { drizzle } from 'drizzle-orm/node-postgres';
|
||||
import pg from 'pg';
|
||||
|
||||
import {
|
||||
runEmbedGrantsNow,
|
||||
setEmbedGrantsDeps,
|
||||
} from './workflows/embed-grants.js';
|
||||
import { runEnrichOrgsNow, setEnrichOrgsDeps } from './workflows/enrich-orgs.js';
|
||||
import {
|
||||
runExpireGrantsNow,
|
||||
@@ -43,6 +47,7 @@ const RUNNERS: Record<string, () => Promise<void>> = {
|
||||
ingestNhdojOrgs: runIngestNhdojOrgsNow,
|
||||
enrichOrgs: runEnrichOrgsNow,
|
||||
expireGrants: runExpireGrantsNow,
|
||||
embedGrants: runEmbedGrantsNow,
|
||||
};
|
||||
|
||||
const FIRST_RUN_ORDER = [
|
||||
@@ -51,6 +56,7 @@ const FIRST_RUN_ORDER = [
|
||||
'ingestNhdojOrgs',
|
||||
'expireGrants',
|
||||
'enrichOrgs',
|
||||
'embedGrants',
|
||||
];
|
||||
|
||||
if (process.env.DATABASE_URL == null) {
|
||||
@@ -85,6 +91,7 @@ async function main() {
|
||||
setIngestPndRssDeps({ db });
|
||||
setIngestNhdojOrgsDeps({ db });
|
||||
setEnrichOrgsDeps({ db });
|
||||
setEmbedGrantsDeps({ db });
|
||||
|
||||
DBOS.setConfig({
|
||||
name: 'helmdocs-outreach-worker',
|
||||
|
||||
170
apps/outreach-worker/src/workflows/embed-grants.ts
Normal file
170
apps/outreach-worker/src/workflows/embed-grants.ts
Normal file
@@ -0,0 +1,170 @@
|
||||
/**
|
||||
* Nightly grant-embedding workflow — Stage 2, step 1.
|
||||
*
|
||||
* Embeds open grants' synopses with gemini-embedding-001 (1536 dims,
|
||||
* RETRIEVAL_DOCUMENT task type — org profiles embed as RETRIEVAL_QUERY on
|
||||
* the other side of the mission-fit comparison) and stores the vectors in
|
||||
* `grants.synopsis_embedding`, where the hnsw cosine index serves the
|
||||
* scoring engine's mission-fit subscore.
|
||||
*
|
||||
* Runs after the ingest crons (04:15 vs 03:00/03:30) so fresh grants embed
|
||||
* the same night they land. This is the pipeline's first paid AI call —
|
||||
* order of magnitude: ~500 tokens/grant, so a 200-grant batch is a few
|
||||
* cents of Vertex spend.
|
||||
*
|
||||
* Registration follows `ingest-grants.ts` exactly (dual workflow+scheduled
|
||||
* registration, module-scope deps registry, globalThis guard).
|
||||
*/
|
||||
import { DBOS, SchedulerMode } from '@dbos-inc/dbos-sdk';
|
||||
import {
|
||||
buildGrantEmbeddingText,
|
||||
generateDocumentEmbeddings,
|
||||
} from '@novelpad/outreach-ai';
|
||||
import type { schema } from '@novelpad/outreach-core';
|
||||
import {
|
||||
serverListGrantsNeedingEmbedding,
|
||||
serverSetGrantEmbeddings,
|
||||
type GrantNeedingEmbedding,
|
||||
} from '@novelpad/outreach-core/server';
|
||||
import type { NodePgDatabase } from 'drizzle-orm/node-postgres';
|
||||
|
||||
export type OutreachDb = NodePgDatabase<typeof schema>;
|
||||
|
||||
/** Per-run cap: bounds nightly spend; the backlog drains across nights. */
|
||||
const EMBED_BATCH_LIMIT = 500;
|
||||
/** gemini-embedding-001 batch endpoint cap handled in the ai package (100). */
|
||||
const STORE_CHUNK_SIZE = 100;
|
||||
|
||||
export interface EmbedGrantsDeps {
|
||||
readonly db: OutreachDb;
|
||||
}
|
||||
|
||||
let registeredDeps: EmbedGrantsDeps | null = null;
|
||||
|
||||
export function setEmbedGrantsDeps(deps: EmbedGrantsDeps): void {
|
||||
registeredDeps = deps;
|
||||
}
|
||||
|
||||
function getEmbedGrantsDeps(): EmbedGrantsDeps {
|
||||
if (registeredDeps == null) {
|
||||
throw new Error(
|
||||
'EmbedGrantsDeps not registered. Call setEmbedGrantsDeps() before DBOS.launch().',
|
||||
);
|
||||
}
|
||||
return registeredDeps;
|
||||
}
|
||||
|
||||
async function listGrantsNeedingEmbedding(
|
||||
db: OutreachDb,
|
||||
): Promise<GrantNeedingEmbedding[]> {
|
||||
return serverListGrantsNeedingEmbedding(db, { limit: EMBED_BATCH_LIMIT });
|
||||
}
|
||||
const listGrantsNeedingEmbeddingStep = DBOS.registerStep(
|
||||
listGrantsNeedingEmbedding,
|
||||
{ name: 'listGrantsNeedingEmbedding', retriesAllowed: true, maxAttempts: 3 },
|
||||
);
|
||||
|
||||
async function embedGrantChunk(
|
||||
grants: GrantNeedingEmbedding[],
|
||||
): Promise<number[][]> {
|
||||
const texts = grants.map((grant) => buildGrantEmbeddingText(grant));
|
||||
return generateDocumentEmbeddings(texts);
|
||||
}
|
||||
const embedGrantChunkStep = DBOS.registerStep(embedGrantChunk, {
|
||||
name: 'embedGrantChunk',
|
||||
retriesAllowed: true,
|
||||
maxAttempts: 3,
|
||||
});
|
||||
|
||||
async function storeGrantEmbeddings(
|
||||
db: OutreachDb,
|
||||
grants: GrantNeedingEmbedding[],
|
||||
vectors: number[][],
|
||||
): Promise<void> {
|
||||
await serverSetGrantEmbeddings(
|
||||
db,
|
||||
grants.map((grant, i) => ({ grantId: grant.id, embedding: vectors[i]! })),
|
||||
);
|
||||
}
|
||||
const storeGrantEmbeddingsStep = DBOS.registerStep(storeGrantEmbeddings, {
|
||||
name: 'storeGrantEmbeddings',
|
||||
retriesAllowed: true,
|
||||
maxAttempts: 3,
|
||||
});
|
||||
|
||||
async function runEmbedGrants(): Promise<void> {
|
||||
const { db } = getEmbedGrantsDeps();
|
||||
|
||||
const pending = await listGrantsNeedingEmbeddingStep(db);
|
||||
if (pending.length === 0) {
|
||||
console.log('[embed-grants] nothing to embed');
|
||||
return;
|
||||
}
|
||||
|
||||
let embedded = 0;
|
||||
// Chunked embed→store so a mid-run failure resumes from the last stored
|
||||
// chunk instead of re-paying for the whole batch.
|
||||
for (let i = 0; i < pending.length; i += STORE_CHUNK_SIZE) {
|
||||
const chunk = pending.slice(i, i + STORE_CHUNK_SIZE);
|
||||
const vectors = await embedGrantChunkStep(chunk);
|
||||
if (vectors.length !== chunk.length) {
|
||||
throw new Error(
|
||||
`[embed-grants] embedding count mismatch: ${vectors.length} vectors for ${chunk.length} grants`,
|
||||
);
|
||||
}
|
||||
await storeGrantEmbeddingsStep(db, chunk, vectors);
|
||||
embedded += chunk.length;
|
||||
}
|
||||
|
||||
console.log(
|
||||
`[embed-grants] embedded=${embedded} (of ${pending.length} pending this run)`,
|
||||
);
|
||||
}
|
||||
|
||||
const g = globalThis as unknown as {
|
||||
__outreachEmbedGrantsRegistered?: boolean;
|
||||
__outreachEmbedGrantsHandle?: (
|
||||
scheduledTime: Date,
|
||||
startedAt: Date,
|
||||
) => Promise<void>;
|
||||
};
|
||||
|
||||
if (!g.__outreachEmbedGrantsRegistered) {
|
||||
g.__outreachEmbedGrantsRegistered = true;
|
||||
|
||||
const embedGrants = async (_scheduledTime: Date, _startedAt: Date) => {
|
||||
try {
|
||||
await runEmbedGrants();
|
||||
} catch (err) {
|
||||
console.error('[embed-grants] pass failed:', err);
|
||||
throw err;
|
||||
}
|
||||
};
|
||||
|
||||
// Must be registered as BOTH a workflow and a scheduled function,
|
||||
// referencing the same function object — see module doc comment.
|
||||
g.__outreachEmbedGrantsHandle = DBOS.registerWorkflow(embedGrants, {
|
||||
name: 'embedGrants',
|
||||
});
|
||||
DBOS.registerScheduled(embedGrants, {
|
||||
crontab: '15 4 * * *',
|
||||
name: 'embedGrants',
|
||||
mode: SchedulerMode.ExactlyOncePerInterval,
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Starts one durable run of this workflow immediately through DBOS —
|
||||
* the exact production path (workflow + checkpointed steps), used by
|
||||
* `run-once.ts` for supervised/manual passes. Requires deps injected and
|
||||
* `DBOS.launch()` completed.
|
||||
*/
|
||||
export function runEmbedGrantsNow(): Promise<void> {
|
||||
const handle = g.__outreachEmbedGrantsHandle;
|
||||
if (handle == null) {
|
||||
throw new Error(
|
||||
'embedGrants is not registered; was this module imported before DBOS.launch()?',
|
||||
);
|
||||
}
|
||||
return handle(new Date(), new Date());
|
||||
}
|
||||
Reference in New Issue
Block a user