Auto-Registering Extractor System (#223)
* initial commit? * Address PR feedback on extractor discovery and startup resilience * Address latest PR review comments * fix city resolution fallback when input parses empty * address PR feedback on extractor registry and pipeline validation * address copilot comments on manifests and registry startup * fix extractor discovery export handling and env isolation in tests * enforce duplicate manifest id failures in strict mode * Fix remaining extractor registry and runtime review comments * docs * docs * test all, logic remains in extractors * Address PR review feedback on extractor registry and validation * Revert extractor moduleResolution to bundler * Enforce shared city filtering across all discovery sources * Deduplicate extractor strict city post-filtering
This commit is contained in:
@@ -0,0 +1,94 @@
|
||||
import { resolveSearchCities } from "@shared/search-cities.js";
|
||||
import type {
|
||||
ExtractorManifest,
|
||||
ExtractorProgressEvent,
|
||||
} from "@shared/types/extractors";
|
||||
import { runHiringCafe } from "./src/run";
|
||||
|
||||
function toProgress(event: {
|
||||
type: string;
|
||||
termIndex: number;
|
||||
termTotal: number;
|
||||
searchTerm: string;
|
||||
pageNo?: number;
|
||||
totalCollected?: number;
|
||||
}): ExtractorProgressEvent {
|
||||
if (event.type === "term_start") {
|
||||
return {
|
||||
phase: "list",
|
||||
termsProcessed: Math.max(event.termIndex - 1, 0),
|
||||
termsTotal: event.termTotal,
|
||||
currentUrl: event.searchTerm,
|
||||
detail: `Hiring Cafe: term ${event.termIndex}/${event.termTotal} (${event.searchTerm})`,
|
||||
};
|
||||
}
|
||||
|
||||
if (event.type === "page_fetched") {
|
||||
const pageNo = (event.pageNo ?? 0) + 1;
|
||||
const totalCollected = event.totalCollected ?? 0;
|
||||
return {
|
||||
phase: "list",
|
||||
termsProcessed: Math.max(event.termIndex - 1, 0),
|
||||
termsTotal: event.termTotal,
|
||||
listPagesProcessed: pageNo,
|
||||
jobPagesEnqueued: totalCollected,
|
||||
jobPagesProcessed: totalCollected,
|
||||
currentUrl: `page ${pageNo}`,
|
||||
detail: `Hiring Cafe: term ${event.termIndex}/${event.termTotal}, page ${pageNo} (${totalCollected} collected)`,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
phase: "list",
|
||||
termsProcessed: event.termIndex,
|
||||
termsTotal: event.termTotal,
|
||||
currentUrl: event.searchTerm,
|
||||
detail: `Hiring Cafe: completed term ${event.termIndex}/${event.termTotal} (${event.searchTerm})`,
|
||||
};
|
||||
}
|
||||
|
||||
export const manifest: ExtractorManifest = {
|
||||
id: "hiringcafe",
|
||||
displayName: "Hiring Cafe",
|
||||
providesSources: ["hiringcafe"],
|
||||
async run(context) {
|
||||
if (context.shouldCancel?.()) {
|
||||
return { success: true, jobs: [] };
|
||||
}
|
||||
|
||||
const maxJobsPerTerm = context.settings.jobspyResultsWanted
|
||||
? parseInt(context.settings.jobspyResultsWanted, 10)
|
||||
: 200;
|
||||
|
||||
const result = await runHiringCafe({
|
||||
country: context.selectedCountry,
|
||||
countryKey: context.selectedCountry,
|
||||
searchTerms: context.searchTerms,
|
||||
locations: resolveSearchCities({
|
||||
single:
|
||||
context.settings.searchCities ?? context.settings.jobspyLocation,
|
||||
}),
|
||||
maxJobsPerTerm,
|
||||
onProgress: (event) => {
|
||||
if (context.shouldCancel?.()) return;
|
||||
|
||||
context.onProgress?.(toProgress(event));
|
||||
},
|
||||
});
|
||||
|
||||
if (!result.success) {
|
||||
return {
|
||||
success: false,
|
||||
jobs: [],
|
||||
error: result.error,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
success: true,
|
||||
jobs: result.jobs,
|
||||
};
|
||||
},
|
||||
};
|
||||
|
||||
export default manifest;
|
||||
@@ -0,0 +1,307 @@
|
||||
import { spawn, spawnSync } from "node:child_process";
|
||||
import { mkdir, readFile, rm } from "node:fs/promises";
|
||||
import { createRequire } from "node:module";
|
||||
import { dirname, join } from "node:path";
|
||||
import { createInterface } from "node:readline";
|
||||
import { fileURLToPath } from "node:url";
|
||||
import { normalizeCountryKey } from "@shared/location-support.js";
|
||||
import {
|
||||
resolveSearchCities,
|
||||
shouldApplyStrictCityFilter,
|
||||
} from "@shared/search-cities.js";
|
||||
import type { CreateJobInput } from "@shared/types/jobs";
|
||||
import {
|
||||
toNumberOrNull,
|
||||
toStringOrNull,
|
||||
} from "@shared/utils/type-conversion.js";
|
||||
|
||||
const srcDir = dirname(fileURLToPath(import.meta.url));
|
||||
const EXTRACTOR_DIR = join(srcDir, "..");
|
||||
const DATASET_PATH = join(EXTRACTOR_DIR, "storage/datasets/default/jobs.json");
|
||||
const STORAGE_DATASET_DIR = join(EXTRACTOR_DIR, "storage/datasets/default");
|
||||
const JOBOPS_PROGRESS_PREFIX = "JOBOPS_PROGRESS ";
|
||||
const require = createRequire(import.meta.url);
|
||||
const TSX_CLI_PATH = resolveTsxCliPath();
|
||||
|
||||
type HiringCafeRawJob = Record<string, unknown>;
|
||||
|
||||
export type HiringCafeProgressEvent =
|
||||
| {
|
||||
type: "term_start";
|
||||
termIndex: number;
|
||||
termTotal: number;
|
||||
searchTerm: string;
|
||||
}
|
||||
| {
|
||||
type: "page_fetched";
|
||||
termIndex: number;
|
||||
termTotal: number;
|
||||
searchTerm: string;
|
||||
pageNo: number;
|
||||
resultsOnPage: number;
|
||||
totalCollected: number;
|
||||
}
|
||||
| {
|
||||
type: "term_complete";
|
||||
termIndex: number;
|
||||
termTotal: number;
|
||||
searchTerm: string;
|
||||
jobsFoundTerm: number;
|
||||
};
|
||||
|
||||
export interface RunHiringCafeOptions {
|
||||
searchTerms?: string[];
|
||||
country?: string;
|
||||
countryKey?: string;
|
||||
locations?: string[];
|
||||
locationRadiusMiles?: number;
|
||||
maxJobsPerTerm?: number;
|
||||
onProgress?: (event: HiringCafeProgressEvent) => void;
|
||||
}
|
||||
|
||||
export interface HiringCafeResult {
|
||||
success: boolean;
|
||||
jobs: CreateJobInput[];
|
||||
error?: string;
|
||||
}
|
||||
|
||||
export function shouldApplyStrictLocationFilter(
|
||||
location: string,
|
||||
countryKey: string,
|
||||
): boolean {
|
||||
return shouldApplyStrictCityFilter(location, countryKey);
|
||||
}
|
||||
|
||||
function resolveTsxCliPath(): string | null {
|
||||
try {
|
||||
return require.resolve("tsx/dist/cli.mjs");
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
function canRunNpmCommand(): boolean {
|
||||
const result = spawnSync("npm", ["--version"], { stdio: "ignore" });
|
||||
return !result.error && result.status === 0;
|
||||
}
|
||||
|
||||
function parseProgressLine(line: string): HiringCafeProgressEvent | null {
|
||||
if (!line.startsWith(JOBOPS_PROGRESS_PREFIX)) return null;
|
||||
const raw = line.slice(JOBOPS_PROGRESS_PREFIX.length).trim();
|
||||
|
||||
let parsed: Record<string, unknown>;
|
||||
try {
|
||||
parsed = JSON.parse(raw) as Record<string, unknown>;
|
||||
} catch {
|
||||
return null;
|
||||
}
|
||||
|
||||
const event = toStringOrNull(parsed.event);
|
||||
const termIndex = toNumberOrNull(parsed.termIndex);
|
||||
const termTotal = toNumberOrNull(parsed.termTotal);
|
||||
const searchTerm = toStringOrNull(parsed.searchTerm) ?? "";
|
||||
if (!event || termIndex === null || termTotal === null) return null;
|
||||
|
||||
if (event === "term_start") {
|
||||
return { type: "term_start", termIndex, termTotal, searchTerm };
|
||||
}
|
||||
|
||||
if (event === "page_fetched") {
|
||||
const pageNo = toNumberOrNull(parsed.pageNo);
|
||||
if (pageNo === null) return null;
|
||||
return {
|
||||
type: "page_fetched",
|
||||
termIndex,
|
||||
termTotal,
|
||||
searchTerm,
|
||||
pageNo,
|
||||
resultsOnPage: toNumberOrNull(parsed.resultsOnPage) ?? 0,
|
||||
totalCollected: toNumberOrNull(parsed.totalCollected) ?? 0,
|
||||
};
|
||||
}
|
||||
|
||||
if (event === "term_complete") {
|
||||
return {
|
||||
type: "term_complete",
|
||||
termIndex,
|
||||
termTotal,
|
||||
searchTerm,
|
||||
jobsFoundTerm: toNumberOrNull(parsed.jobsFoundTerm) ?? 0,
|
||||
};
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
function mapHiringCafeRow(row: HiringCafeRawJob): CreateJobInput | null {
|
||||
const jobUrl = toStringOrNull(row.jobUrl);
|
||||
if (!jobUrl) return null;
|
||||
|
||||
return {
|
||||
source: "hiringcafe",
|
||||
sourceJobId: toStringOrNull(row.sourceJobId) ?? undefined,
|
||||
title: toStringOrNull(row.title) ?? "Unknown Title",
|
||||
employer: toStringOrNull(row.employer) ?? "Unknown Employer",
|
||||
jobUrl,
|
||||
applicationLink: toStringOrNull(row.applicationLink) ?? jobUrl,
|
||||
location: toStringOrNull(row.location) ?? undefined,
|
||||
salary: toStringOrNull(row.salary) ?? undefined,
|
||||
datePosted: toStringOrNull(row.datePosted) ?? undefined,
|
||||
jobDescription: toStringOrNull(row.jobDescription) ?? undefined,
|
||||
jobType: toStringOrNull(row.jobType) ?? undefined,
|
||||
};
|
||||
}
|
||||
|
||||
async function readDataset(): Promise<CreateJobInput[]> {
|
||||
const content = await readFile(DATASET_PATH, "utf-8");
|
||||
const parsed = JSON.parse(content) as unknown;
|
||||
if (!Array.isArray(parsed)) return [];
|
||||
|
||||
const jobs: CreateJobInput[] = [];
|
||||
const seen = new Set<string>();
|
||||
for (const value of parsed) {
|
||||
if (!value || typeof value !== "object" || Array.isArray(value)) continue;
|
||||
const mapped = mapHiringCafeRow(value as HiringCafeRawJob);
|
||||
if (!mapped) continue;
|
||||
const dedupeKey = mapped.sourceJobId || mapped.jobUrl;
|
||||
if (seen.has(dedupeKey)) continue;
|
||||
seen.add(dedupeKey);
|
||||
jobs.push(mapped);
|
||||
}
|
||||
return jobs;
|
||||
}
|
||||
|
||||
async function clearStorageDataset(): Promise<void> {
|
||||
await rm(STORAGE_DATASET_DIR, { recursive: true, force: true });
|
||||
await mkdir(STORAGE_DATASET_DIR, { recursive: true });
|
||||
}
|
||||
|
||||
export async function runHiringCafe(
|
||||
options: RunHiringCafeOptions = {},
|
||||
): Promise<HiringCafeResult> {
|
||||
const searchTerms =
|
||||
options.searchTerms && options.searchTerms.length > 0
|
||||
? options.searchTerms
|
||||
: ["web developer"];
|
||||
const country = (options.country || "united kingdom").trim().toLowerCase();
|
||||
const countryKey = normalizeCountryKey(options.countryKey ?? "");
|
||||
const maxJobsPerTerm = options.maxJobsPerTerm ?? 200;
|
||||
const locationRadiusMiles = Math.max(
|
||||
1,
|
||||
Math.floor(options.locationRadiusMiles ?? 1),
|
||||
);
|
||||
const locations = resolveSearchCities({
|
||||
list: options.locations,
|
||||
env: process.env.HIRING_CAFE_LOCATION_QUERY,
|
||||
});
|
||||
const runLocations = locations.length > 0 ? locations : [null];
|
||||
const termTotal = searchTerms.length * runLocations.length;
|
||||
|
||||
const useNpmCommand = canRunNpmCommand();
|
||||
if (!useNpmCommand && !TSX_CLI_PATH) {
|
||||
return {
|
||||
success: false,
|
||||
jobs: [],
|
||||
error: "Unable to execute Hiring Cafe extractor (npm/tsx unavailable)",
|
||||
};
|
||||
}
|
||||
|
||||
try {
|
||||
const jobs: CreateJobInput[] = [];
|
||||
const seen = new Set<string>();
|
||||
|
||||
for (let runIndex = 0; runIndex < runLocations.length; runIndex += 1) {
|
||||
const location = runLocations[runIndex];
|
||||
const strictLocationFilter =
|
||||
location !== null &&
|
||||
shouldApplyStrictLocationFilter(location, countryKey);
|
||||
|
||||
await clearStorageDataset();
|
||||
|
||||
await new Promise<void>((resolve, reject) => {
|
||||
const extractorEnv = {
|
||||
...process.env,
|
||||
JOBOPS_EMIT_PROGRESS: "1",
|
||||
HIRING_CAFE_SEARCH_TERMS: JSON.stringify(searchTerms),
|
||||
HIRING_CAFE_COUNTRY: country,
|
||||
HIRING_CAFE_MAX_JOBS_PER_TERM: String(maxJobsPerTerm),
|
||||
HIRING_CAFE_OUTPUT_JSON: DATASET_PATH,
|
||||
HIRING_CAFE_LOCATION_QUERY: strictLocationFilter ? location : "",
|
||||
HIRING_CAFE_LOCATION_RADIUS_MILES: strictLocationFilter
|
||||
? String(locationRadiusMiles)
|
||||
: "",
|
||||
};
|
||||
|
||||
const child = useNpmCommand
|
||||
? spawn("npm", ["run", "start"], {
|
||||
cwd: EXTRACTOR_DIR,
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
env: extractorEnv,
|
||||
})
|
||||
: (() => {
|
||||
const tsxCliPath = TSX_CLI_PATH;
|
||||
if (!tsxCliPath) {
|
||||
throw new Error(
|
||||
"Unable to execute Hiring Cafe extractor (npm/tsx unavailable)",
|
||||
);
|
||||
}
|
||||
|
||||
return spawn(process.execPath, [tsxCliPath, "src/main.ts"], {
|
||||
cwd: EXTRACTOR_DIR,
|
||||
stdio: ["ignore", "pipe", "pipe"],
|
||||
env: extractorEnv,
|
||||
});
|
||||
})();
|
||||
|
||||
const handleLine = (line: string, stream: NodeJS.WriteStream) => {
|
||||
const progressEvent = parseProgressLine(line);
|
||||
if (progressEvent) {
|
||||
const termOffset = runIndex * searchTerms.length;
|
||||
options.onProgress?.({
|
||||
...progressEvent,
|
||||
termIndex: termOffset + progressEvent.termIndex,
|
||||
termTotal,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
stream.write(`${line}\n`);
|
||||
};
|
||||
|
||||
const stdoutRl = child.stdout
|
||||
? createInterface({ input: child.stdout })
|
||||
: null;
|
||||
const stderrRl = child.stderr
|
||||
? createInterface({ input: child.stderr })
|
||||
: null;
|
||||
|
||||
stdoutRl?.on("line", (line) => handleLine(line, process.stdout));
|
||||
stderrRl?.on("line", (line) => handleLine(line, process.stderr));
|
||||
|
||||
child.on("close", (code) => {
|
||||
stdoutRl?.close();
|
||||
stderrRl?.close();
|
||||
if (code === 0) resolve();
|
||||
else
|
||||
reject(new Error(`Hiring Cafe extractor exited with code ${code}`));
|
||||
});
|
||||
child.on("error", reject);
|
||||
});
|
||||
|
||||
const runJobs = await readDataset();
|
||||
const filtered = runJobs;
|
||||
|
||||
for (const job of filtered) {
|
||||
const key = job.sourceJobId || job.jobUrl;
|
||||
if (seen.has(key)) continue;
|
||||
seen.add(key);
|
||||
jobs.push(job);
|
||||
}
|
||||
}
|
||||
|
||||
return { success: true, jobs };
|
||||
} catch (error) {
|
||||
const message = error instanceof Error ? error.message : "Unknown error";
|
||||
return { success: false, jobs: [], error: message };
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,15 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { shouldApplyStrictLocationFilter } from "../src/run";
|
||||
|
||||
describe("hiringcafe location query strictness", () => {
|
||||
it("enables strict filtering when city differs from country", () => {
|
||||
expect(shouldApplyStrictLocationFilter("Leeds", "united kingdom")).toBe(
|
||||
true,
|
||||
);
|
||||
});
|
||||
|
||||
it("disables strict filtering when location is country-level", () => {
|
||||
expect(shouldApplyStrictLocationFilter("UK", "united kingdom")).toBe(false);
|
||||
expect(shouldApplyStrictLocationFilter("United States", "us")).toBe(false);
|
||||
});
|
||||
});
|
||||
@@ -1,13 +1,17 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"module": "NodeNext",
|
||||
"moduleResolution": "NodeNext",
|
||||
"module": "ESNext",
|
||||
"moduleResolution": "bundler",
|
||||
"target": "ES2022",
|
||||
"outDir": "dist",
|
||||
"strict": true,
|
||||
"noUnusedLocals": false,
|
||||
"lib": ["ES2022", "DOM"],
|
||||
"types": ["node"]
|
||||
"types": ["node"],
|
||||
"baseUrl": ".",
|
||||
"paths": {
|
||||
"@shared/*": ["../../shared/src/*"]
|
||||
}
|
||||
},
|
||||
"include": ["./src/**/*"]
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user