1import { createHash } from 'node:crypto';
2
3export const ENGINE_VERSION = '1.0.0';
4const SUPPORTED_REQUIRED_FIELDS = ['id', 'itemId', 'author', 'rating', 'text', 'publishedAt'] as const;
5type RequiredField = (typeof SUPPORTED_REQUIRED_FIELDS)[number];
6export interface ReviewRecord {
7 id?: string;
8 itemId?: string;
9 author?: string;
10 rating?: number;
11 text?: string;
12 publishedAt?: string;
13}
14export interface Finding {
15 recordType: 'quality-finding';
16 batchId: string;
17 code: 'MISSING_FIELD' | 'RATING_RANGE' | 'INVALID_TIMESTAMP' | 'FUTURE_TIMESTAMP' | 'DUPLICATE_ID' | 'REPEATED_CONTENT';
18 severity: 'ERROR' | 'WARNING';
19 reviewIds: string[];
20 field: string | null;
21 message: string;
22 fingerprint: string;
23 engineVersion: string;
24}
25export interface RunResult {
26 inputHash: string;
27 findings: Finding[];
28 summary: {
29 recordType: 'quality-summary';
30 status: 'PASSED' | 'FAILED';
31 batchId: string;
32 reviewCount: number;
33 findingCount: number;
34 returnedFindings: number;
35 findingsTruncated: boolean;
36 errorCount: number;
37 warningCount: number;
38 repeatedContentGroups: number;
39 duplicateIdGroups: number;
40 inputHash: string;
41 engineVersion: string;
42 };
43}
44export class DomainError extends Error {
45 public constructor(public readonly code: string, message: string) {
46 super(message);
47 }
48}
49function canonical(value: unknown): string {
50 if (Array.isArray(value)) return `[${value.map(canonical).join(',')}]`;
51 if (value && typeof value === 'object') {
52 return `{${Object.entries(value as Record<string, unknown>)
53 .sort(([a], [b]) => a.localeCompare(b))
54 .map(([key, child]) => `${JSON.stringify(key)}:${canonical(child)}`)
55 .join(',')}}`;
56 }
57 return JSON.stringify(value);
58}
59function hash(value: unknown): string {
60 return createHash('sha256').update(canonical(value)).digest('hex');
61}
62function object(value: unknown, label: string): Record<string, unknown> {
63 if (!value || typeof value !== 'object' || Array.isArray(value))
64 throw new DomainError('INVALID_INPUT', `${label} must be an object`);
65 return value as Record<string, unknown>;
66}
67function integer(value: unknown, fallback: number, min: number, max: number, label: string): number {
68 const actual = value ?? fallback;
69 if (!Number.isInteger(actual) || (actual as number) < min || (actual as number) > max)
70 throw new DomainError('INVALID_INPUT', `${label} must be an integer from ${min} to ${max}`);
71 return actual as number;
72}
73function normalizeText(value: string): string {
74 return value.normalize('NFKC').toLocaleLowerCase('en-US').replace(/[^\p{L}\p{N}]+/gu, ' ').trim().replace(/\s+/g, ' ');
75}
76function validTimestamp(value: string): boolean {
77 if (!/^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}(?:\.\d{1,3})?Z$/.test(value)) return false;
78 const parsed = Date.parse(value);
79 return Number.isFinite(parsed) && new Date(parsed).toISOString().startsWith(value.slice(0, 19));
80}
81export function processInput(value: unknown): RunResult {
82 const input = object(value, 'input');
83 if (typeof input.batchId !== 'string' || !input.batchId.trim() || input.batchId.length > 120)
84 throw new DomainError('INVALID_INPUT', 'batchId is required');
85 const batchId = input.batchId;
86 const maximum = integer(input.maxRecords, 20_000, 1, 20_000, 'maxRecords');
87 const findingLimit = integer(input.maxFindings, 1_000, 1, 1_000, 'maxFindings');
88 if (!Array.isArray(input.reviews) || input.reviews.length < 1 || input.reviews.length > maximum)
89 throw new DomainError('RESOURCE_LIMIT_EXCEEDED', `reviews must contain 1 to ${maximum} records`);
90 if (Buffer.byteLength(canonical(input.reviews), 'utf8') > 10_000_000)
91 throw new DomainError('RESOURCE_LIMIT_EXCEEDED', 'reviews exceed 10 MB');
92 const ratingMin = input.ratingMin ?? 1;
93 const ratingMax = input.ratingMax ?? 5;
94 if (typeof ratingMin !== 'number' || !Number.isFinite(ratingMin) ||
95 typeof ratingMax !== 'number' || !Number.isFinite(ratingMax) || ratingMin >= ratingMax)
96 throw new DomainError('INVALID_INPUT', 'ratingMin and ratingMax must be finite numbers with min < max');
97 const required = input.requiredFields ?? SUPPORTED_REQUIRED_FIELDS;
98 if (!Array.isArray(required) || required.some((field) => !SUPPORTED_REQUIRED_FIELDS.includes(field as RequiredField)))
99 throw new DomainError('INVALID_INPUT', 'requiredFields contains an unsupported field');
100 if (new Set(required).size !== required.length)
101 throw new DomainError('INVALID_INPUT', 'requiredFields must be unique');
102 let observedAt: number | null = null;
103 if (input.observedAt !== undefined) {
104 if (typeof input.observedAt !== 'string' || !validTimestamp(input.observedAt))
105 throw new DomainError('INVALID_INPUT', 'observedAt must be a valid UTC timestamp');
106 observedAt = Date.parse(input.observedAt);
107 }
108 const reviews = input.reviews.map((candidate, index) => {
109 const row = object(candidate, `reviews[${index}]`);
110 for (const field of SUPPORTED_REQUIRED_FIELDS) {
111 const item = row[field];
112 if (item !== undefined && item !== null && field !== 'rating' && typeof item !== 'string')
113 throw new DomainError('INVALID_INPUT', `reviews[${index}].${field} must be a string`);
114 }
115 if (row.rating !== undefined && row.rating !== null &&
116 (typeof row.rating !== 'number' || !Number.isFinite(row.rating)))
117 throw new DomainError('INVALID_INPUT', `reviews[${index}].rating must be a finite number`);
118 for (const field of ['id', 'itemId', 'author', 'text', 'publishedAt'] as const) {
119 if (typeof row[field] === 'string' && row[field].length > (field === 'text' ? 20_000 : 1_000))
120 throw new DomainError('RESOURCE_LIMIT_EXCEEDED', `reviews[${index}].${field} is too large`);
121 }
122 return row as ReviewRecord;
123 });
124 const all: Finding[] = [];
125 const add = (code: Finding['code'], severity: Finding['severity'], ids: string[], field: string | null, message: string): void => {
126 all.push({
127 recordType: 'quality-finding',
128 batchId: batchId.trim(),
129 code,
130 severity,
131 reviewIds: ids,
132 field,
133 message,
134 fingerprint: hash({ code, ids, field, message }).slice(0, 24),
135 engineVersion: ENGINE_VERSION,
136 });
137 };
138 reviews.forEach((review, index) => {
139 const id = typeof review.id === 'string' && review.id.trim() ? review.id.trim() : `index:${index}`;
140 for (const field of required as RequiredField[]) {
141 const item = review[field];
142 if (item === undefined || item === null || (typeof item === 'string' && !item.trim()))
143 add('MISSING_FIELD', 'ERROR', [id], field, `${field} is required`);
144 }
145 if (typeof review.rating === 'number' && (review.rating < ratingMin || review.rating > ratingMax))
146 add('RATING_RANGE', 'ERROR', [id], 'rating', `rating must be from ${ratingMin} to ${ratingMax}`);
147 if (typeof review.publishedAt === 'string' && review.publishedAt.trim()) {
148 if (!validTimestamp(review.publishedAt))
149 add('INVALID_TIMESTAMP', 'ERROR', [id], 'publishedAt', 'publishedAt must be a valid UTC timestamp');
150 else if (observedAt !== null && Date.parse(review.publishedAt) > observedAt)
151 add('FUTURE_TIMESTAMP', 'ERROR', [id], 'publishedAt', 'publishedAt is later than observedAt');
152 }
153 });
154 const groupBy = (selector: (review: ReviewRecord) => string | null): Map<string, string[]> => {
155 const groups = new Map<string, string[]>();
156 reviews.forEach((review, index) => {
157 const key = selector(review);
158 if (!key) return;
159 const id = review.id?.trim() || `index:${index}`;
160 groups.set(key, [...(groups.get(key) ?? []), id]);
161 });
162 return groups;
163 };
164 const duplicateIds = [...groupBy((review) => review.id?.trim() || null)]
165 .filter(([, ids]) => ids.length > 1)
166 .sort(([a], [b]) => a.localeCompare(b));
167 for (const [id, ids] of duplicateIds)
168 add('DUPLICATE_ID', 'ERROR', ids, 'id', `review ID ${id} occurs ${ids.length} times`);
169 const repeated = [...groupBy((review) => {
170 const normalized = review.text ? normalizeText(review.text) : '';
171 return normalized.length >= 8 ? normalized : null;
172 })].filter(([, ids]) => ids.length > 1).sort(([a], [b]) => a.localeCompare(b));
173 for (const [, ids] of repeated)
174 add('REPEATED_CONTENT', 'WARNING', [...ids].sort(), 'text', `normalized review text repeats across ${ids.length} records`);
175 const sorted = all.sort((a, b) =>
176 a.reviewIds[0]!.localeCompare(b.reviewIds[0]!) || a.code.localeCompare(b.code) || (a.field ?? '').localeCompare(b.field ?? ''),
177 );
178 const inputHash = hash({
179 batchId: batchId.trim(),
180 reviews,
181 ratingMin,
182 ratingMax,
183 requiredFields: required,
184 observedAt: input.observedAt ?? null,
185 });
186 return {
187 inputHash,
188 findings: sorted.slice(0, findingLimit),
189 summary: {
190 recordType: 'quality-summary',
191 status: sorted.some((finding) => finding.severity === 'ERROR') ? 'FAILED' : 'PASSED',
192 batchId: batchId.trim(),
193 reviewCount: reviews.length,
194 findingCount: sorted.length,
195 returnedFindings: Math.min(sorted.length, findingLimit),
196 findingsTruncated: sorted.length > findingLimit,
197 errorCount: sorted.filter((finding) => finding.severity === 'ERROR').length,
198 warningCount: sorted.filter((finding) => finding.severity === 'WARNING').length,
199 repeatedContentGroups: repeated.length,
200 duplicateIdGroups: duplicateIds.length,
201 inputHash,
202 engineVersion: ENGINE_VERSION,
203 },
204 };
205}