1import { createHash } from 'node:crypto';
2
3export const ENGINE_VERSION = '1.0.0';
4const DEFAULT_TRACKED_FIELDS = ['title', 'company', 'location', 'url', 'description', 'employmentType', 'salary'];
5export type JobRecord = Record<string, string | number | boolean | null | undefined> & {
6 sourceId?: string;
7 title: string;
8 company: string;
9 location?: string;
10 url?: string;
11};
12export interface Change {
13 recordType: 'job-change';
14 comparisonId: string;
15 stableKey: string;
16 changeType: 'ADDED' | 'REMOVED' | 'CHANGED';
17 changedFields: { field: string; before: unknown; after: unknown }[];
18 baseline: JobRecord | null;
19 current: JobRecord | null;
20 engineVersion: string;
21}
22export interface DuplicateGroup {
23 recordType: 'duplicate-group';
24 comparisonId: string;
25 snapshot: 'baseline' | 'current';
26 stableKey: string;
27 indexes: number[];
28 count: number;
29 engineVersion: string;
30}
31export interface RunResult {
32 inputHash: string;
33 changes: Change[];
34 duplicates: DuplicateGroup[];
35 summary: {
36 recordType: 'job-summary';
37 status: 'SUCCEEDED';
38 comparisonId: string;
39 baselineRecords: number;
40 currentRecords: number;
41 added: number;
42 removed: number;
43 changed: number;
44 unchanged: number;
45 duplicateGroups: number;
46 inputHash: string;
47 engineVersion: string;
48 };
49}
50export class DomainError extends Error {
51 public constructor(public readonly code: string, message: string) {
52 super(message);
53 }
54}
55function canonical(value: unknown): string {
56 if (Array.isArray(value)) return `[${value.map(canonical).join(',')}]`;
57 if (value && typeof value === 'object') {
58 return `{${Object.entries(value as Record<string, unknown>)
59 .sort(([a], [b]) => a.localeCompare(b))
60 .map(([key, child]) => `${JSON.stringify(key)}:${canonical(child)}`)
61 .join(',')}}`;
62 }
63 return JSON.stringify(value);
64}
65function hash(value: unknown): string {
66 return createHash('sha256').update(canonical(value)).digest('hex');
67}
68function object(value: unknown, label: string): Record<string, unknown> {
69 if (!value || typeof value !== 'object' || Array.isArray(value))
70 throw new DomainError('INVALID_INPUT', `${label} must be an object`);
71 return value as Record<string, unknown>;
72}
73function text(value: unknown): string {
74 return typeof value === 'string'
75 ? value.normalize('NFKC').toLocaleLowerCase('en-US').trim().replace(/\s+/g, ' ')
76 : String(value ?? '');
77}
78function canonicalUrl(value: string): string {
79 let url: URL;
80 try {
81 url = new URL(value);
82 } catch {
83 throw new DomainError('INVALID_INPUT', 'job url must be a valid absolute HTTP or HTTPS URL');
84 }
85 if (!['http:', 'https:'].includes(url.protocol))
86 throw new DomainError('INVALID_INPUT', 'job url must use HTTP or HTTPS');
87 url.hash = '';
88 url.hostname = url.hostname.toLowerCase().replace(/^www\./, '');
89 for (const key of [...url.searchParams.keys()]) {
90 if (/^(utm_|gclid$|fbclid$|ref$)/i.test(key)) url.searchParams.delete(key);
91 }
92 url.searchParams.sort();
93 url.pathname = url.pathname.replace(/\/+$/, '') || '/';
94 return url.toString();
95}
96export function stableJobKey(record: JobRecord): string {
97 if (record.sourceId && text(record.sourceId)) return `source:${text(record.sourceId)}`;
98 if (record.url && text(record.url)) return `url:${canonicalUrl(record.url)}`;
99 const fallback = [record.company, record.title, record.location].map(text);
100 if (!fallback[0] || !fallback[1] || !fallback[2])
101 throw new DomainError('INVALID_INPUT', 'job needs sourceId, url, or company/title/location');
102 return `tuple:${fallback.join('|')}`;
103}
104function parseRows(value: unknown, label: string, maximum: number): JobRecord[] {
105 if (!Array.isArray(value) || value.length > maximum)
106 throw new DomainError('RESOURCE_LIMIT_EXCEEDED', `${label} must be an array with at most ${maximum} records`);
107 return value.map((candidate, index) => {
108 const row = object(candidate, `${label}[${index}]`);
109 if (typeof row.title !== 'string' || !row.title.trim() || row.title.length > 500)
110 throw new DomainError('INVALID_INPUT', `${label}[${index}].title is invalid`);
111 if (typeof row.company !== 'string' || !row.company.trim() || row.company.length > 500)
112 throw new DomainError('INVALID_INPUT', `${label}[${index}].company is invalid`);
113 for (const [key, item] of Object.entries(row)) {
114 if (item !== null && !['string', 'number', 'boolean', 'undefined'].includes(typeof item))
115 throw new DomainError('INVALID_INPUT', `${label}[${index}].${key} must be scalar`);
116 if (typeof item === 'string' && item.length > 20_000)
117 throw new DomainError('RESOURCE_LIMIT_EXCEEDED', `${label}[${index}].${key} is too large`);
118 }
119 const record = row as JobRecord;
120 stableJobKey(record);
121 return record;
122 });
123}
124function indexRows(rows: JobRecord[], snapshot: 'baseline' | 'current'): {
125 canonical: Map<string, JobRecord>;
126 duplicates: Omit<DuplicateGroup, 'comparisonId'>[];
127} {
128 const grouped = new Map<string, { record: JobRecord; indexes: number[] }>();
129 rows.forEach((record, index) => {
130 const key = stableJobKey(record);
131 const group = grouped.get(key);
132 if (group) group.indexes.push(index);
133 else grouped.set(key, { record, indexes: [index] });
134 });
135 return {
136 canonical: new Map([...grouped].map(([key, value]) => [key, value.record])),
137 duplicates: [...grouped]
138 .filter(([, value]) => value.indexes.length > 1)
139 .map(([stableKey, value]) => ({
140 recordType: 'duplicate-group' as const,
141 snapshot,
142 stableKey,
143 indexes: value.indexes,
144 count: value.indexes.length,
145 engineVersion: ENGINE_VERSION,
146 }))
147 .sort((a, b) => a.stableKey.localeCompare(b.stableKey)),
148 };
149}
150export function processInput(value: unknown): RunResult {
151 const input = object(value, 'input');
152 if (typeof input.comparisonId !== 'string' || !input.comparisonId.trim() || input.comparisonId.length > 120)
153 throw new DomainError('INVALID_INPUT', 'comparisonId is required');
154 const comparisonId = input.comparisonId;
155 const maximum = input.maxRecordsPerSnapshot ?? 1_000;
156 if (!Number.isInteger(maximum) || (maximum as number) < 1 || (maximum as number) > 1_000)
157 throw new DomainError('INVALID_INPUT', 'maxRecordsPerSnapshot must be an integer from 1 to 1000');
158 const baseline = parseRows(input.baseline, 'baseline', maximum as number);
159 const current = parseRows(input.current, 'current', maximum as number);
160 if (baseline.length + current.length === 0)
161 throw new DomainError('INVALID_INPUT', 'at least one snapshot must contain a record');
162 if (Buffer.byteLength(canonical({ baseline, current }), 'utf8') > 10_000_000)
163 throw new DomainError('RESOURCE_LIMIT_EXCEEDED', 'snapshot input exceeds 10 MB');
164 const tracked = input.trackedFields ?? DEFAULT_TRACKED_FIELDS;
165 if (!Array.isArray(tracked) || tracked.length < 1 || tracked.length > 30 ||
166 tracked.some((field) => typeof field !== 'string' || !field || field.length > 100))
167 throw new DomainError('INVALID_INPUT', 'trackedFields must contain 1 to 30 field names');
168 if (new Set(tracked).size !== tracked.length)
169 throw new DomainError('INVALID_INPUT', 'trackedFields must be unique');
170 const oldIndex = indexRows(baseline, 'baseline');
171 const newIndex = indexRows(current, 'current');
172 const allKeys = [...new Set([...oldIndex.canonical.keys(), ...newIndex.canonical.keys()])].sort();
173 let unchanged = 0;
174 const changes: Change[] = [];
175 for (const key of allKeys) {
176 const before = oldIndex.canonical.get(key) ?? null;
177 const after = newIndex.canonical.get(key) ?? null;
178 let changeType: Change['changeType'];
179 let changedFields: Change['changedFields'];
180 if (!before) {
181 changeType = 'ADDED';
182 changedFields = [];
183 } else if (!after) {
184 changeType = 'REMOVED';
185 changedFields = [];
186 } else {
187 changedFields = (tracked as string[]).flatMap((field) =>
188 canonical(before[field]) === canonical(after[field])
189 ? []
190 : [{ field, before: before[field] ?? null, after: after[field] ?? null }],
191 );
192 if (!changedFields.length) {
193 unchanged += 1;
194 continue;
195 }
196 changeType = 'CHANGED';
197 }
198 changes.push({
199 recordType: 'job-change',
200 comparisonId: comparisonId.trim(),
201 stableKey: key,
202 changeType,
203 changedFields,
204 baseline: before,
205 current: after,
206 engineVersion: ENGINE_VERSION,
207 });
208 }
209 const duplicates = [...oldIndex.duplicates, ...newIndex.duplicates].map((item) => ({
210 ...item,
211 comparisonId: comparisonId.trim(),
212 }));
213 const inputHash = hash({
214 comparisonId: comparisonId.trim(),
215 baseline,
216 current,
217 trackedFields: tracked,
218 });
219 return {
220 inputHash,
221 changes,
222 duplicates,
223 summary: {
224 recordType: 'job-summary',
225 status: 'SUCCEEDED',
226 comparisonId: comparisonId.trim(),
227 baselineRecords: baseline.length,
228 currentRecords: current.length,
229 added: changes.filter((item) => item.changeType === 'ADDED').length,
230 removed: changes.filter((item) => item.changeType === 'REMOVED').length,
231 changed: changes.filter((item) => item.changeType === 'CHANGED').length,
232 unchanged,
233 duplicateGroups: duplicates.length,
234 inputHash,
235 engineVersion: ENGINE_VERSION,
236 },
237 };
238}