feat: complete dashcam catalog ingestion demo
Master Production Deploy / production (push) Failing after 49s

This commit is contained in:
2026-07-31 02:01:58 +03:00
parent ae6c35145f
commit 2ce58fcbc0
62 changed files with 6341 additions and 178 deletions
+28
View File
@@ -0,0 +1,28 @@
import { describe, expect, it } from 'vitest';
import { getAdminAuthConfig, isBasicAuthAuthorized } from './adminAuth.ts';
function basic(username: string, password: string) {
return `Basic ${Buffer.from(`${username}:${password}`).toString('base64')}`;
}
describe('ingestion admin auth', () => {
it('requires an admin password', () => {
expect(() => getAdminAuthConfig({ INGESTION_ADMIN_USERNAME: 'admin' })).toThrow(
'INGESTION_ADMIN_PASSWORD is required',
);
});
it('accepts matching basic auth credentials', () => {
const config = { password: 'secret', username: 'operator' };
expect(isBasicAuthAuthorized(basic('operator', 'secret'), config)).toBe(true);
});
it('rejects missing or wrong basic auth credentials', () => {
const config = { password: 'secret', username: 'operator' };
expect(isBasicAuthAuthorized(undefined, config)).toBe(false);
expect(isBasicAuthAuthorized(basic('operator', 'wrong'), config)).toBe(false);
expect(isBasicAuthAuthorized(basic('wrong', 'secret'), config)).toBe(false);
});
});
+45
View File
@@ -0,0 +1,45 @@
import { timingSafeEqual } from 'node:crypto';
export type AdminAuthConfig = {
password: string;
username: string;
};
export const adminAuthRealm = 'videoreg ingestion admin';
export function getAdminAuthConfig(env: Record<string, string | undefined> = process.env): AdminAuthConfig {
const username = env.INGESTION_ADMIN_USERNAME ?? 'admin';
const password = env.INGESTION_ADMIN_PASSWORD;
if (!password) {
throw new Error('INGESTION_ADMIN_PASSWORD is required for ingestion-admin');
}
return { password, username };
}
export function isBasicAuthAuthorized(authorizationHeader: string | undefined, config: AdminAuthConfig) {
if (!authorizationHeader?.startsWith('Basic ')) {
return false;
}
const encodedCredentials = authorizationHeader.slice('Basic '.length).trim();
const decodedCredentials = Buffer.from(encodedCredentials, 'base64').toString('utf8');
const separatorIndex = decodedCredentials.indexOf(':');
if (separatorIndex < 0) {
return false;
}
const username = decodedCredentials.slice(0, separatorIndex);
const password = decodedCredentials.slice(separatorIndex + 1);
return safeEqual(username, config.username) && safeEqual(password, config.password);
}
function safeEqual(left: string, right: string) {
const leftBuffer = Buffer.from(left);
const rightBuffer = Buffer.from(right);
return leftBuffer.length === rightBuffer.length && timingSafeEqual(leftBuffer, rightBuffer);
}
+479
View File
@@ -0,0 +1,479 @@
import http from 'node:http';
import {
adminAuthRealm,
getAdminAuthConfig,
isBasicAuthAuthorized,
type AdminAuthConfig,
} from './adminAuth.ts';
import { adminUiCss, adminUiHtml, adminUiJs } from './adminUi.ts';
import { createPgPool } from './db.ts';
function sendJson(response: http.ServerResponse, status: number, payload: unknown) {
response.writeHead(status, { 'content-type': 'application/json; charset=utf-8' });
response.end(`${JSON.stringify(payload, null, 2)}\n`);
}
function sendText(response: http.ServerResponse, status: number, contentType: string, payload: string) {
response.writeHead(status, {
'cache-control': 'no-store',
'content-type': contentType,
});
response.end(payload);
}
function sendBasicAuthChallenge(response: http.ServerResponse) {
response.writeHead(401, {
'cache-control': 'no-store',
'content-type': 'text/plain; charset=utf-8',
'www-authenticate': `Basic realm="${adminAuthRealm}", charset="UTF-8"`,
});
response.end('Authentication required\n');
}
async function readJsonBody(request: http.IncomingMessage) {
const chunks: Buffer[] = [];
for await (const chunk of request) {
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
}
if (chunks.length === 0) return {};
return JSON.parse(Buffer.concat(chunks).toString('utf8')) as Record<string, unknown>;
}
function isTrustStatus(value: unknown): value is 'trusted' | 'unknown' | 'untrusted' {
return value === 'trusted' || value === 'unknown' || value === 'untrusted';
}
function isModerationStatus(value: unknown): value is 'needs_review' | 'approved' | 'rejected' | 'duplicate' | 'variant' | 'oem_clone' {
return ['needs_review', 'approved', 'rejected', 'duplicate', 'variant', 'oem_clone'].includes(String(value));
}
function isCanonicalText(value: unknown): value is string {
return typeof value === 'string' && value.trim().length > 0 && value.trim().length <= 300;
}
function isProtectedAdminPath(pathname: string) {
return (
pathname === '/admin' ||
pathname.startsWith('/admin/') ||
pathname === '/queue' ||
pathname === '/moderation' ||
pathname === '/manufacturers' ||
pathname.startsWith('/manufacturers/')
);
}
function adminApiPath(pathname: string) {
if (pathname === '/admin/api') return '/';
if (pathname.startsWith('/admin/api/')) return pathname.slice('/admin/api'.length);
return pathname;
}
type CreateAdminServerOptions = {
auth?: AdminAuthConfig;
};
export function createAdminServer(options: CreateAdminServerOptions = {}) {
const pool = createPgPool();
const auth = options.auth ?? getAdminAuthConfig();
return http.createServer(async (request, response) => {
const url = new URL(request.url ?? '/', 'http://localhost');
try {
if (request.method === 'GET' && url.pathname === '/') {
response.writeHead(302, { location: '/admin' });
response.end();
return;
}
if (isProtectedAdminPath(url.pathname) && !isBasicAuthAuthorized(request.headers.authorization, auth)) {
sendBasicAuthChallenge(response);
return;
}
if (request.method === 'GET' && url.pathname === '/admin') {
sendText(response, 200, 'text/html; charset=utf-8', adminUiHtml);
return;
}
if (request.method === 'GET' && url.pathname === '/admin/styles.css') {
sendText(response, 200, 'text/css; charset=utf-8', adminUiCss);
return;
}
if (request.method === 'GET' && url.pathname === '/admin/app.js') {
sendText(response, 200, 'text/javascript; charset=utf-8', adminUiJs);
return;
}
const apiPath = adminApiPath(url.pathname);
if (apiPath === '/health') {
await pool.query('select 1');
sendJson(response, 200, { status: 'ok' });
return;
}
if (apiPath === '/queue') {
const result = await pool.query(
`
select state, job_type as "jobType", count(*)::int as count
from crawl_jobs
group by state, job_type
order by state, job_type
`,
);
sendJson(response, 200, { jobs: result.rows });
return;
}
if (apiPath === '/moderation') {
const result = await pool.query(
`
select state, reason, count(*)::int as count
from moderation_queue
group by state, reason
order by state, reason
`,
);
sendJson(response, 200, { queue: result.rows });
return;
}
if (apiPath === '/moderation/products' && request.method === 'GET') {
const requestedStatus = url.searchParams.get('status') ?? 'needs_review';
if (!isModerationStatus(requestedStatus)) {
sendJson(response, 400, { error: 'invalid moderation status' });
return;
}
const result = await pool.query(
`
select
products.id,
products.brand,
products.model,
products.canonical_name as "canonicalName",
products.title,
products.moderation_status as "moderationStatus",
products.updated_at as "updatedAt",
count(distinct product_sources.id)::int as "sourceCount",
count(distinct product_source_specs.id)::int as "rawSpecCount",
count(distinct moderation_conflicts.id) filter (where moderation_conflicts.state = 'pending')::int as "pendingConflictCount"
from products
left join product_sources on product_sources.product_id = products.id
left join product_source_specs on product_source_specs.product_id = products.id
left join moderation_conflicts on moderation_conflicts.product_id = products.id
where products.moderation_status = $1
group by products.id
order by products.updated_at desc, products.brand asc, products.model asc
`,
[requestedStatus],
);
sendJson(response, 200, { products: result.rows });
return;
}
const moderationProductMatch = apiPath.match(/^\/moderation\/products\/([^/]+)$/);
if (moderationProductMatch && request.method === 'GET') {
const productId = decodeURIComponent(moderationProductMatch[1]);
const productResult = await pool.query(
`
select
id,
brand,
model as "canonicalModel",
canonical_name as "canonicalName",
title,
moderation_status as "moderationStatus",
model_override as "canonicalModelOverride",
canonical_name_override as "canonicalNameOverride",
updated_at as "updatedAt"
from products
where id = $1
`,
[productId],
);
if (!productResult.rows[0]) {
sendJson(response, 404, { error: 'product_not_found' });
return;
}
const [sources, specs, conflicts] = await Promise.all([
pool.query(
`
select
source,
url,
source_product_id as "sourceProductId",
stable_content_hash as "stableContentHash",
last_fetched_at as "lastFetchedAt"
from product_sources
where product_id = $1
order by source, url
`,
[productId],
),
pool.query(
`
select
source,
source_url as "sourceUrl",
raw_group as "group",
raw_label as "label",
raw_value as "value",
normalized_key as "normalizedKey",
normalized_value as "normalizedValue"
from product_source_specs
where product_id = $1
order by source, raw_group nulls first, raw_label
`,
[productId],
),
pool.query(
`
select
conflict_type as "conflictType",
key,
values_json as "values",
state
from moderation_conflicts
where product_id = $1
order by state, key
`,
[productId],
),
]);
sendJson(response, 200, {
product: {
...productResult.rows[0],
conflicts: conflicts.rows,
rawSpecs: specs.rows,
sources: sources.rows,
},
});
return;
}
if (moderationProductMatch && request.method === 'PATCH') {
const body = await readJsonBody(request);
const hasModerationStatus = Object.hasOwn(body, 'moderationStatus');
const hasCanonicalModel = Object.hasOwn(body, 'canonicalModel');
const hasCanonicalName = Object.hasOwn(body, 'canonicalName');
if (hasModerationStatus && !isModerationStatus(body.moderationStatus)) {
sendJson(response, 400, { error: 'invalid moderationStatus' });
return;
}
if (hasCanonicalModel && !isCanonicalText(body.canonicalModel)) {
sendJson(response, 400, { error: 'canonicalModel must be a non-empty string up to 300 characters' });
return;
}
if (hasCanonicalName && !isCanonicalText(body.canonicalName)) {
sendJson(response, 400, { error: 'canonicalName must be a non-empty string up to 300 characters' });
return;
}
if (!hasModerationStatus && !hasCanonicalModel && !hasCanonicalName) {
sendJson(response, 400, { error: 'at least one product field is required' });
return;
}
const productId = decodeURIComponent(moderationProductMatch[1]);
const result = await pool.query(
`
update products
set moderation_status = case when $2::boolean then $3::text else moderation_status end,
model = case when $4::boolean then $5::text else model end,
canonical_name = case when $6::boolean then $7::text else canonical_name end,
model_override = case when $4::boolean then true else model_override end,
canonical_name_override = case when $6::boolean then true else canonical_name_override end,
updated_at = now()
where id = $1
returning
id,
model as "canonicalModel",
canonical_name as "canonicalName",
moderation_status as "moderationStatus",
model_override as "canonicalModelOverride",
canonical_name_override as "canonicalNameOverride",
updated_at as "updatedAt"
`,
[
productId,
hasModerationStatus,
hasModerationStatus ? body.moderationStatus : null,
hasCanonicalModel,
hasCanonicalModel ? String(body.canonicalModel).trim() : null,
hasCanonicalName,
hasCanonicalName ? String(body.canonicalName).trim() : null,
],
);
if (!result.rows[0]) {
sendJson(response, 404, { error: 'product_not_found' });
return;
}
if (hasModerationStatus) {
await pool.query(
`
update moderation_queue
set state = case when $2 = 'needs_review' then 'pending' else 'resolved' end,
updated_at = now()
where product_id = $1
`,
[productId, body.moderationStatus],
);
if (body.moderationStatus !== 'needs_review') {
await pool.query(
`
update crawl_jobs job
set state = $2,
updated_at = now()
from product_sources source
where source.product_id = $1
and job.source = source.source
and job.url = source.url
and job.job_type = 'product'
and job.state = 'needs_review'
`,
[productId, body.moderationStatus === 'approved' ? 'approved' : 'needs_manual_review'],
);
}
}
sendJson(response, 200, { product: result.rows[0] });
return;
}
if (apiPath === '/manufacturers') {
const result = await pool.query(
`
select
manufacturers.slug,
manufacturers.name,
manufacturers.trust_status as "trustStatus",
manufacturers.notes,
count(products.id)::int as "productCount",
max(products.updated_at) as "lastProductUpdatedAt"
from manufacturers
left join products on products.manufacturer_id = manufacturers.id
group by manufacturers.id
order by manufacturers.trust_status desc, manufacturers.name asc
`,
);
sendJson(response, 200, { manufacturers: result.rows });
return;
}
const manufacturerMatch = apiPath.match(/^\/manufacturers\/([^/]+)$/);
if (manufacturerMatch && request.method === 'GET') {
const result = await pool.query(
`
select
manufacturers.slug,
manufacturers.name,
manufacturers.trust_status as "trustStatus",
manufacturers.notes,
coalesce(
jsonb_agg(
jsonb_build_object(
'id', products.id,
'canonicalName', products.canonical_name,
'title', products.title,
'moderationStatus', products.moderation_status,
'updatedAt', products.updated_at
)
order by products.updated_at desc
) filter (where products.id is not null),
'[]'::jsonb
) as products
from manufacturers
left join products on products.manufacturer_id = manufacturers.id
where manufacturers.slug = $1
group by manufacturers.id
`,
[decodeURIComponent(manufacturerMatch[1])],
);
if (!result.rows[0]) {
sendJson(response, 404, { error: 'manufacturer_not_found' });
return;
}
sendJson(response, 200, { manufacturer: result.rows[0] });
return;
}
if (manufacturerMatch && request.method === 'PATCH') {
const body = await readJsonBody(request);
if (!isTrustStatus(body.trustStatus)) {
sendJson(response, 400, { error: 'trustStatus must be trusted, unknown, or untrusted' });
return;
}
const hasNotes = Object.hasOwn(body, 'notes');
if (hasNotes && body.notes !== null && typeof body.notes !== 'string') {
sendJson(response, 400, { error: 'notes must be a string or null' });
return;
}
const result = await pool.query(
`
update manufacturers
set trust_status = $2,
notes = case when $3 then $4 else notes end,
updated_at = now()
where slug = $1
returning slug, name, trust_status as "trustStatus", notes
`,
[
decodeURIComponent(manufacturerMatch[1]),
body.trustStatus,
hasNotes,
hasNotes ? body.notes : null,
],
);
if (!result.rows[0]) {
sendJson(response, 404, { error: 'manufacturer_not_found' });
return;
}
sendJson(response, 200, { manufacturer: result.rows[0] });
return;
}
sendJson(response, 404, { error: 'not_found' });
} catch (error) {
sendJson(response, 500, { error: (error as Error).message });
}
});
}
if (import.meta.url === `file://${process.argv[1]}`) {
const port = Number(process.env.INGESTION_ADMIN_PORT ?? 4101);
createAdminServer().listen(port, () => {
console.log(`ingestion-admin listening on ${port}`);
});
}
+17
View File
@@ -0,0 +1,17 @@
import { describe, expect, it } from 'vitest';
import { adminUiHtml, adminUiJs } from './adminUi.ts';
describe('ingestion admin UI', () => {
it('is served as a standalone ingestion app', () => {
expect(adminUiHtml).toContain('/admin/styles.css');
expect(adminUiHtml).toContain('/admin/app.js');
expect(adminUiHtml).not.toContain('/api/admin');
expect(adminUiJs).toContain("const adminApiBasePath = '/admin/api'");
expect(adminUiJs).toContain('data-product-form');
expect(adminUiJs).not.toContain("requestJson('/manufacturers");
});
it('ships valid browser JavaScript', () => {
expect(() => new Function(adminUiJs)).not.toThrow();
});
});
+790
View File
@@ -0,0 +1,790 @@
export const adminUiHtml = `<!doctype html>
<html lang="en">
<head>
<meta charset="utf-8" />
<meta name="viewport" content="width=device-width, initial-scale=1" />
<title>Videoreg ingestion admin</title>
<link rel="stylesheet" href="/admin/styles.css" />
</head>
<body>
<main class="shell">
<header class="topbar">
<div>
<p class="eyebrow">catalog moderation</p>
<h1>Ingestion admin</h1>
</div>
<nav class="actions" aria-label="Admin navigation">
<a class="button ghost" href="/admin/api/health">Health</a>
<a class="button ghost" href="/admin/api/queue">Queue</a>
<a class="button ghost" href="/admin/api/moderation">Moderation</a>
<a class="button" href="#/moderation">Products</a>
<a class="button" href="#/">Manufacturers</a>
</nav>
</header>
<div id="app" class="app" aria-live="polite"></div>
</main>
<script src="/admin/app.js" type="module"></script>
</body>
</html>
`;
export const adminUiCss = `
:root {
color-scheme: light;
--surface-0: #f3f0e8;
--surface-1: #fffdf8;
--surface-2: #e9e4d8;
--text-primary: #202528;
--text-muted: #667076;
--text-subtle: #858076;
--border-subtle: #d7cec0;
--border-strong: #bdb3a4;
--accent: #56798a;
--accent-strong: #385f70;
--accent-soft: #dbe7e9;
--warning: #8f7044;
--warning-soft: #efe5d6;
--danger: #884b46;
--danger-soft: #eadbd8;
--success: #58705c;
--success-soft: #dde8dd;
--radius-sm: 4px;
--radius-md: 6px;
--radius-lg: 8px;
--shadow-soft: 0 14px 34px rgb(32 37 40 / 8%);
}
* {
box-sizing: border-box;
}
body {
margin: 0;
background: var(--surface-0);
color: var(--text-primary);
font-family: Inter, ui-sans-serif, system-ui, -apple-system, BlinkMacSystemFont, "Segoe UI", sans-serif;
}
a {
color: inherit;
text-decoration: none;
}
button,
input,
select,
textarea {
font: inherit;
}
.shell {
width: min(1180px, calc(100% - 32px));
margin: 0 auto;
padding: 24px 0 40px;
}
.topbar {
display: flex;
align-items: flex-start;
justify-content: space-between;
gap: 24px;
border-bottom: 1px solid var(--border-subtle);
padding-bottom: 18px;
}
.eyebrow {
margin: 0 0 6px;
color: var(--accent-strong);
font-size: 0.78rem;
font-weight: 750;
text-transform: uppercase;
}
h1,
h2,
h3,
p {
overflow-wrap: anywhere;
}
h1 {
margin: 0;
font-size: 2rem;
line-height: 1.1;
}
h2 {
margin: 0;
font-size: 1.1rem;
}
.actions {
display: flex;
flex-wrap: wrap;
gap: 8px;
justify-content: flex-end;
}
.button {
min-height: 38px;
display: inline-flex;
align-items: center;
justify-content: center;
border: 1px solid var(--border-strong);
border-radius: var(--radius-md);
background: var(--surface-1);
color: var(--text-primary);
cursor: pointer;
padding: 8px 12px;
font-size: 0.86rem;
font-weight: 700;
}
.button.primary {
border-color: var(--accent);
background: var(--accent);
color: #fff;
}
.button.ghost {
border-color: transparent;
background: transparent;
color: var(--text-muted);
}
.button:hover {
background: var(--surface-2);
}
.button.primary:hover {
background: var(--accent-strong);
}
.app {
display: grid;
gap: 20px;
padding-top: 22px;
}
.panel {
border: 1px solid var(--border-subtle);
border-radius: var(--radius-lg);
background: var(--surface-1);
box-shadow: var(--shadow-soft);
}
.panel-header {
display: flex;
align-items: flex-start;
justify-content: space-between;
gap: 16px;
border-bottom: 1px solid var(--border-subtle);
padding: 16px;
}
.panel-body {
padding: 16px;
}
.stats-grid {
display: grid;
grid-template-columns: repeat(4, minmax(120px, 1fr));
gap: 10px;
}
.metric {
min-height: 76px;
border: 1px solid var(--border-subtle);
border-radius: var(--radius-lg);
background: var(--surface-2);
padding: 12px;
}
.metric span {
display: block;
color: var(--text-muted);
font-size: 0.76rem;
font-weight: 750;
text-transform: uppercase;
}
.metric strong {
display: block;
margin-top: 7px;
font-size: 1.45rem;
}
.table-wrap {
overflow-x: auto;
}
table {
min-width: 760px;
width: 100%;
border-collapse: collapse;
text-align: left;
font-size: 0.9rem;
}
th,
td {
border-bottom: 1px solid var(--border-subtle);
padding: 14px 16px;
vertical-align: top;
}
thead th {
background: var(--surface-1);
color: var(--text-muted);
font-size: 0.74rem;
font-weight: 750;
text-transform: uppercase;
}
tbody tr:hover {
background: color-mix(in srgb, var(--surface-2) 48%, transparent);
}
tbody tr:last-child th,
tbody tr:last-child td {
border-bottom: 0;
}
.badge {
display: inline-flex;
min-height: 24px;
align-items: center;
border-radius: var(--radius-sm);
border: 1px solid var(--border-subtle);
padding: 4px 8px;
font-size: 0.72rem;
font-weight: 750;
}
.badge.trusted {
background: var(--success-soft);
color: var(--success);
}
.badge.unknown {
background: var(--warning-soft);
color: var(--warning);
}
.badge.untrusted {
background: var(--danger-soft);
color: var(--danger);
}
.muted {
color: var(--text-muted);
}
.mono {
font-family: "IBM Plex Mono", ui-monospace, SFMono-Regular, Menlo, Monaco, Consolas, monospace;
}
.layout {
display: grid;
grid-template-columns: minmax(280px, 380px) minmax(0, 1fr);
gap: 20px;
align-items: start;
}
.form-grid {
display: grid;
gap: 14px;
}
.field {
display: grid;
gap: 6px;
}
.field span {
color: var(--text-muted);
font-size: 0.82rem;
font-weight: 700;
}
.field select,
.field textarea {
width: 100%;
border: 1px solid var(--border-subtle);
border-radius: var(--radius-md);
background: var(--surface-1);
color: var(--text-primary);
padding: 9px 10px;
}
.field select {
min-height: 42px;
}
.field textarea {
min-height: 120px;
resize: vertical;
}
.status-line {
min-height: 24px;
color: var(--text-muted);
font-size: 0.88rem;
font-weight: 700;
}
.error,
.empty {
border: 1px dashed var(--border-subtle);
border-radius: var(--radius-lg);
color: var(--text-muted);
padding: 22px;
}
.error strong {
display: block;
color: var(--danger);
margin-bottom: 6px;
}
@media (max-width: 900px) {
.topbar,
.layout {
display: grid;
grid-template-columns: 1fr;
}
.actions {
justify-content: flex-start;
}
.stats-grid {
grid-template-columns: repeat(2, minmax(0, 1fr));
}
}
@media (max-width: 560px) {
.shell {
width: min(100% - 20px, 1180px);
padding-top: 16px;
}
h1 {
font-size: 1.55rem;
}
.stats-grid {
grid-template-columns: 1fr;
}
}
`;
export const adminUiJs = `
const app = document.querySelector('#app');
const adminApiBasePath = '/admin/api';
const dateFormatter = new Intl.DateTimeFormat('en', {
dateStyle: 'medium',
timeStyle: 'short',
});
function escapeHtml(value) {
return String(value ?? '')
.replaceAll('&', '&amp;')
.replaceAll('<', '&lt;')
.replaceAll('>', '&gt;')
.replaceAll('"', '&quot;')
.replaceAll("'", '&#039;');
}
function formatDate(value) {
if (!value) return 'Never';
return dateFormatter.format(new Date(value));
}
function formatValue(value) {
if (value === undefined || value === null) return '—';
if (typeof value === 'string') return value;
return JSON.stringify(value);
}
function statusLabel(status) {
if (status === 'trusted') return 'Trusted';
if (status === 'untrusted') return 'Untrusted';
if (status === 'needs_review') return 'Needs review';
if (status === 'approved') return 'Approved';
if (status === 'rejected') return 'Rejected';
if (status === 'duplicate') return 'Duplicate';
if (status === 'variant') return 'Variant';
if (status === 'oem_clone') return 'OEM clone';
return 'Unknown';
}
async function requestJson(path, options) {
const response = await fetch(path, {
credentials: 'same-origin',
headers: {
accept: 'application/json',
...(options?.headers ?? {}),
},
...options,
});
const payload = await response.json().catch(() => undefined);
if (!response.ok) {
throw new Error(payload?.error ?? response.status + ' ' + response.statusText);
}
return payload;
}
function renderLoading(label) {
app.innerHTML = '<section class="panel"><div class="panel-body muted">' + escapeHtml(label) + '</div></section>';
}
function renderError(title, error) {
app.innerHTML =
'<section class="panel"><div class="panel-body"><div class="error">' +
'<strong>' + escapeHtml(title) + '</strong>' +
'<span>' + escapeHtml(error.message ?? error) + '</span>' +
'</div></div></section>';
}
function renderMetric(label, value) {
return '<div class="metric"><span>' + escapeHtml(label) + '</span><strong>' + escapeHtml(value) + '</strong></div>';
}
function renderBadge(status) {
return '<span class="badge ' + escapeHtml(status) + '">' + escapeHtml(statusLabel(status)) + '</span>';
}
function renderManufacturers(manufacturers) {
const stats = {
total: manufacturers.length,
trusted: manufacturers.filter((item) => item.trustStatus === 'trusted').length,
unknown: manufacturers.filter((item) => item.trustStatus === 'unknown').length,
untrusted: manufacturers.filter((item) => item.trustStatus === 'untrusted').length,
};
const rows = manufacturers
.map(
(manufacturer) =>
'<tr>' +
'<th scope="row"><div><strong>' + escapeHtml(manufacturer.name) + '</strong><br />' +
'<span class="mono muted">' + escapeHtml(manufacturer.slug) + '</span></div></th>' +
'<td>' + renderBadge(manufacturer.trustStatus) + '</td>' +
'<td class="mono">' + escapeHtml(manufacturer.productCount) + '</td>' +
'<td>' + escapeHtml(formatDate(manufacturer.lastProductUpdatedAt)) + '</td>' +
'<td class="muted">' + escapeHtml(manufacturer.notes || 'No notes') + '</td>' +
'<td><a class="button" href="#/manufacturers/' + encodeURIComponent(manufacturer.slug) + '">Edit</a></td>' +
'</tr>',
)
.join('');
app.innerHTML =
'<section class="stats-grid" aria-label="Manufacturer moderation stats">' +
renderMetric('Total', stats.total) +
renderMetric('Trusted', stats.trusted) +
renderMetric('Unknown', stats.unknown) +
renderMetric('Untrusted', stats.untrusted) +
'</section>' +
'<section class="panel"><div class="panel-header"><h2>Manufacturers</h2></div>' +
(manufacturers.length
? '<div class="table-wrap"><table><thead><tr>' +
'<th scope="col">Brand</th><th scope="col">Trust</th><th scope="col">Products</th>' +
'<th scope="col">Updated</th><th scope="col">Notes</th><th scope="col">Actions</th>' +
'</tr></thead><tbody>' + rows + '</tbody></table></div>'
: '<div class="panel-body"><div class="empty">No manufacturers discovered yet</div></div>') +
'</section>';
}
function renderManufacturer(manufacturer) {
const products = manufacturer.products ?? [];
const productRows = products
.map(
(product) =>
'<tr>' +
'<th scope="row"><strong>' + escapeHtml(product.canonicalName || product.title) + '</strong><br />' +
'<span class="mono muted">' + escapeHtml(product.id) + '</span></th>' +
'<td>' + escapeHtml(product.moderationStatus) + '</td>' +
'<td>' + escapeHtml(formatDate(product.updatedAt)) + '</td>' +
'</tr>',
)
.join('');
app.innerHTML =
'<div class="layout">' +
'<section class="panel">' +
'<div class="panel-header"><div><h2>Moderation</h2><p class="muted">Mark unreliable brands so source conflicts are visible during moderation.</p></div>' +
renderBadge(manufacturer.trustStatus) + '</div>' +
'<div class="panel-body">' +
'<form class="form-grid" data-manufacturer-form data-slug="' + escapeHtml(manufacturer.slug) + '">' +
'<label class="field"><span>Trust status</span><select name="trustStatus">' +
'<option value="trusted"' + (manufacturer.trustStatus === 'trusted' ? ' selected' : '') + '>Trusted</option>' +
'<option value="unknown"' + (manufacturer.trustStatus === 'unknown' ? ' selected' : '') + '>Unknown</option>' +
'<option value="untrusted"' + (manufacturer.trustStatus === 'untrusted' ? ' selected' : '') + '>Untrusted</option>' +
'</select></label>' +
'<label class="field"><span>Notes</span><textarea name="notes" rows="6" placeholder="Why this manufacturer is trusted or not trusted">' +
escapeHtml(manufacturer.notes ?? '') + '</textarea></label>' +
'<div class="actions"><button class="button primary" type="submit">Save</button>' +
'<a class="button ghost" href="#/">Back</a></div>' +
'<div class="status-line" data-save-state></div>' +
'</form></div></section>' +
'<section class="panel"><div class="panel-header"><div><h2>Products</h2><p class="muted">Products currently linked to this manufacturer.</p></div>' +
'<span class="badge unknown">' + escapeHtml(products.length) + '</span></div>' +
(products.length
? '<div class="table-wrap"><table><thead><tr><th scope="col">Product</th><th scope="col">Moderation</th><th scope="col">Updated</th></tr></thead><tbody>' +
productRows + '</tbody></table></div>'
: '<div class="panel-body"><div class="empty">No products linked to this manufacturer</div></div>') +
'</section></div>';
}
function renderModerationProducts(products) {
const rows = products
.map(
(product) =>
'<tr>' +
'<th scope="row"><strong>' + escapeHtml(product.canonicalName || product.title) + '</strong><br />' +
'<span class="muted">' + escapeHtml(product.title) + '</span></th>' +
'<td>' + renderBadge(product.moderationStatus) + '</td>' +
'<td class="mono">' + escapeHtml(product.sourceCount) + '</td>' +
'<td class="mono">' + escapeHtml(product.rawSpecCount) + '</td>' +
'<td class="mono">' + escapeHtml(product.pendingConflictCount) + '</td>' +
'<td>' +
'<div class="actions">' +
'<a class="button" href="#/moderation/' + encodeURIComponent(product.id) + '">Edit</a>' +
'<button class="button primary" type="button" data-product-status="approved" data-product-id="' + escapeHtml(product.id) + '">Approve</button>' +
'<button class="button" type="button" data-product-status="rejected" data-product-id="' + escapeHtml(product.id) + '">Reject</button>' +
'</div>' +
'</td>' +
'</tr>',
)
.join('');
app.innerHTML =
'<section class="panel"><div class="panel-header"><div><h2>Product moderation</h2><p class="muted">Approve source cards after checking title, canonical model and source conflicts.</p></div>' +
'<span class="badge unknown">' + escapeHtml(products.length) + '</span></div>' +
(products.length
? '<div class="table-wrap"><table><thead><tr>' +
'<th scope="col">Product</th><th scope="col">Status</th><th scope="col">Sources</th>' +
'<th scope="col">Raw specs</th><th scope="col">Conflicts</th><th scope="col">Actions</th>' +
'</tr></thead><tbody>' + rows + '</tbody></table></div>'
: '<div class="panel-body"><div class="empty">No products need moderation</div></div>') +
'</section>';
}
function renderProductEditor(product) {
const sourceRows = (product.sources ?? [])
.map(
(source) =>
'<tr><td>' + escapeHtml(source.source) + '</td><td class="mono">' +
escapeHtml(source.sourceProductId ?? '—') + '</td><td><a class="source-link" href="' +
escapeHtml(source.url) + '" target="_blank" rel="noreferrer">Source card</a></td><td>' +
escapeHtml(formatDate(source.lastFetchedAt)) + '</td></tr>',
)
.join('');
const rawSpecRows = (product.rawSpecs ?? [])
.map(
(spec) =>
'<tr><td>' + escapeHtml(spec.source) + '</td><td>' + escapeHtml(spec.group ?? '—') + '</td><th scope="row">' +
escapeHtml(spec.label) + '</th><td>' + escapeHtml(formatValue(spec.value)) + '</td><td class="mono">' +
escapeHtml(spec.normalizedKey ?? '—') + '</td><td>' + escapeHtml(formatValue(spec.normalizedValue)) + '</td></tr>',
)
.join('');
const conflictRows = (product.conflicts ?? [])
.map(
(conflict) =>
'<tr><th scope="row">' + escapeHtml(conflict.key) + '</th><td>' + renderBadge(conflict.state) + '</td><td class="mono">' +
escapeHtml(formatValue(conflict.values)) + '</td></tr>',
)
.join('');
app.innerHTML =
'<div class="layout">' +
'<section class="panel"><div class="panel-header"><div><h2>Product editor</h2><p class="muted">' +
escapeHtml(product.title) + '</p></div>' + renderBadge(product.moderationStatus) + '</div><div class="panel-body">' +
'<form class="form-grid" data-product-form data-product-id="' + escapeHtml(product.id) + '">' +
'<label class="field"><span>Canonical model</span><input name="canonicalModel" value="' +
escapeHtml(product.canonicalModel) + '" required maxlength="300" /></label>' +
'<label class="field"><span>Canonical name</span><input name="canonicalName" value="' +
escapeHtml(product.canonicalName) + '" required maxlength="300" /></label>' +
'<label class="field"><span>Moderation status</span><select name="moderationStatus">' +
['needs_review', 'approved', 'rejected', 'duplicate', 'variant', 'oem_clone']
.map((status) => '<option value="' + status + '"' + (status === product.moderationStatus ? ' selected' : '') + '>' + escapeHtml(statusLabel(status)) + '</option>')
.join('') +
'</select></label><div class="actions"><button class="button primary" type="submit">Save</button>' +
'<a class="button ghost" href="#/moderation">Back</a></div><div class="status-line" data-save-state></div></form></div></section>' +
'<section class="panel"><div class="panel-header"><div><h2>Sources</h2></div><span class="badge unknown">' +
escapeHtml((product.sources ?? []).length) + '</span></div>' +
(sourceRows ? '<div class="table-wrap"><table><thead><tr><th>Source</th><th>ID</th><th>Card</th><th>Fetched</th></tr></thead><tbody>' + sourceRows + '</tbody></table></div>' : '<div class="panel-body"><div class="empty">No source records</div></div>') +
'</section></div>' +
'<section class="panel"><div class="panel-header"><div><h2>Raw specifications</h2></div><span class="badge unknown">' +
escapeHtml((product.rawSpecs ?? []).length) + '</span></div>' +
(rawSpecRows ? '<div class="table-wrap"><table><thead><tr><th>Source</th><th>Group</th><th>Label</th><th>Raw value</th><th>Normalized key</th><th>Normalized value</th></tr></thead><tbody>' + rawSpecRows + '</tbody></table></div>' : '<div class="panel-body"><div class="empty">No raw specifications</div></div>') +
'</section>' +
'<section class="panel"><div class="panel-header"><div><h2>Source conflicts</h2></div><span class="badge unknown">' +
escapeHtml((product.conflicts ?? []).length) + '</span></div>' +
(conflictRows ? '<div class="table-wrap"><table><thead><tr><th>Key</th><th>Status</th><th>Values</th></tr></thead><tbody>' + conflictRows + '</tbody></table></div>' : '<div class="panel-body"><div class="empty">No source conflicts</div></div>') +
'</section>';
}
async function loadManufacturers() {
renderLoading('Loading manufacturers');
try {
const payload = await requestJson(adminApiBasePath + '/manufacturers');
renderManufacturers(payload.manufacturers ?? []);
} catch (error) {
renderError('Manufacturers could not be loaded', error);
}
}
async function loadManufacturer(slug) {
renderLoading('Loading manufacturer');
try {
const payload = await requestJson(adminApiBasePath + '/manufacturers/' + encodeURIComponent(slug));
renderManufacturer(payload.manufacturer);
} catch (error) {
renderError('Manufacturer could not be loaded', error);
}
}
async function loadModerationProducts() {
renderLoading('Loading product moderation');
try {
const payload = await requestJson(adminApiBasePath + '/moderation/products?status=needs_review');
renderModerationProducts(payload.products ?? []);
} catch (error) {
renderError('Product moderation could not be loaded', error);
}
}
async function loadModerationProduct(productId) {
renderLoading('Loading product');
try {
const payload = await requestJson(adminApiBasePath + '/moderation/products/' + encodeURIComponent(productId));
renderProductEditor(payload.product);
} catch (error) {
renderError('Product could not be loaded', error);
}
}
function route() {
const hash = window.location.hash || '#/';
const match = hash.match(/^#\\/manufacturers\\/(.+)$/);
if (match) {
loadManufacturer(decodeURIComponent(match[1]));
return;
}
const productMatch = hash.match(/^#\\/moderation\\/([^/]+)$/);
if (productMatch) {
loadModerationProduct(decodeURIComponent(productMatch[1]));
return;
}
if (hash === '#/moderation') {
loadModerationProducts();
return;
}
loadManufacturers();
}
app.addEventListener('submit', async (event) => {
const form = event.target.closest('[data-manufacturer-form]');
if (!form) return;
event.preventDefault();
const state = form.querySelector('[data-save-state]');
const slug = form.dataset.slug;
const formData = new FormData(form);
const trustStatus = formData.get('trustStatus');
const notesValue = String(formData.get('notes') ?? '').trim();
state.textContent = 'Saving';
try {
await requestJson(adminApiBasePath + '/manufacturers/' + encodeURIComponent(slug), {
method: 'PATCH',
headers: {
'content-type': 'application/json',
},
body: JSON.stringify({
trustStatus,
notes: notesValue || null,
}),
});
state.textContent = 'Saved';
await loadManufacturer(slug);
} catch (error) {
state.textContent = error.message ?? String(error);
}
});
app.addEventListener('click', async (event) => {
const button = event.target.closest('[data-product-status]');
if (!button) return;
const productId = button.dataset.productId;
const moderationStatus = button.dataset.productStatus;
if (!productId || !moderationStatus) return;
button.disabled = true;
try {
await requestJson(adminApiBasePath + '/moderation/products/' + encodeURIComponent(productId), {
method: 'PATCH',
headers: {
'content-type': 'application/json',
},
body: JSON.stringify({ moderationStatus }),
});
await loadModerationProducts();
} catch (error) {
button.disabled = false;
window.alert(error.message ?? String(error));
}
});
app.addEventListener('submit', async (event) => {
const form = event.target.closest('[data-product-form]');
if (!form) return;
event.preventDefault();
const state = form.querySelector('[data-save-state]');
const productId = form.dataset.productId;
const formData = new FormData(form);
if (!productId) return;
state.textContent = 'Saving';
try {
await requestJson(adminApiBasePath + '/moderation/products/' + encodeURIComponent(productId), {
method: 'PATCH',
headers: {
'content-type': 'application/json',
},
body: JSON.stringify({
canonicalModel: String(formData.get('canonicalModel') ?? '').trim(),
canonicalName: String(formData.get('canonicalName') ?? '').trim(),
moderationStatus: formData.get('moderationStatus'),
}),
});
state.textContent = 'Saved';
await loadModerationProduct(productId);
} catch (error) {
state.textContent = error.message ?? String(error);
}
});
window.addEventListener('hashchange', route);
route();
`;
+34
View File
@@ -0,0 +1,34 @@
import { describe, expect, it } from 'vitest';
import { buildSourceSeeds, yandexMarketDashcamCategoryUrl } from './config.ts';
describe('ingestion config', () => {
it('always includes implemented Yandex Market source seed', () => {
expect(buildSourceSeeds({})).toEqual([
{
categoryUrl: yandexMarketDashcamCategoryUrl,
implemented: true,
source: 'yandex_market',
},
]);
});
it('enables an explicitly configured MVideo source seed in the worker queue', () => {
expect(
buildSourceSeeds({
MVIDEO_DASHCAM_CATEGORY_URL: 'https://www.mvideo.ru/avtomobilnye-videoregistratory',
YANDEX_MARKET_DASHCAM_CATEGORY_URL: 'https://market.yandex.ru/custom',
}),
).toEqual([
{
categoryUrl: 'https://market.yandex.ru/custom',
implemented: true,
source: 'yandex_market',
},
{
categoryUrl: 'https://www.mvideo.ru/avtomobilnye-videoregistratory',
implemented: true,
source: 'mvideo',
},
]);
});
});
+57
View File
@@ -0,0 +1,57 @@
import path from 'node:path';
import type { MarketSource } from './types.ts';
export const yandexMarketDashcamCategoryUrl =
'https://market.yandex.ru/catalog--avtomobilnye-videoregistratory/82798275/list?hid=82798269';
export const mvideoDashcamCategoryUrl =
'https://www.mvideo.ru/product-list-page?q=%D0%B2%D0%B8%D0%B4%D0%B5%D0%BE%D1%80%D0%B5%D0%B3%D0%B8%D1%81%D1%82%D1%80%D0%B0%D1%82%D0%BE%D1%80%D1%8B';
export type SourceSeed = {
categoryUrl: string;
implemented: boolean;
source: MarketSource;
};
function numberEnv(name: string, fallback: number) {
const value = Number(process.env[name]);
return Number.isFinite(value) ? value : fallback;
}
export function buildSourceSeeds(env: Record<string, string | undefined> = process.env): SourceSeed[] {
const seeds: SourceSeed[] = [
{
categoryUrl: env.YANDEX_MARKET_DASHCAM_CATEGORY_URL ?? yandexMarketDashcamCategoryUrl,
implemented: true,
source: 'yandex_market',
},
];
if (env.MVIDEO_DASHCAM_CATEGORY_URL) {
seeds.push({
categoryUrl: env.MVIDEO_DASHCAM_CATEGORY_URL,
implemented: true,
source: 'mvideo',
});
}
return seeds;
}
const sourceSeeds = buildSourceSeeds();
export const ingestionConfig = {
databaseUrl: process.env.INGESTION_DATABASE_URL ?? 'postgres://videoreg:videoreg@localhost:5432/videoreg_ingestion',
snapshotDir: process.env.INGESTION_SNAPSHOT_DIR ?? path.join(process.cwd(), 'ingestion', 'data', 'snapshots'),
exportDir: process.env.INGESTION_EXPORT_DIR ?? path.join(process.cwd(), 'ingestion', 'data', 'exports'),
seedCategoryUrl: process.env.YANDEX_MARKET_DASHCAM_CATEGORY_URL ?? yandexMarketDashcamCategoryUrl,
sourceSeeds,
implementedSourceSeeds: sourceSeeds.filter((seed) => seed.implemented),
fetchDelayMinMs: numberEnv('INGESTION_FETCH_DELAY_MIN_MS', 5_000),
fetchDelayMaxMs: numberEnv('INGESTION_FETCH_DELAY_MAX_MS', 15_000),
categoryDailyBudget: numberEnv('INGESTION_CATEGORY_DAILY_BUDGET', 25),
productDailyBudget: numberEnv('INGESTION_PRODUCT_DAILY_BUDGET', 100),
sitemapDailyBudget: numberEnv('INGESTION_SITEMAP_DAILY_BUDGET', 25),
workerBatchSize: numberEnv('INGESTION_WORKER_BATCH_SIZE', 1),
workerIdleMs: numberEnv('INGESTION_WORKER_IDLE_MS', 30_000),
};
+10
View File
@@ -0,0 +1,10 @@
import pg from 'pg';
import { ingestionConfig } from './config.ts';
const { Pool } = pg;
export function createPgPool(connectionString = ingestionConfig.databaseUrl) {
return new Pool({
connectionString,
});
}
+12
View File
@@ -0,0 +1,12 @@
import { createPgPool } from './db.ts';
import { exportApprovedCatalog } from './exporter.ts';
const pool = createPgPool();
try {
const manifest = await exportApprovedCatalog(pool);
console.log(JSON.stringify(manifest, null, 2));
} finally {
await pool.end();
}
+184
View File
@@ -0,0 +1,184 @@
import { createHash } from 'node:crypto';
import { mkdir, rename, writeFile } from 'node:fs/promises';
import path from 'node:path';
import type { Pool } from 'pg';
import {
DashcamCatalogExportSchema,
type DashcamCatalogExport,
type MarketSource,
} from '../../web-front/src/entities/dashcam-catalog/model/schema.ts';
import { ingestionConfig } from './config.ts';
type ProductRow = {
id: string;
brand: string;
model: string;
canonical_name: string;
title: string;
discovered_at: Date;
updated_at: Date;
};
type SpecRow = {
product_id: string;
key: string;
value: string | number | boolean;
};
type ImageRow = {
product_id: string;
artifact_url: string;
original_url: string;
source: MarketSource;
sha256: string | null;
width: number | null;
height: number | null;
};
type SourceRow = {
product_id: string;
source: MarketSource;
source_product_id: string | null;
url: string;
stable_content_hash: string;
last_fetched_at: Date;
};
function sha256(buffer: Buffer) {
return createHash('sha256').update(buffer).digest('hex');
}
function exportId(now: Date) {
return `catalog_${now.toISOString().replace(/[-:.]/g, '').replace('T', '_').replace('Z', '')}`;
}
function groupRows<T extends { product_id: string }>(rows: T[]) {
const grouped = new Map<string, T[]>();
for (const row of rows) {
grouped.set(row.product_id, [...(grouped.get(row.product_id) ?? []), row]);
}
return grouped;
}
async function writeFileAtomically(filePath: string, contents: Buffer) {
const temporaryPath = `${filePath}.${process.pid}.tmp`;
await mkdir(path.dirname(filePath), { recursive: true });
await writeFile(temporaryPath, contents);
await rename(temporaryPath, filePath);
}
export async function exportApprovedCatalog(pool: Pool, now = new Date()) {
const createdAt = now.toISOString();
const id = exportId(now);
const [products, specs, images, sources] = await Promise.all([
pool.query<ProductRow>(
`
select id, brand, model, canonical_name, title, discovered_at, updated_at
from products
where moderation_status = 'approved'
order by brand asc, model asc, id asc
`,
),
pool.query<SpecRow>('select product_id, key, value from product_specs'),
pool.query<ImageRow>(
`
select product_id, artifact_url, original_url, source, sha256, width, height
from product_images
where status = 'approved'
and artifact_url is not null
order by product_id, sort_order asc, original_url asc
`,
),
pool.query<SourceRow>(
`
select product_id, source, source_product_id, url, stable_content_hash, last_fetched_at
from product_sources
where product_id is not null
and stable_content_hash is not null
and last_fetched_at is not null
order by product_id, source, url
`,
),
]);
const specsByProduct = groupRows(specs.rows);
const imagesByProduct = groupRows(images.rows);
const sourcesByProduct = groupRows(sources.rows);
const catalog: DashcamCatalogExport = {
schemaVersion: 1,
exportId: id,
createdAt,
products: products.rows
.map((product) => ({
id: product.id,
brand: product.brand,
model: product.model,
canonicalName: product.canonical_name,
title: product.title,
specs: Object.fromEntries((specsByProduct.get(product.id) ?? []).map((spec) => [spec.key, spec.value])),
images: (imagesByProduct.get(product.id) ?? []).map((image) => ({
url: image.artifact_url,
originalUrl: image.original_url,
source: image.source,
...(image.sha256 ? { sha256: image.sha256 } : {}),
...(image.width ? { width: image.width } : {}),
...(image.height ? { height: image.height } : {}),
})),
sources: (sourcesByProduct.get(product.id) ?? []).map((source) => ({
source: source.source,
...(source.source_product_id ? { sourceProductId: source.source_product_id } : {}),
url: source.url,
stableContentHash: source.stable_content_hash,
fetchedAt: source.last_fetched_at.toISOString(),
})),
discoveredAt: product.discovered_at.toISOString(),
updatedAt: product.updated_at.toISOString(),
}))
.filter((product) => product.sources.length > 0),
};
const validatedCatalog = DashcamCatalogExportSchema.parse(catalog);
const artifactBuffer = Buffer.from(`${JSON.stringify(validatedCatalog, null, 2)}\n`);
const artifactHash = sha256(artifactBuffer);
const artifactRelativePath = path.posix.join('artifacts', id, 'catalog.json');
const artifactPath = path.join(ingestionConfig.exportDir, artifactRelativePath);
const manifest = {
schemaVersion: 1 as const,
exportId: id,
createdAt,
artifactPath: artifactRelativePath,
sha256: artifactHash,
itemCount: validatedCatalog.products.length,
};
const manifestBuffer = Buffer.from(`${JSON.stringify(manifest, null, 2)}\n`);
await writeFileAtomically(artifactPath, artifactBuffer);
await writeFileAtomically(path.join(ingestionConfig.exportDir, 'manifest.json'), manifestBuffer);
await pool.query(
`
insert into export_runs (id, schema_version, created_at, artifact_path, manifest_path, sha256, item_count)
values ($1, 1, $2, $3, $4, $5, $6)
on conflict (id) do nothing
`,
[id, createdAt, artifactPath, path.join(ingestionConfig.exportDir, 'manifest.json'), artifactHash, manifest.itemCount],
);
await pool.query(
`
update crawl_jobs job
set state = 'exported',
updated_at = now()
from product_sources source
join products product on product.id = source.product_id
where product.moderation_status = 'approved'
and job.source = source.source
and job.url = source.url
and job.job_type = 'product'
and job.state in ('needs_review', 'approved')
`,
);
return manifest;
}
+107
View File
@@ -0,0 +1,107 @@
import { createHash } from 'node:crypto';
import { mkdir, writeFile } from 'node:fs/promises';
import path from 'node:path';
import { ingestionConfig } from './config.ts';
import type { CrawlFailureReason, FetchClassification, MarketSource } from './types.ts';
type FetchSnapshotInput = {
source: MarketSource;
url: string;
snapshotDir?: string;
fetcher?: typeof fetch;
headers?: HeadersInit;
accept?: string;
};
export type RawFetchSnapshot = {
source: MarketSource;
url: string;
fetchedAt: string;
httpStatus: number;
rawSnapshotHash: string;
rawHtmlPath: string;
html: string;
setCookieHeaders: string[];
};
function getSetCookieHeaders(response: Response) {
const headers = response.headers as Headers & { getSetCookie?: () => string[] };
if (typeof headers.getSetCookie === 'function') return headers.getSetCookie();
const single = response.headers.get('set-cookie');
return single ? [single] : [];
}
function classifyResponse(status: number, body: string): FetchClassification {
const lowered = body.slice(0, 80_000).toLowerCase();
if (status === 404) return { ok: false, reason: 'not_found', retryable: false };
if (status === 429) return { ok: false, reason: 'rate_limited', retryable: true };
if (status === 403) return { ok: false, reason: 'blocked', retryable: true };
if (status >= 500) return { ok: false, reason: 'rate_limited', retryable: true };
if (/captcha|showcaptcha|smartcaptcha|подтвердите, что вы не робот/.test(lowered)) {
return { ok: false, reason: 'captcha', retryable: true };
}
if (status < 200 || status >= 300) return { ok: false, reason: 'needs_manual_review', retryable: false };
return { ok: true };
}
function hash(value: string) {
return createHash('sha256').update(value).digest('hex');
}
function snapshotPath(snapshotDir: string, source: MarketSource, fetchedAt: string, digest: string) {
const day = fetchedAt.slice(0, 10);
return path.join(snapshotDir, source, day, `${digest}.html`);
}
export async function fetchRawSnapshot(input: FetchSnapshotInput): Promise<RawFetchSnapshot> {
const fetcher = input.fetcher ?? fetch;
const response = await fetcher(input.url, {
headers: {
accept: input.accept ?? 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8',
'accept-language': 'ru,en;q=0.8',
'user-agent': 'videoreg-ingestion/1.0 (+https://videoreg.ru)',
...(input.headers ?? {}),
},
redirect: 'follow',
});
const html = await response.text();
const classification = classifyResponse(response.status, html);
if (!classification.ok) {
const error = new Error(`Fetch classified as ${classification.reason} for status ${response.status}`);
(error as Error & { reason?: CrawlFailureReason }).reason = classification.reason;
throw error;
}
const fetchedAt = new Date().toISOString();
const rawSnapshotHash = hash(html);
const rawHtmlPath = snapshotPath(input.snapshotDir ?? ingestionConfig.snapshotDir, input.source, fetchedAt, rawSnapshotHash);
await mkdir(path.dirname(rawHtmlPath), { recursive: true });
await writeFile(rawHtmlPath, html, 'utf8');
return {
source: input.source,
url: response.url || input.url,
fetchedAt,
httpStatus: response.status,
rawSnapshotHash,
rawHtmlPath,
html,
setCookieHeaders: getSetCookieHeaders(response),
};
}
export function getFetchFailureReason(error: unknown) {
return (error as Error & { reason?: CrawlFailureReason }).reason;
}
export function randomDelayMs(minMs = ingestionConfig.fetchDelayMinMs, maxMs = ingestionConfig.fetchDelayMaxMs) {
return Math.floor(minMs + Math.random() * Math.max(0, maxMs - minMs));
}
+87
View File
@@ -0,0 +1,87 @@
const namedEntities: Record<string, string> = {
amp: '&',
gt: '>',
lt: '<',
quot: '"',
apos: "'",
nbsp: ' ',
};
export function decodeHtml(value: string) {
return value.replace(/&(#x[0-9a-f]+|#\d+|[a-z]+);/gi, (entity, code: string) => {
const normalized = code.toLowerCase();
if (normalized.startsWith('#x')) {
return String.fromCodePoint(Number.parseInt(normalized.slice(2), 16));
}
if (normalized.startsWith('#')) {
return String.fromCodePoint(Number.parseInt(normalized.slice(1), 10));
}
return namedEntities[normalized] ?? entity;
});
}
export function stripTags(value: string) {
return decodeHtml(value.replace(/<[^>]*>/g, ' ')).replace(/\s+/g, ' ').trim();
}
export function normalizeWhitespace(value: string) {
return value.trim().replace(/\s+/g, ' ');
}
export function uniqueStable<T>(values: T[]) {
return [...new Set(values)];
}
function parseJsonPayload(raw: string) {
return JSON.parse(decodeHtml(raw).trim()) as unknown;
}
function flattenJsonLd(input: unknown): unknown[] {
if (Array.isArray(input)) {
return input.flatMap(flattenJsonLd);
}
if (input && typeof input === 'object') {
const object = input as Record<string, unknown>;
const graph = object['@graph'];
return graph ? [object, ...flattenJsonLd(graph)] : [object];
}
return [];
}
export function extractJsonLdObjects(html: string) {
const scripts = [...html.matchAll(/<script[^>]+type=["']application\/ld\+json["'][^>]*>([\s\S]*?)<\/script>/gi)];
const objects: unknown[] = [];
for (const [, raw] of scripts) {
try {
objects.push(...flattenJsonLd(parseJsonPayload(raw)));
} catch {
// Ignore malformed analytics or partial JSON-LD blocks.
}
}
return objects;
}
export function getObjectString(value: unknown) {
if (typeof value === 'string') return normalizeWhitespace(value);
if (value && typeof value === 'object' && typeof (value as { name?: unknown }).name === 'string') {
return normalizeWhitespace((value as { name: string }).name);
}
return undefined;
}
export function readMetaContent(html: string, property: string) {
const escaped = property.replace(/[.*+?^${}()|[\]\\]/g, '\\$&');
const pattern = new RegExp(`<meta[^>]+(?:property|name)=["']${escaped}["'][^>]+content=["']([^"']+)["'][^>]*>`, 'i');
const match = html.match(pattern);
return match ? decodeHtml(match[1]).trim() : undefined;
}
+51
View File
@@ -0,0 +1,51 @@
import { readFile, readdir } from 'node:fs/promises';
import path from 'node:path';
import { fileURLToPath } from 'node:url';
import { createPgPool } from './db.ts';
const currentFile = fileURLToPath(import.meta.url);
const migrationsDir = path.resolve(path.dirname(currentFile), '..', 'migrations');
export async function runMigrations() {
const pool = createPgPool();
const client = await pool.connect();
try {
await client.query("select pg_advisory_lock(hashtext('videoreg_ingestion_migrations'))");
await client.query(`
create table if not exists schema_migrations (
filename text primary key,
applied_at timestamptz not null default now()
)
`);
const files = (await readdir(migrationsDir)).filter((file) => file.endsWith('.sql')).sort();
for (const file of files) {
const applied = await client.query('select 1 from schema_migrations where filename = $1', [file]);
if (applied.rowCount) continue;
const sql = await readFile(path.join(migrationsDir, file), 'utf8');
await client.query('begin');
try {
await client.query(sql);
await client.query('insert into schema_migrations (filename) values ($1)', [file]);
await client.query('commit');
console.log(`applied migration ${file}`);
} catch (error) {
await client.query('rollback');
throw error;
}
}
} finally {
await client.query("select pg_advisory_unlock(hashtext('videoreg_ingestion_migrations'))").catch(() => undefined);
client.release();
await pool.end();
}
}
if (import.meta.url === `file://${process.argv[1]}`) {
await runMigrations();
}
+300
View File
@@ -0,0 +1,300 @@
import type { NormalizedMarketProduct, ParsedMarketProduct, ParsedTechnicalSpec } from './types.ts';
const brandAliases = new Map<string, string>([
['ibox', 'iBOX'],
['sho-me', 'Sho-Me'],
['trendvision', 'TrendVision'],
['fujida', 'Fujida'],
['intego', 'INTEGO'],
['spawnson', 'SPAWNSON'],
]);
const specAliases = new Map<string, string>([
['количество камер', 'channels'],
['количество каналов записи видео', 'video_channels'],
['количество каналов записи звука', 'audio_channels'],
['качество видеосъемки/ fps', 'resolution_front'],
['максимальное разрешение видеосъемки', 'resolution_max'],
['макс. разрешение видеосъемки', 'resolution_max'],
['максимальное разрешение камеры заднего вида', 'resolution_max'],
['разрешение видеозаписи (при макс. частоте)', 'resolution_front'],
['разрешение видеозаписи', 'resolution_front'],
['макс. частота кадров', 'frame_rate_max_fps'],
['частота кадров при макс. разрешении', 'frame_rate_max_fps'],
['частота при макс. разрешении камеры заднего вида', 'frame_rate_max_fps'],
['угол обзора (диагональ)', 'view_angle_degrees'],
['угол обзора (ширина)', 'view_angle_width_degrees'],
['угол обзора (высота)', 'view_angle_height_degrees'],
['угол обзора (основная камера)', 'view_angle_degrees'],
['угол обзора', 'view_angle_degrees'],
['gps', 'gps'],
['встроенный gps', 'gps'],
['наличие gps модуля', 'gps'],
['поддержка gps', 'gps'],
['глонасс', 'glonass'],
['встроенный глонасс', 'glonass'],
['наличие глонасс', 'glonass'],
['поддержка глонасс', 'glonass'],
['wi-fi', 'wifi'],
['wifi', 'wifi'],
['поддержка wi-fi', 'wifi'],
['поддержка wifi', 'wifi'],
['трансляция видео через wi-fi', 'wifi'],
['трансляция видео через wifi', 'wifi'],
['управление с помощью смартфона', 'wifi'],
['ночной режим', 'night_mode'],
['режим ночной съемки', 'night_mode'],
['экран', 'screen'],
['запись времени и даты', 'timestamp_recording'],
['запись скорости', 'speed_recording'],
['запись скорости движения', 'speed_recording'],
['встроенный микрофон', 'microphone'],
['запись звука', 'microphone'],
['встроенный динамик', 'speaker'],
['с конденсатором', 'capacitor'],
['максимальный размер карты памяти', 'memory_card_max_gb'],
['макс. объем карты памяти', 'memory_card_max_gb'],
['максимальная емкость карты памяти', 'memory_card_max_gb'],
['матрица', 'sensor_type'],
['размер матрицы', 'sensor_size'],
['количество мегапикселей матрицы', 'sensor_megapixels'],
['поддержка hd', 'hd_support'],
['режим записи', 'recording_mode'],
['запись роликов без разрывов', 'gapless_recording'],
['циклическая запись', 'gapless_recording'],
['непрерывная запись', 'gapless_recording'],
['запись события в отдельный файл', 'event_file_recording'],
['запись событий в отдельный файл', 'event_file_recording'],
['формат записи', 'recording_format'],
['используемый видеокодек', 'video_codec'],
['питание от аккумулятора', 'battery_power'],
['резервный источник питания', 'battery_power'],
['питание от бортовой сети автомобиля', 'car_power'],
['питание от прикуривателя', 'car_power'],
['подключение внешних камер', 'external_cameras'],
['подключение камеры заднего вида', 'external_cameras'],
['дополнительная камера', 'external_cameras'],
['камера заднего вида', 'external_cameras'],
['поддержка карт памяти', 'memory_card_support'],
['база камер', 'speedcam'],
['предупреждение о камерах слежения', 'speedcam'],
['speedcam', 'speedcam'],
['обнаружение радаров', 'radar_detection'],
['магнитный держатель', 'magnetic_mount'],
['крепление на присоске', 'suction_mount'],
['голосовые подсказки', 'voice_prompts'],
['с радар-детектором', 'radar_detector'],
]);
function normalizeKey(value: string) {
return value.trim().toLowerCase().replace(/\s+/g, ' ');
}
function parseBoolean(value: unknown) {
if (typeof value === 'boolean') return value;
if (typeof value !== 'string') return undefined;
const normalized = normalizeKey(value);
if (['да', 'есть', 'true', 'yes'].includes(normalized)) return true;
if (['нет', 'false', 'no', 'отсутствует', 'none'].includes(normalized)) return false;
return undefined;
}
function parseFirstNumber(value: unknown) {
if (typeof value === 'number') return value;
if (typeof value !== 'string') return undefined;
const match = value.replace(',', '.').match(/\d+(?:\.\d+)?/);
return match ? Number(match[0]) : undefined;
}
function normalizeBatteryPower(value: string | number | boolean) {
const booleanValue = parseBoolean(value);
if (booleanValue !== undefined) return booleanValue;
if (typeof value === 'number') return value > 0;
if (typeof value !== 'string') return undefined;
const normalized = normalizeKey(value);
const hasBattery = /аккум|battery|li[-\s]?ion|литий/.test(normalized);
const hasCapacitor = /конденс|capacitor|supercap/.test(normalized);
if (hasBattery && hasCapacitor) return 'hybrid';
if (hasBattery) return 'battery';
if (hasCapacitor) return 'capacitor';
if (['нет', 'false', 'no', 'отсутствует', 'none'].includes(normalized)) return false;
return undefined;
}
function normalizeResolution(value: unknown) {
if (typeof value !== 'string') return undefined;
const match = value.replace(/[×х]/gi, 'x').match(/\d{3,5}\s*x\s*\d{3,5}/i);
return match ? match[0].replace(/\s+/g, '').toLowerCase() : undefined;
}
function normalizeSpecValue(key: string, value: string | number | boolean) {
if (key === 'battery_power') {
return normalizeBatteryPower(value);
}
const booleanValue = parseBoolean(value);
if (booleanValue !== undefined) {
return booleanValue;
}
if (
[
'gps',
'glonass',
'wifi',
'night_mode',
'screen',
'timestamp_recording',
'speed_recording',
'microphone',
'speaker',
'capacitor',
'speedcam',
'radar_detection',
].includes(key)
) {
if (typeof value === 'string' && value.trim()) {
return true;
}
return value;
}
if (
[
'channels',
'video_channels',
'audio_channels',
'frame_rate_max_fps',
'view_angle_degrees',
'view_angle_width_degrees',
'view_angle_height_degrees',
'memory_card_max_gb',
'sensor_megapixels',
].includes(key)
) {
return parseFirstNumber(value);
}
if (key === 'resolution_front' || key === 'resolution_max') {
return normalizeResolution(value);
}
return value;
}
function normalizeBrand(value: string | undefined, title: string) {
const raw = value ?? [...brandAliases.values()].find((brand) => title.toLowerCase().includes(brand.toLowerCase()));
if (!raw) return 'Unknown';
return brandAliases.get(raw.toLowerCase()) ?? raw.trim();
}
export function cleanModelFromTitle(title: string, brand: string) {
const genericWords = [
'видеорегистратор-зеркало',
'зеркало-видеорегистратор',
'видеорегистратор для автомобиля',
'видеорегистратор автомобильный',
'автомобильный видеорегистратор',
'видеорегистратор',
'зеркало',
'цвет черный',
'черный',
];
let model = title;
for (const word of genericWords) {
model = model.replace(new RegExp(word, 'gi'), ' ');
}
model = model.replace(new RegExp(`["«]?${brand.replace(/[.*+?^${}()|[\]\\]/g, '\\$&')}["»]?`, 'gi'), ' ');
model = model.replace(new RegExp(brand.replace(/[.*+?^${}()|[\]\\]/g, '\\$&'), 'gi'), ' ');
return (
model
.replace(/[«»"]/g, ' ')
.replace(/\bwi[\s-]?fi\b/gi, 'Wi-Fi')
.replace(/\s*,\s*/g, ' ')
.replace(/^[\s,+–-]+/g, '')
.replace(/[,+–-]\s*$/g, '')
.replace(/\s+/g, ' ')
.trim() || title.trim()
);
}
function normalizeTechnicalSpec(spec: ParsedTechnicalSpec) {
const normalizedKey = specAliases.get(normalizeKey(spec.label));
const normalizedValue = normalizedKey ? normalizeSpecValue(normalizedKey, spec.value) : undefined;
return {
...spec,
...(normalizedKey ? { normalizedKey } : {}),
...(normalizedValue !== undefined ? { normalizedValue } : {}),
};
}
function buildCanonicalModel(model: string) {
const compact = model
.replace(/\s+\+\s+(?:внутрисалонная|задняя)\s+камера.*$/i, '')
.replace(/\s+\+\s+камера.*$/i, '')
.replace(/\s+с\s+(?:базой камер|сигнатурным|радар-детектором|gps|глонасс).*$/i, '')
.replace(/\s+автомобильн(?:ый|ая|ое|ые)(?:\s+с\b.*)?$/i, '')
.replace(/\s+(?:зеркало|с креплением на зеркало)(?:\s+с\b.*)?$/i, '')
.replace(/[\s,]+(?:ночной режим|ночная съемка)$/i, '')
.replace(/\s+/g, ' ')
.trim();
return compact || model;
}
function createCanonicalName(brand: string, canonicalModel: string) {
if (canonicalModel.toLowerCase().startsWith(brand.toLowerCase())) return canonicalModel;
return `${brand} ${canonicalModel}`.replace(/\s+/g, ' ').trim();
}
export function normalizeMarketProduct(product: ParsedMarketProduct): NormalizedMarketProduct {
const brand = normalizeBrand(product.brand, product.title);
const model = product.model?.trim() || cleanModelFromTitle(product.title, brand);
const canonicalModel = buildCanonicalModel(model);
const rawSpecs = product.rawSpecs.map(normalizeTechnicalSpec);
const specs: NormalizedMarketProduct['specs'] = {};
for (const spec of rawSpecs) {
if (spec.normalizedKey && spec.normalizedValue !== undefined) {
specs[spec.normalizedKey] = spec.normalizedValue;
}
}
return {
source: product.source,
sourceProductId: product.sourceProductId,
sourceUrl: product.sourceUrl,
title: product.title,
brand,
model,
canonicalModel,
canonicalName: createCanonicalName(brand, canonicalModel),
specs,
rawSpecs,
images: product.images,
aggregateRating: product.aggregateRating,
};
}
// Kept as an import-compatible alias for existing Yandex parser callers.
export const normalizeYandexMarketProduct = normalizeMarketProduct;
+188
View File
@@ -0,0 +1,188 @@
import type { Pool } from 'pg';
import type { CrawlFailureReason, CrawlJobRecord, CrawlJobType, MarketSource } from './types.ts';
type CrawlJobRow = {
id: string;
source: MarketSource;
url: string;
job_type: CrawlJobType;
state: CrawlJobRecord['state'];
attempts: number;
next_run_at: Date;
last_error: string | null;
};
type BudgetBucket = CrawlJobType;
const retryableFailureReasons = new Set<CrawlFailureReason>(['rate_limited', 'blocked', 'captcha']);
function toRecord(row: CrawlJobRow): CrawlJobRecord {
return {
id: row.id,
source: row.source,
url: row.url,
jobType: row.job_type,
state: row.state,
attempts: row.attempts,
nextRunAt: row.next_run_at.toISOString(),
...(row.last_error ? { lastError: row.last_error } : {}),
};
}
export function dailyBudgetForJob(jobType: CrawlJobType) {
if (jobType === 'category') return Number(process.env.INGESTION_CATEGORY_DAILY_BUDGET ?? 25);
if (jobType === 'sitemap') return Number(process.env.INGESTION_SITEMAP_DAILY_BUDGET ?? 25);
return Number(process.env.INGESTION_PRODUCT_DAILY_BUDGET ?? 100);
}
export class CrawlQueue {
private readonly pool: Pool;
constructor(pool: Pool) {
this.pool = pool;
}
async enqueue(input: {
source: MarketSource;
url: string;
jobType: CrawlJobType;
nextRunAt?: Date;
state?: CrawlJobRecord['state'];
}) {
await this.pool.query(
`
insert into crawl_jobs (source, url, job_type, state, next_run_at)
values ($1, $2, $3, $4, $5)
on conflict (source, url, job_type) do update
set
next_run_at = least(crawl_jobs.next_run_at, excluded.next_run_at),
updated_at = now()
`,
[input.source, input.url, input.jobType, input.state ?? 'pending', input.nextRunAt ?? new Date()],
);
}
async leaseNext(now = new Date()) {
const result = await this.pool.query<CrawlJobRow>(
`
with candidate as (
select job.id
from crawl_jobs job
left join source_pauses pause
on pause.source = job.source
and pause.paused_until > $1
where job.state = 'pending'
and job.next_run_at <= $1
and pause.source is null
order by job.next_run_at asc, job.created_at asc
for update of job skip locked
limit 1
)
update crawl_jobs job
set state = 'fetching',
attempts = attempts + 1,
updated_at = now()
from candidate
where job.id = candidate.id
returning job.*
`,
[now],
);
return result.rows[0] ? toRecord(result.rows[0]) : undefined;
}
async reserveDailyBudget(source: MarketSource, bucket: BudgetBucket, dailyLimit: number, now = new Date()) {
const day = now.toISOString().slice(0, 10);
const result = await this.pool.query<{ reserved: boolean }>(
`
with existing as (
select used_count, daily_limit
from source_daily_budgets
where source = $1 and budget_bucket = $2 and budget_day = $3
for update
),
upserted as (
insert into source_daily_budgets (source, budget_bucket, budget_day, daily_limit, used_count)
select $1, $2, $3, $4, 1
where not exists (select 1 from existing)
on conflict do nothing
returning true as reserved
),
updated as (
update source_daily_budgets
set used_count = used_count + 1,
daily_limit = $4,
updated_at = now()
where source = $1
and budget_bucket = $2
and budget_day = $3
and exists (select 1 from existing)
and used_count < daily_limit
returning true as reserved
)
select coalesce(
(select reserved from upserted),
(select reserved from updated),
false
) as reserved
`,
[source, bucket, day, dailyLimit],
);
return result.rows[0]?.reserved ?? false;
}
async markState(id: string, state: CrawlJobRecord['state']) {
await this.pool.query(
`
update crawl_jobs
set state = $2,
updated_at = now()
where id = $1
`,
[id, state],
);
}
async markFailure(job: CrawlJobRecord, reason: CrawlFailureReason, message: string) {
const nextRunAt = retryableFailureReasons.has(reason)
? new Date(Date.now() + Math.min(24 * 60 * 60 * 1000, 2 ** job.attempts * 60_000))
: undefined;
await this.pool.query(
`
update crawl_jobs
set state = $2,
last_error = $3,
next_run_at = coalesce($4, next_run_at),
updated_at = now()
where id = $1
`,
[job.id, reason, message.slice(0, 2000), nextRunAt],
);
await this.pool.query(
`
insert into crawl_errors (job_id, source, url, reason, message)
values ($1, $2, $3, $4, $5)
`,
[job.id, job.source, job.url, reason, message.slice(0, 2000)],
);
}
async pauseSource(source: MarketSource, reason: CrawlFailureReason, pausedUntil: Date) {
await this.pool.query(
`
insert into source_pauses (source, reason, paused_until)
values ($1, $2, $3)
on conflict (source) do update
set reason = excluded.reason,
paused_until = excluded.paused_until,
updated_at = now()
`,
[source, reason, pausedUntil],
);
}
}
+69
View File
@@ -0,0 +1,69 @@
import type { Pool } from 'pg';
import { describe, expect, it, vi } from 'vitest';
import { IngestionRepository } from './repository.ts';
import type { NormalizedMarketProduct } from './types.ts';
const product: NormalizedMarketProduct = {
source: 'mvideo',
sourceProductId: '400484046',
sourceUrl: 'https://www.mvideo.ru/products/test-400484046',
title: 'Verification Canonical Link 1',
brand: 'Verification',
model: 'Canonical Link 1',
canonicalModel: 'Canonical Link 1',
canonicalName: 'Verification Canonical Link 1',
specs: { night_mode: false },
rawSpecs: [
{
group: 'Verification',
label: 'Night mode',
value: 'нет',
normalizedKey: 'night_mode',
normalizedValue: false,
},
],
images: [],
};
describe('IngestionRepository', () => {
it('links an exact canonical-name match and serializes source conflicts as JSONB', async () => {
const conflictValues = [
{ source: 'mvideo', value: false },
{ source: 'yandex_market', value: true },
];
const query = vi.fn(async (statement: string) => {
if (statement.includes('select id\n from products')) {
return { rows: [{ id: 'ym_5905353579' }], rowCount: 1 };
}
if (statement.includes('insert into manufacturers')) {
return { rows: [{ id: 'manufacturer-id' }], rowCount: 1 };
}
if (statement.includes('select\n normalized_key')) {
return {
rows: [{ normalized_key: 'night_mode', values_json: conflictValues }],
rowCount: 1,
};
}
return { rows: [], rowCount: 0 };
});
const repository = new IngestionRepository({ query } as unknown as Pool);
const productId = await repository.upsertNormalizedProduct(
product,
'stable-content-hash',
'2026-07-30T23:00:00.000Z',
);
const calls = query.mock.calls as unknown as Array<[string, unknown[] | undefined]>;
const sourceCall = calls.find(([statement]) => statement.includes('insert into product_sources'));
const conflictCall = calls.find(([statement]) => statement.includes('insert into moderation_conflicts'));
const queueCalls = calls.filter(([statement]) => statement.includes('insert into moderation_queue'));
expect(productId).toBe('ym_5905353579');
expect(sourceCall?.[1]?.[0]).toBe('ym_5905353579');
expect(JSON.parse(String(conflictCall?.[1]?.[2]))).toEqual(conflictValues);
expect(queueCalls.at(-1)?.[1]).toEqual(['ym_5905353579', 'same_product']);
});
});
+317
View File
@@ -0,0 +1,317 @@
import { createHash } from 'node:crypto';
import type { Pool } from 'pg';
import type { DiscoveredProduct, NormalizedMarketProduct, ProductSourceSnapshot } from './types.ts';
function publicProductId(product: NormalizedMarketProduct) {
const sourcePrefix = product.source === 'mvideo' ? 'mv' : 'ym';
if (product.sourceProductId) return `${sourcePrefix}_${product.sourceProductId}`;
return `${sourcePrefix}_${Buffer.from(product.sourceUrl).toString('base64url').slice(0, 32)}`;
}
function valueHash(value: unknown) {
return createHash('sha256').update(JSON.stringify(value)).digest('hex');
}
function manufacturerSlug(name: string) {
const slug = name
.trim()
.toLowerCase()
.replace(/[^\p{Letter}\p{Number}]+/gu, '-')
.replace(/^-+|-+$/g, '');
return slug || 'unknown';
}
export class IngestionRepository {
private readonly pool: Pool;
constructor(pool: Pool) {
this.pool = pool;
}
async recordSnapshot(snapshot: ProductSourceSnapshot) {
await this.pool.query(
`
insert into source_snapshots (
source, url, fetched_at, http_status, raw_snapshot_hash,
stable_content_hash, raw_html_path, parsed_json
)
values ($1, $2, $3, $4, $5, $6, $7, $8)
`,
[
snapshot.source,
snapshot.url,
snapshot.fetchedAt,
snapshot.httpStatus,
snapshot.rawSnapshotHash,
snapshot.stableContentHash ?? null,
snapshot.rawHtmlPath ?? null,
snapshot.parsedJson,
],
);
}
async recordDiscoveredProducts(products: DiscoveredProduct[]) {
for (const product of products) {
await this.pool.query(
`
insert into product_sources (
source, url, source_product_id, category_url, title_preview,
first_seen_at, updated_at, state
)
values ($1, $2, $3, $4, $5, $6, now(), 'discovered')
on conflict (source, url) do update
set source_product_id = coalesce(product_sources.source_product_id, excluded.source_product_id),
category_url = coalesce(excluded.category_url, product_sources.category_url),
title_preview = coalesce(excluded.title_preview, product_sources.title_preview),
updated_at = now()
`,
[
product.source,
product.productUrl,
product.sourceProductId ?? null,
product.categoryUrl ?? null,
product.titlePreview ?? null,
product.discoveredAt,
],
);
}
}
async upsertNormalizedProduct(product: NormalizedMarketProduct, stableContentHash: string, fetchedAt: string) {
const sourceScopedProductId = publicProductId(product);
const matchingProduct = await this.pool.query<{ id: string }>(
`
select id
from products
where lower(canonical_name) = lower($1)
and id <> $2
order by discovered_at asc, id asc
limit 1
`,
[product.canonicalName, sourceScopedProductId],
);
const productId = matchingProduct.rows[0]?.id ?? sourceScopedProductId;
const linkedToExistingCanonicalProduct = productId !== sourceScopedProductId;
const manufacturerId = await this.upsertManufacturer(product.brand);
await this.pool.query(
`
insert into products (id, manufacturer_id, brand, model, canonical_name, title, moderation_status, discovered_at, updated_at)
values ($1, $2, $3, $4, $5, $6, 'needs_review', $7, $7)
on conflict (id) do update
set manufacturer_id = excluded.manufacturer_id,
brand = excluded.brand,
model = case when products.model_override then products.model else excluded.model end,
canonical_name = case when products.canonical_name_override then products.canonical_name else excluded.canonical_name end,
title = excluded.title,
updated_at = excluded.updated_at
`,
[productId, manufacturerId, product.brand, product.canonicalModel, product.canonicalName, product.title, fetchedAt],
);
await this.pool.query(
`
insert into product_sources (
product_id, source, url, source_product_id, stable_content_hash,
last_fetched_at, parsed_json, state, updated_at
)
values ($1, $2, $3, $4, $5, $6, $7, 'parsed', now())
on conflict (source, url) do update
set product_id = excluded.product_id,
source_product_id = coalesce(product_sources.source_product_id, excluded.source_product_id),
stable_content_hash = excluded.stable_content_hash,
last_fetched_at = excluded.last_fetched_at,
parsed_json = excluded.parsed_json,
state = excluded.state,
updated_at = now()
`,
[productId, product.source, product.sourceUrl, product.sourceProductId ?? null, stableContentHash, fetchedAt, product],
);
for (const [key, value] of Object.entries(product.specs)) {
await this.pool.query(
`
insert into product_specs (product_id, key, value, updated_at)
values ($1, $2, $3, now())
on conflict (product_id, key) do update
set value = excluded.value,
updated_at = now()
`,
[productId, key, JSON.stringify(value)],
);
}
await this.pool.query(
`
delete from product_source_specs
where product_id = $1
and source = $2
and source_url = $3
`,
[productId, product.source, product.sourceUrl],
);
for (const spec of product.rawSpecs) {
const comparableValue = spec.normalizedValue ?? spec.value;
await this.pool.query(
`
insert into product_source_specs (
product_id, source, source_url, source_product_id, raw_group, raw_label,
raw_value, normalized_key, normalized_value, value_hash, fetched_at
)
values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11)
on conflict (product_id, source, source_url, raw_label) do update
set raw_group = excluded.raw_group,
raw_value = excluded.raw_value,
normalized_key = excluded.normalized_key,
normalized_value = excluded.normalized_value,
value_hash = excluded.value_hash,
fetched_at = excluded.fetched_at,
updated_at = now()
`,
[
productId,
product.source,
product.sourceUrl,
product.sourceProductId ?? null,
spec.group ?? null,
spec.label,
JSON.stringify(spec.value),
spec.normalizedKey ?? null,
spec.normalizedValue === undefined ? null : JSON.stringify(spec.normalizedValue),
valueHash(comparableValue),
fetchedAt,
],
);
}
await this.refreshSpecConflicts(productId);
for (const imageUrl of product.images) {
await this.pool.query(
`
insert into product_images (product_id, original_url, source, status)
values ($1, $2, $3, 'needs_review')
on conflict (product_id, original_url) do nothing
`,
[productId, imageUrl, product.source],
);
}
if (product.aggregateRating) {
await this.pool.query(
`
insert into source_ratings (product_id, source, source_product_id, rating_value, rating_scale, rating_count, fetched_at)
values ($1, $2, $3, $4, $5, $6, $7)
`,
[
productId,
product.source,
product.sourceProductId ?? null,
product.aggregateRating.value,
product.aggregateRating.scale ?? null,
product.aggregateRating.ratingCount ?? null,
fetchedAt,
],
);
}
await this.pool.query(
`
insert into moderation_queue (product_id, reason, state)
values ($1, $2, 'pending')
on conflict (product_id, reason) do nothing
`,
[productId, linkedToExistingCanonicalProduct ? 'same_product' : 'new_product'],
);
return productId;
}
private async upsertManufacturer(name: string) {
const result = await this.pool.query<{ id: string }>(
`
insert into manufacturers (slug, name)
values ($1, $2)
on conflict (slug) do update
set name = excluded.name,
updated_at = now()
returning id
`,
[manufacturerSlug(name), name],
);
return result.rows[0].id;
}
private async refreshSpecConflicts(productId: string) {
await this.pool.query(
`
update moderation_conflicts
set state = 'resolved',
updated_at = now()
where product_id = $1
and conflict_type = 'spec_value'
`,
[productId],
);
const conflicts = await this.pool.query<{
normalized_key: string;
values_json: unknown;
}>(
`
select
normalized_key,
jsonb_agg(
jsonb_build_object(
'source', source,
'url', source_url,
'sourceProductId', source_product_id,
'value', coalesce(normalized_value, raw_value),
'rawLabel', raw_label,
'rawGroup', raw_group
)
order by source, source_url, raw_label
) as values_json
from product_source_specs
where product_id = $1
and normalized_key is not null
group by normalized_key
having count(distinct value_hash) > 1
`,
[productId],
);
for (const conflict of conflicts.rows) {
await this.pool.query(
`
insert into moderation_conflicts (product_id, conflict_type, key, values_json, state)
values ($1, 'spec_value', $2, $3, 'pending')
on conflict (product_id, conflict_type, key) do update
set values_json = excluded.values_json,
state = 'pending',
updated_at = now()
`,
[productId, conflict.normalized_key, JSON.stringify(conflict.values_json)],
);
}
if (conflicts.rowCount) {
await this.pool.query(
`
insert into moderation_queue (product_id, reason, state)
values ($1, 'conflicting_specs', 'pending')
on conflict (product_id, reason) do update
set state = 'pending',
updated_at = now()
`,
[productId],
);
}
}
}
+284
View File
@@ -0,0 +1,284 @@
import { ingestionConfig } from './config.ts';
import { normalizeYandexMarketProduct } from './normalizer.ts';
import { createStableContentHash } from './stableHash.ts';
import {
buildCookieHeader,
discoverMvideoProductIds,
extractMvideoBffBasePath,
extractMvideoSearchQuery,
parseMvideoProduct,
} from './sources/mvideo.ts';
type SmokeProductResult =
| {
ok: true;
index: number;
sourceProductId: string;
productUrl: string;
productApiUrl: string;
httpStatus: number;
title: string;
brand: string;
model: string;
canonicalModel: string;
canonicalName: string;
sources: Array<{
source: 'mvideo';
sourceProductId?: string;
url: string;
fetchedAt: string;
stableContentHash: string;
title: string;
aggregateRating?: {
value: number;
scale?: number;
ratingCount?: number;
};
}>;
technicalSpecs: {
normalized: Record<string, string | number | boolean>;
raw: Array<{
group?: string;
label: string;
value: string | number | boolean;
normalizedKey?: string;
normalizedValue?: string | number | boolean;
}>;
};
images: string[];
stableContentHash: string;
moderation: {
conflicts: [];
};
internalSignals: {
aggregateRating?: {
value: number;
scale?: number;
ratingCount?: number;
};
rawSpecKeys: string[];
};
}
| {
ok: false;
index: number;
sourceProductId: string;
productApiUrl: string;
error: string;
};
function numberEnv(name: string, fallback: number) {
const value = Number(process.env[name]);
return Number.isFinite(value) ? value : fallback;
}
function sleep(ms: number) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
function setCookieHeaders(response: Response) {
const headers = response.headers as Headers & { getSetCookie?: () => string[] };
if (typeof headers.getSetCookie === 'function') return headers.getSetCookie();
const single = response.headers.get('set-cookie');
return single ? [single] : [];
}
async function fetchText(url: string, init?: RequestInit) {
const response = await fetch(url, {
headers: {
accept: 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8',
'accept-language': 'ru,en;q=0.8',
'user-agent': 'Mozilla/5.0 (compatible; videoreg-ingestion-smoke/1.0)',
...(init?.headers ?? {}),
},
redirect: 'follow',
...init,
});
const text = await response.text();
if (!response.ok) {
throw new Error(`HTTP ${response.status} for ${url}`);
}
return {
response,
url: response.url || url,
status: response.status,
text,
};
}
async function fetchJson(url: string, init?: RequestInit) {
const response = await fetch(url, {
headers: {
accept: 'application/json,text/plain,*/*',
'accept-language': 'ru,en;q=0.8',
'user-agent': 'Mozilla/5.0 (compatible; videoreg-ingestion-smoke/1.0)',
...(init?.headers ?? {}),
},
redirect: 'follow',
...init,
});
const text = await response.text();
let payload: unknown;
try {
payload = JSON.parse(text);
} catch {
throw new Error(`Expected JSON response from ${url}, got non-JSON body`);
}
if (!response.ok) {
throw new Error(`HTTP ${response.status} for ${url}`);
}
return {
url: response.url || url,
status: response.status,
payload,
};
}
async function fetchSmokeProduct(
productId: string,
index: number,
bffBasePath: string,
cookieHeader: string,
): Promise<SmokeProductResult> {
const detailsBaseUrl = bffBasePath.endsWith('/') ? bffBasePath : `${bffBasePath}/`;
const detailsUrl = new URL('product-details', detailsBaseUrl);
detailsUrl.searchParams.set('productId', productId);
try {
const response = await fetchJson(detailsUrl.toString(), {
headers: {
cookie: cookieHeader,
},
});
const fetchedAt = new Date().toISOString();
const parsed = parseMvideoProduct(response.payload);
const normalized = normalizeYandexMarketProduct(parsed);
const stableContentHash = createStableContentHash(normalized);
return {
ok: true,
index,
sourceProductId: productId,
productUrl: normalized.sourceUrl,
productApiUrl: detailsUrl.toString(),
httpStatus: response.status,
title: normalized.title,
brand: normalized.brand,
model: normalized.model,
canonicalModel: normalized.canonicalModel,
canonicalName: normalized.canonicalName,
sources: [
{
source: 'mvideo',
sourceProductId: normalized.sourceProductId,
url: normalized.sourceUrl,
fetchedAt,
stableContentHash,
title: normalized.title,
aggregateRating: normalized.aggregateRating,
},
],
technicalSpecs: {
normalized: normalized.specs,
raw: normalized.rawSpecs,
},
images: normalized.images,
stableContentHash,
moderation: {
conflicts: [],
},
internalSignals: {
aggregateRating: normalized.aggregateRating,
rawSpecKeys: Object.keys(parsed.specs),
},
};
} catch (error) {
return {
ok: false,
index,
sourceProductId: productId,
productApiUrl: detailsUrl.toString(),
error: (error as Error).message,
};
}
}
const limit = Math.max(1, numberEnv('INGESTION_SMOKE_LIMIT', 5));
const delayMs = Math.max(0, numberEnv('INGESTION_SMOKE_DELAY_MS', ingestionConfig.fetchDelayMinMs));
const categoryUrl =
process.env.MVIDEO_DASHCAM_CATEGORY_URL ?? 'https://www.mvideo.ru/product-list-page?q=%D0%B2%D0%B8%D0%B4%D0%B5%D0%BE%D1%80%D0%B5%D0%B3%D0%B8%D1%81%D1%82%D1%80%D0%B0%D1%82%D0%BE%D1%80%D1%8B';
const startedAt = new Date().toISOString();
console.error(`Fetching MVideo category page: ${categoryUrl}`);
const categoryResponse = await fetchText(categoryUrl);
const cookieHeader = buildCookieHeader(setCookieHeaders(categoryResponse.response));
if (!cookieHeader) {
throw new Error('MVideo category response did not provide cookies required for BFF calls');
}
const bffBasePath = extractMvideoBffBasePath(categoryResponse.text);
const searchQuery = extractMvideoSearchQuery(categoryResponse.url);
const searchBaseUrl = bffBasePath.endsWith('/') ? bffBasePath : `${bffBasePath}/`;
const searchUrl = new URL('products/v2/search', searchBaseUrl);
searchUrl.searchParams.set('query', searchQuery);
searchUrl.searchParams.set('offset', '0');
searchUrl.searchParams.set('limit', String(limit));
console.error(`Fetching MVideo search: ${searchUrl.toString()}`);
const searchResponse = await fetchJson(searchUrl.toString(), {
headers: {
cookie: cookieHeader,
},
});
const discoveredProductIds = discoverMvideoProductIds(searchResponse.payload);
const selectedProductIds = discoveredProductIds.slice(0, limit);
const products: SmokeProductResult[] = [];
for (const [index, sourceProductId] of selectedProductIds.entries()) {
if (index > 0 && delayMs > 0) {
console.error(`Waiting ${delayMs}ms before next product fetch`);
await sleep(delayMs);
}
console.error(`Fetching MVideo product ${index + 1}/${selectedProductIds.length}: ${sourceProductId}`);
products.push(await fetchSmokeProduct(sourceProductId, index + 1, bffBasePath, cookieHeader));
}
console.log(
JSON.stringify(
{
source: 'mvideo',
mode: 'smoke',
categoryUrl,
startedAt,
finishedAt: new Date().toISOString(),
limit,
delayMs,
category: {
url: categoryResponse.url,
httpStatus: categoryResponse.status,
htmlLength: categoryResponse.text.length,
},
search: {
query: searchQuery,
url: searchUrl.toString(),
httpStatus: searchResponse.status,
discoveredCount: discoveredProductIds.length,
selectedCount: selectedProductIds.length,
},
products,
},
null,
2,
),
);
+204
View File
@@ -0,0 +1,204 @@
import { ingestionConfig } from './config.ts';
import { normalizeYandexMarketProduct } from './normalizer.ts';
import { discoverYandexMarketProductsFromCategory, parseYandexMarketProduct } from './sources/yandexMarket.ts';
import { createStableContentHash } from './stableHash.ts';
type SmokeProductResult =
| {
ok: true;
index: number;
url: string;
httpStatus: number;
htmlLength: number;
sourceProductId?: string;
sources: Array<{
source: 'yandex_market';
sourceProductId?: string;
url: string;
fetchedAt: string;
stableContentHash: string;
title: string;
aggregateRating?: {
value: number;
scale?: number;
ratingCount?: number;
};
}>;
title: string;
brand: string;
model: string;
canonicalModel: string;
canonicalName: string;
technicalSpecs: {
normalized: Record<string, string | number | boolean>;
raw: Array<{
group?: string;
label: string;
value: string | number | boolean;
normalizedKey?: string;
normalizedValue?: string | number | boolean;
}>;
};
images: string[];
stableContentHash: string;
moderation: {
conflicts: [];
};
internalSignals: {
aggregateRating?: {
value: number;
scale?: number;
ratingCount?: number;
};
rawSpecKeys: string[];
};
}
| {
ok: false;
index: number;
url: string;
error: string;
};
function numberEnv(name: string, fallback: number) {
const value = Number(process.env[name]);
return Number.isFinite(value) ? value : fallback;
}
function sleep(ms: number) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
function classifyHtml(status: number, html: string) {
const lowered = html.slice(0, 80_000).toLowerCase();
if (status === 429) return 'rate_limited';
if (status === 403) return 'blocked';
if (/captcha|showcaptcha|smartcaptcha|подтвердите, что вы не робот/.test(lowered)) return 'captcha';
if (status === 404) return 'not_found';
if (status < 200 || status >= 300) return `http_${status}`;
return undefined;
}
async function fetchHtml(url: string) {
const response = await fetch(url, {
headers: {
accept: 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8',
'accept-language': 'ru,en;q=0.8',
},
redirect: 'follow',
});
const html = await response.text();
const failure = classifyHtml(response.status, html);
if (failure) {
throw new Error(`Fetch classified as ${failure} for status ${response.status}`);
}
return {
url: response.url || url,
status: response.status,
html,
};
}
async function fetchSmokeProduct(url: string, index: number): Promise<SmokeProductResult> {
try {
const response = await fetchHtml(url);
const fetchedAt = new Date().toISOString();
const parsed = parseYandexMarketProduct(response.html, url);
const normalized = normalizeYandexMarketProduct(parsed);
const stableContentHash = createStableContentHash(normalized);
return {
ok: true,
index,
url,
httpStatus: response.status,
htmlLength: response.html.length,
sourceProductId: normalized.sourceProductId,
sources: [
{
source: 'yandex_market',
sourceProductId: normalized.sourceProductId,
url: normalized.sourceUrl,
fetchedAt,
stableContentHash,
title: normalized.title,
aggregateRating: normalized.aggregateRating,
},
],
title: normalized.title,
brand: normalized.brand,
model: normalized.model,
canonicalModel: normalized.canonicalModel,
canonicalName: normalized.canonicalName,
technicalSpecs: {
normalized: normalized.specs,
raw: normalized.rawSpecs,
},
images: normalized.images,
stableContentHash,
moderation: {
conflicts: [],
},
internalSignals: {
aggregateRating: normalized.aggregateRating,
rawSpecKeys: Object.keys(parsed.specs),
},
};
} catch (error) {
return {
ok: false,
index,
url,
error: (error as Error).message,
};
}
}
const limit = Math.max(1, numberEnv('INGESTION_SMOKE_LIMIT', 5));
const delayMs = Math.max(0, numberEnv('INGESTION_SMOKE_DELAY_MS', ingestionConfig.fetchDelayMinMs));
const categoryUrl = process.env.YANDEX_MARKET_DASHCAM_CATEGORY_URL ?? ingestionConfig.seedCategoryUrl;
const startedAt = new Date().toISOString();
console.error(`Fetching category: ${categoryUrl}`);
const categoryResponse = await fetchHtml(categoryUrl);
const discovered = discoverYandexMarketProductsFromCategory(categoryResponse.html, categoryUrl, startedAt);
const selected = discovered.slice(0, limit);
const products: SmokeProductResult[] = [];
for (const [index, product] of selected.entries()) {
if (index > 0 && delayMs > 0) {
console.error(`Waiting ${delayMs}ms before next product fetch`);
await sleep(delayMs);
}
console.error(`Fetching product ${index + 1}/${selected.length}: ${product.productUrl}`);
products.push(await fetchSmokeProduct(product.productUrl, index + 1));
}
console.log(
JSON.stringify(
{
source: 'yandex_market',
mode: 'smoke',
categoryUrl,
startedAt,
finishedAt: new Date().toISOString(),
limit,
delayMs,
category: {
httpStatus: categoryResponse.status,
htmlLength: categoryResponse.html.length,
discoveredCount: discovered.length,
selectedCount: selected.length,
},
products,
},
null,
2,
),
);
+171
View File
@@ -0,0 +1,171 @@
import { describe, expect, it } from 'vitest';
import { normalizeYandexMarketProduct } from '../normalizer.ts';
import {
buildCookieHeader,
discoverMvideoProductIds,
extractMvideoBffBasePath,
extractMvideoSearchQuery,
parseMvideoProduct,
} from './mvideo.ts';
describe('MVideo parser', () => {
it('extracts search query and BFF path from category HTML', () => {
const html = `<script>window.MVID_CONFIG={basePath:'https://www.mvideo.ru/bff'};</script>`;
expect(extractMvideoBffBasePath(html)).toBe('https://www.mvideo.ru/bff');
expect(
extractMvideoSearchQuery(
'https://www.mvideo.ru/product-list-page?q=%D0%B2%D0%B8%D0%B4%D0%B5%D0%BE%D1%80%D0%B5%D0%B3%D0%B8%D1%81%D1%82%D1%80%D0%B0%D1%82%D0%BE%D1%80',
),
).toBe('видеорегистратор');
expect(extractMvideoSearchQuery('https://www.mvideo.ru/category/unknown')).toBe('видеорегистраторы');
});
it('builds cookie header and discovers stable product ids from search response', () => {
const cookieHeader = buildCookieHeader([
'MVID_CITY_ID=CityCZ_975; Path=/; Secure',
'MVID_REGION_ID=1; Path=/; Secure',
'MVID_CITY_ID=CityCZ_975; Path=/; Secure',
]);
expect(cookieHeader).toContain('MVID_CITY_ID=CityCZ_975');
expect(cookieHeader).toContain('MVID_REGION_ID=1');
expect(
discoverMvideoProductIds({
body: {
products: ['4254724', 4254724, '400481258', 'bad-id'],
},
}),
).toEqual(['4254724', '400481258']);
expect(
discoverMvideoProductIds({
body: {
type: 'redirect',
url: '/products/videoregistrator-navitel-ar202-nv-400484046',
},
}),
).toEqual(['400484046']);
});
it('maps product details to parsed and normalized shape with full raw specs', () => {
const parsed = parseMvideoProduct({
body: {
productId: '4254724',
name: 'Видеорегистратор Navitel R9 DUAL',
nameTranslit: 'videoregistrator-navitel-r9-dual',
brandName: 'Navitel',
modelName: 'R9 DUAL',
images: ['Big/4254724bb.jpg', '/Big/4254724bb1.jpg'],
rating: {
score: 4.8,
total: 17,
},
properties: {
all: [
{
name: 'Основные характеристики',
properties: [
{
name: 'Поддержка Wi-Fi',
value: 'Да',
},
{
name: 'Качество видеосъемки/ FPS',
value: 'FullHD (1920x1080 Пикс) 30 кадр/сек',
},
{
name: 'Частота при макс. разрешении камеры заднего вида',
value: '30 кадр/сек',
},
{
name: 'Макс. разрешение видеосъемки',
value: '2560x1440',
},
{
name: 'База камер',
value: 'Да',
},
{
name: 'Запись звука',
value: 'Да',
},
{
name: 'Резервный источник питания',
value: 'аккумулятор',
},
],
},
{
name: 'Оптика',
properties: [
{
name: 'Угол обзора (основная камера)',
value: '170',
measure: '°',
},
],
},
],
},
},
});
const normalized = normalizeYandexMarketProduct(parsed);
expect(parsed).toMatchObject({
source: 'mvideo',
sourceProductId: '4254724',
sourceUrl: 'https://www.mvideo.ru/products/videoregistrator-navitel-r9-dual-4254724',
brand: 'Navitel',
model: 'R9 DUAL',
title: 'Видеорегистратор Navitel R9 DUAL',
aggregateRating: {
value: 4.8,
scale: 5,
ratingCount: 17,
},
});
expect(parsed.images).toEqual([
'https://img.mvideo.ru/Big/4254724bb.jpg',
'https://www.mvideo.ru/Big/4254724bb1.jpg',
]);
expect(parsed.rawSpecs).toEqual(
expect.arrayContaining([
expect.objectContaining({
group: 'Основные характеристики',
label: 'Поддержка Wi-Fi',
value: 'Да',
}),
expect.objectContaining({
group: 'Оптика',
label: 'Угол обзора (основная камера)',
value: '170 °',
}),
]),
);
expect(normalized.specs).toEqual(
expect.objectContaining({
wifi: true,
speedcam: true,
microphone: true,
battery_power: 'battery',
resolution_front: '1920x1080',
frame_rate_max_fps: 30,
view_angle_degrees: 170,
}),
);
expect(normalized.rawSpecs).toEqual(
expect.arrayContaining([
expect.objectContaining({
label: 'Поддержка Wi-Fi',
normalizedKey: 'wifi',
normalizedValue: true,
}),
expect.objectContaining({
label: 'Макс. разрешение видеосъемки',
normalizedKey: 'resolution_max',
normalizedValue: '2560x1440',
}),
]),
);
});
});
+238
View File
@@ -0,0 +1,238 @@
import { createHash } from 'node:crypto';
import { normalizeWhitespace, uniqueStable } from '../html.ts';
import type { ParsedMarketProduct, ParsedTechnicalSpec } from '../types.ts';
export const mvideoOrigin = 'https://www.mvideo.ru';
const productIdPattern = /^\d+$/;
type ProductDetailsProperty = {
name?: unknown;
value?: unknown;
measure?: unknown;
};
type ProductDetailsPropertyGroup = {
name?: unknown;
properties?: unknown;
};
function stringValue(value: unknown) {
if (typeof value !== 'string') return undefined;
const normalized = normalizeWhitespace(value);
return normalized ? normalized : undefined;
}
function normalizeSourceImageUrl(url: string) {
const trimmed = url.trim();
if (trimmed.startsWith('http://') || trimmed.startsWith('https://')) return trimmed;
if (trimmed.startsWith('//')) return `https:${trimmed}`;
if (trimmed.startsWith('/')) return `${mvideoOrigin}${trimmed}`;
return `https://img.mvideo.ru/${trimmed.replace(/^\/+/, '')}`;
}
function createFallbackProductId(sourceUrl: string, title: string) {
return createHash('sha1').update(`${sourceUrl}:${title}`).digest('hex').slice(0, 16);
}
function parseSpecValue(value: unknown, measure: unknown): string | number | boolean | undefined {
if (typeof value === 'boolean' || typeof value === 'number') return value;
const normalizedValue = stringValue(value);
if (!normalizedValue) return undefined;
const normalizedMeasure = stringValue(measure);
if (!normalizedMeasure) return normalizedValue;
if (normalizedValue.toLowerCase().includes(normalizedMeasure.toLowerCase())) return normalizedValue;
return `${normalizedValue} ${normalizedMeasure}`;
}
function specsRecord(technicalSpecs: ParsedTechnicalSpec[]) {
const record: Record<string, string | number | boolean> = {};
for (const spec of technicalSpecs) {
record[spec.label] = spec.value;
}
return record;
}
export function extractMvideoBffBasePath(html: string) {
const directMatch = html.match(/basePath\s*:\s*['"]([^'"]+)['"]/i)?.[1];
const jsonMatch = html.match(/"basePath"\s*:\s*"([^"]+)"/i)?.[1];
const value = directMatch ?? jsonMatch;
if (!value) return `${mvideoOrigin}/bff`;
if (value.startsWith('http://') || value.startsWith('https://')) return value;
if (value.startsWith('/')) return `${mvideoOrigin}${value}`;
return `${mvideoOrigin}/${value.replace(/^\/+/, '')}`;
}
export function extractMvideoSearchQuery(categoryUrl: string) {
try {
const url = new URL(categoryUrl);
const query = url.searchParams.get('q') ?? url.searchParams.get('query');
if (query?.trim()) return query.trim();
} catch {
// Fall through to default query.
}
return 'видеорегистраторы';
}
export function buildMvideoProductUrl(productId: string) {
return `${mvideoOrigin}/product/${encodeURIComponent(productId)}`;
}
export function buildMvideoBffUrl(basePath: string, endpoint: string) {
const normalizedBasePath = basePath.endsWith('/') ? basePath : `${basePath}/`;
return new URL(endpoint.replace(/^\/+/, ''), normalizedBasePath);
}
export function buildCookieHeader(setCookieHeaders: string[]) {
const cookies = new Map<string, string>();
for (const header of setCookieHeaders) {
const cookiePart = header.split(';')[0]?.trim();
if (!cookiePart) continue;
const separatorIndex = cookiePart.indexOf('=');
if (separatorIndex <= 0) continue;
const name = cookiePart.slice(0, separatorIndex).trim();
const value = cookiePart.slice(separatorIndex + 1).trim();
if (!name || !value) continue;
cookies.set(name, `${name}=${value}`);
}
return [...cookies.values()].join('; ');
}
export function discoverMvideoProductIds(searchPayload: unknown) {
if (!searchPayload || typeof searchPayload !== 'object') {
throw new Error('MVideo search payload is not an object');
}
const body = (searchPayload as { body?: unknown }).body;
const products = body && typeof body === 'object' ? (body as { products?: unknown }).products : undefined;
const productIds = Array.isArray(products)
? products
.map((value) => (typeof value === 'number' || typeof value === 'string' ? String(value).trim() : ''))
.filter((value): value is string => productIdPattern.test(value))
: [];
if (productIds.length > 0) return uniqueStable(productIds);
// MVideo converts exact one-result searches into a redirect instead of a
// one-element products array. The destination still carries the public ID.
const redirectUrl =
typeof (body as { url?: unknown }).url === 'string'
? (body as { url: string }).url
: typeof (body as { redirectUrl?: unknown }).redirectUrl === 'string'
? (body as { redirectUrl: string }).redirectUrl
: undefined;
const redirectProductId = redirectUrl?.match(/(\d{6,13})(?:[/?#]|$)/)?.[1];
if (redirectProductId) return [redirectProductId];
throw new Error('MVideo search payload does not include product identifiers');
}
export function parseMvideoProduct(detailsPayload: unknown): ParsedMarketProduct {
if (!detailsPayload || typeof detailsPayload !== 'object') {
throw new Error('MVideo product payload is not an object');
}
const body = (detailsPayload as { body?: unknown }).body;
if (!body || typeof body !== 'object') {
throw new Error('MVideo product payload has no body');
}
const product = body as Record<string, unknown>;
const title = stringValue(product.name);
if (!title) {
throw new Error('MVideo product title was not found');
}
const sourceProductIdValue = product.productId;
const sourceProductId =
(typeof sourceProductIdValue === 'number' || typeof sourceProductIdValue === 'string'
? String(sourceProductIdValue).trim()
: undefined) ?? undefined;
const slug = stringValue(product.nameTranslit);
const sourceUrl = slug && sourceProductId
? `${mvideoOrigin}/products/${slug}-${sourceProductId}`
: sourceProductId
? `${mvideoOrigin}/product/${sourceProductId}`
: mvideoOrigin;
const brand = stringValue(product.brandName);
const model = stringValue(product.modelName);
const images = uniqueStable(
(Array.isArray(product.images) ? product.images : [])
.filter((value): value is string => typeof value === 'string' && value.trim().length > 0)
.map(normalizeSourceImageUrl),
);
const rawSpecs: ParsedTechnicalSpec[] = [];
const rawGroups = (product.properties as { all?: unknown } | undefined)?.all;
if (Array.isArray(rawGroups)) {
for (const rawGroup of rawGroups as ProductDetailsPropertyGroup[]) {
const groupName = stringValue(rawGroup.name);
const properties = Array.isArray(rawGroup.properties) ? rawGroup.properties : [];
for (const rawSpec of properties as ProductDetailsProperty[]) {
const label = stringValue(rawSpec.name);
const value = parseSpecValue(rawSpec.value, rawSpec.measure);
if (!label || value === undefined) continue;
rawSpecs.push({
...(groupName ? { group: groupName } : {}),
label,
value,
});
}
}
}
const ratingRaw = product.rating;
const aggregateRating =
ratingRaw && typeof ratingRaw === 'object' && Number.isFinite(Number((ratingRaw as { score?: unknown }).score))
? {
value: Number((ratingRaw as { score: number | string }).score),
scale: 5,
...(Number.isFinite(Number((ratingRaw as { total?: unknown }).total))
? { ratingCount: Number((ratingRaw as { total: number | string }).total) }
: {}),
}
: undefined;
return {
source: 'mvideo',
sourceProductId: sourceProductId ?? createFallbackProductId(sourceUrl, title),
title,
...(brand ? { brand } : {}),
...(model ? { model } : {}),
images,
specs: specsRecord(rawSpecs),
rawSpecs,
sourceUrl,
aggregateRating,
};
}
+154
View File
@@ -0,0 +1,154 @@
import { describe, expect, it } from 'vitest';
import { normalizeYandexMarketProduct } from '../normalizer';
import { createStableContentHash } from '../stableHash';
import { discoverYandexMarketProductsFromCategory, parseYandexMarketProduct } from './yandexMarket';
describe('Yandex Market parser', () => {
it('discovers and canonicalizes card and product URLs from category HTML', () => {
const discovered = discoverYandexMarketProductsFromCategory(
`
<a href="/card/ibox-f5-wifi-videoregistrator-s-signaturnym-radar-detektorom/5197387201?do-waremd5=x&amp;cpc=y">F5</a>
<script>{"url":"https://market.yandex.ru/product/312774508?offerid=secret&sku=103274083526"}</script>
`,
'https://market.yandex.ru/catalog--avtomobilnye-videoregistratory/82798275/list?hid=82798269',
'2026-05-16T20:00:00.000Z',
);
expect(discovered).toMatchObject([
{
productUrl:
'https://market.yandex.ru/card/ibox-f5-wifi-videoregistrator-s-signaturnym-radar-detektorom/5197387201',
sourceProductId: '5197387201',
},
{
productUrl: 'https://market.yandex.ru/product/312774508',
sourceProductId: '312774508',
},
]);
});
it('maps product JSON-LD and DOM specs without prices, sellers or reviews', () => {
const html = `
<script type="application/ld+json">
{
"@context": "https://schema.org",
"@type": "Product",
"brand": "iBOX",
"image": "https://avatars.mds.yandex.net/get-mpic/18787582/2a000/orig",
"name": "Видеорегистратор для автомобиля iBOX F5 + WiFi с радар-детектором",
"offers": { "price": "17816", "availability": "https://schema.org/InStock" },
"aggregateRating": { "ratingValue": "4.8", "ratingCount": "17" },
"url": "https://market.yandex.ru/card/ibox-f5-wifi-videoregistrator-s-signaturnym-radar-detektorom/5197387201"
}
</script>
<label class="_3J2L5">
<div><span>Экран</span></div>
<div class="_2yz50"><span>да</span></div>
</label>
<div>
<span data-auto="product-spec">Количество камер</span>
<div class="eXP5k"><div><div class="ds-text"><span>1</span></div></div></div>
<span data-auto="product-spec">Разрешение видеозаписи (при макс. частоте)</span>
<div class="eXP5k"><div><div class="ds-text"><span>2304x1296</span></div></div></div>
<span data-auto="product-spec">Wi-Fi</span>
<div class="eXP5k"><div><div class="ds-text"><span>да</span></div></div></div>
<span data-auto="product-spec">С радар-детектором</span>
<div class="eXP5k"><div><div class="ds-text"><span>нет</span></div></div></div>
<label for="group-collapse-Камера"><span>Камера</span></label>
<span data-auto="product-spec">Угол обзора (диагональ)</span>
<div class="eXP5k"><div><div class="ds-text"><span>170°</span></div></div></div>
<span data-auto="product-spec">Матрица</span>
<div class="eXP5k"><div><div class="ds-text"><span>CMOS</span></div></div></div>
</div>
<span class="ds-text ds-text_weight_med">Основные характеристики</span>
<div data-apiary-widget-name=@card/AboutSpecV2>
<div><span data-auto="product-spec">Режим записи</span></div>
<div class="ds-flex uA8Xn"><span>циклическая</span></div>
</div>
`;
const parsed = parseYandexMarketProduct(
html,
'https://market.yandex.ru/card/ibox-f5-wifi-videoregistrator-s-signaturnym-radar-detektorom/5197387201',
);
const normalized = normalizeYandexMarketProduct(parsed);
expect(parsed).toMatchObject({
source: 'yandex_market',
sourceProductId: '5197387201',
brand: 'iBOX',
title: 'Видеорегистратор для автомобиля iBOX F5 + WiFi с радар-детектором',
aggregateRating: {
value: 4.8,
ratingCount: 17,
},
});
expect(normalized.specs).toEqual({
channels: 1,
resolution_front: '2304x1296',
wifi: true,
radar_detector: false,
view_angle_degrees: 170,
sensor_type: 'CMOS',
recording_mode: 'циклическая',
screen: true,
});
expect(normalized.canonicalModel).toBe('F5 + Wi-Fi');
expect(normalized.canonicalName).toBe('iBOX F5 + Wi-Fi');
expect(normalized.rawSpecs).toEqual(
expect.arrayContaining([
expect.objectContaining({
group: 'Камера',
label: 'Матрица',
value: 'CMOS',
normalizedKey: 'sensor_type',
normalizedValue: 'CMOS',
}),
]),
);
expect(parsed.rawSpecs).toEqual(
expect.arrayContaining([
expect.objectContaining({
group: 'Основные характеристики',
label: 'Режим записи',
value: 'циклическая',
}),
]),
);
expect(createStableContentHash(normalized)).toMatch(/^[a-f0-9]{64}$/);
expect(JSON.stringify(normalized)).not.toContain('17816');
});
it('keeps stable model attributes and removes feature tails from canonical model', () => {
const normalized = normalizeYandexMarketProduct({
source: 'yandex_market',
sourceProductId: '5247127917',
sourceUrl:
'https://market.yandex.ru/card/videoregistrator-avtomobilnyy-12-mp-1296p-s-dvumya-kamerami-2k-spawnson-full-hd-antiblikovyy-obyektiv/5247127917',
title: 'Видеорегистратор SPAWNSON "Spawnson", 24МП, Wi-Fi, ночной режим',
brand: 'SPAWNSON',
images: [],
specs: {},
rawSpecs: [],
});
expect(normalized.model).toBe('24МП Wi-Fi ночной режим');
expect(normalized.canonicalModel).toBe('24МП Wi-Fi');
expect(normalized.canonicalName).toBe('SPAWNSON 24МП Wi-Fi');
});
it('removes mirror form factor and automotive descriptor tails from canonical model', () => {
const normalized = normalizeYandexMarketProduct({
source: 'yandex_market',
sourceProductId: 'mirror-1',
sourceUrl: 'https://market.yandex.ru/product/mirror-1',
title: 'Видеорегистратор-зеркало iBOX Rover 2 автомобильный с базой камер',
brand: 'iBOX',
images: [],
specs: {},
rawSpecs: [],
});
expect(normalized.canonicalModel).toBe('Rover 2');
expect(normalized.canonicalName).toBe('iBOX Rover 2');
});
});
+272
View File
@@ -0,0 +1,272 @@
import { createHash } from 'node:crypto';
import { decodeHtml, extractJsonLdObjects, getObjectString, normalizeWhitespace, readMetaContent, stripTags, uniqueStable } from '../html.ts';
import type { DiscoveredProduct, ParsedMarketProduct, ParsedTechnicalSpec } from '../types.ts';
const marketOrigin = 'https://market.yandex.ru';
const productUrlPattern = /(?:https?:\/\/market\.yandex\.ru)?\/(?:card|product)\/[^\s"'<>\\]+/gi;
function getJsonLdType(object: Record<string, unknown>) {
const type = object['@type'];
if (Array.isArray(type)) return type.map(String);
if (typeof type === 'string') return [type];
return [];
}
function isJsonLdProduct(object: unknown): object is Record<string, unknown> {
return !!object && typeof object === 'object' && getJsonLdType(object as Record<string, unknown>).includes('Product');
}
function normalizeSourceImageUrl(url: string) {
const decoded = decodeHtml(url).trim();
if (decoded.startsWith('//')) return `https:${decoded}`;
return decoded;
}
function canonicalizeMarketUrl(rawUrl: string) {
const decoded = decodeHtml(rawUrl).replace(/,$/, '');
const url = new URL(decoded, marketOrigin);
if (url.hostname !== 'market.yandex.ru') {
throw new Error(`Unsupported Yandex Market host: ${url.hostname}`);
}
const parts = url.pathname.split('/').filter(Boolean);
const type = parts[0];
if (type === 'product' && parts[1]) {
return `${marketOrigin}/product/${parts[1]}`;
}
if (type === 'card' && parts[1] && parts[2]) {
return `${marketOrigin}/card/${parts[1]}/${parts[2]}`;
}
throw new Error(`Unsupported Yandex Market product URL: ${rawUrl}`);
}
function sourceProductIdFromUrl(url: string) {
const parsed = new URL(url);
const parts = parsed.pathname.split('/').filter(Boolean);
if (parts[0] === 'product') return parts[1];
if (parts[0] === 'card') return parts[2];
return undefined;
}
function titleFromCardSlug(url: string) {
const parsed = new URL(url);
const parts = parsed.pathname.split('/').filter(Boolean);
if (parts[0] !== 'card') return undefined;
return parts[1]
?.split('-')
.filter(Boolean)
.map((part) => part.charAt(0).toUpperCase() + part.slice(1))
.join(' ');
}
function normalizeProductUrlCandidate(rawUrl: string) {
try {
const url = canonicalizeMarketUrl(rawUrl);
return {
productUrl: url,
sourceProductId: sourceProductIdFromUrl(url),
titlePreview: titleFromCardSlug(url),
};
} catch {
return undefined;
}
}
function extractRating(value: unknown): ParsedMarketProduct['aggregateRating'] {
if (!value || typeof value !== 'object') return undefined;
const object = value as Record<string, unknown>;
const ratingValue = Number(object.ratingValue);
const ratingCount = Number(object.ratingCount ?? object.reviewCount);
const bestRating = Number(object.bestRating);
if (!Number.isFinite(ratingValue)) return undefined;
return {
value: ratingValue,
...(Number.isFinite(bestRating) ? { scale: bestRating } : {}),
...(Number.isFinite(ratingCount) ? { ratingCount } : {}),
};
}
function extractJsonLdProduct(html: string) {
return extractJsonLdObjects(html).find(isJsonLdProduct);
}
function extractJsonLdImages(product: Record<string, unknown> | undefined) {
const image = product?.image;
const values = Array.isArray(image) ? image : image ? [image] : [];
return values.filter((value): value is string => typeof value === 'string').map(normalizeSourceImageUrl);
}
function extractMetaImages(html: string) {
return [readMetaContent(html, 'og:image')]
.filter((value): value is string => !!value)
.map(normalizeSourceImageUrl);
}
function extractSpecGroup(htmlBeforeSpec: string) {
const matches = [
...htmlBeforeSpec.matchAll(/for=["']group-collapse-([^"']+)["'][\s\S]{0,500}?<span[^>]*>([\s\S]*?)<\/span>/gi),
...htmlBeforeSpec.matchAll(
/<span[^>]*class=["'][^"']*ds-text_weight_med[^"']*["'][^>]*>([\s\S]*?)<\/span>/gi,
),
];
const match = matches
.filter((candidate) => candidate.index !== undefined)
.sort((left, right) => (left.index ?? 0) - (right.index ?? 0))
.at(-1);
return match ? stripTags(match[2] || match[1]) : undefined;
}
function extractTechnicalSpecsFromDom(html: string): ParsedTechnicalSpec[] {
const specs: ParsedTechnicalSpec[] = [];
const seen = new Map<string, number>();
const legacySpecPattern =
/<span[^>]*data-auto=["']product-spec["'][^>]*>([\s\S]*?)<\/span>[\s\S]{0,1400}?<div[^>]+class=["'][^"']*ds-text[^"']*["'][^>]*>\s*<span>([\s\S]*?)<\/span>/gi;
const currentSpecPattern =
/data-apiary-widget-name=@card\/AboutSpecV2[^>]*>[\s\S]*?<span[^>]*data-auto=["']product-spec["'][^>]*>([\s\S]*?)<\/span>[\s\S]*?<div[^>]*class=["'][^"']*(?:uA8Xn|_2yz50)[^"']*["'][^>]*>\s*<span[^>]*>([\s\S]*?)<\/span>/gi;
const compactSpecPattern =
/<label[^>]*>[\s\S]*?<span[^>]*>([\s\S]*?)<\/span>[\s\S]*?<div[^>]*class=["'][^"']*_2yz50[^"']*["'][^>]*>\s*<span[^>]*>([\s\S]*?)<\/span>[\s\S]*?<\/label>/gi;
function addSpec(match: RegExpMatchArray) {
const rawName = match[1];
const rawValue = match[2];
const label = stripTags(rawName);
const value = stripTags(rawValue);
if (!label || !value || label === value) return;
const group = extractSpecGroup(html.slice(Math.max(0, (match.index ?? 0) - 8_000), match.index));
const dedupeKey = `${label}:${value}`;
const existingIndex = seen.get(dedupeKey);
if (existingIndex !== undefined) {
if (group && !specs[existingIndex]?.group) {
specs[existingIndex] = { ...specs[existingIndex], group };
}
return;
}
seen.set(dedupeKey, specs.length);
specs.push({
...(group ? { group } : {}),
label,
value,
});
}
for (const match of html.matchAll(currentSpecPattern)) {
addSpec(match);
}
for (const match of html.matchAll(legacySpecPattern)) {
addSpec(match);
}
for (const match of html.matchAll(compactSpecPattern)) {
addSpec(match);
}
return specs;
}
function specsRecord(technicalSpecs: ParsedTechnicalSpec[]) {
const record: Record<string, string | number | boolean> = {};
for (const spec of technicalSpecs) {
record[spec.label] = spec.value;
}
return record;
}
function extractTitle(html: string, product?: Record<string, unknown>) {
const jsonTitle = getObjectString(product?.name);
const ogTitle = readMetaContent(html, 'og:title');
const titleTag = html.match(/<title[^>]*>([\s\S]*?)<\/title>/i)?.[1];
return jsonTitle ?? ogTitle ?? (titleTag ? stripTags(titleTag) : undefined);
}
function createFallbackProductId(sourceUrl: string, title: string) {
return createHash('sha1').update(`${sourceUrl}:${title}`).digest('hex').slice(0, 16);
}
export function discoverYandexMarketProductsFromCategory(
html: string,
categoryUrl: string,
discoveredAt = new Date().toISOString(),
): DiscoveredProduct[] {
const products = [...html.matchAll(productUrlPattern)]
.map((match) => normalizeProductUrlCandidate(match[0]))
.filter((value): value is NonNullable<ReturnType<typeof normalizeProductUrlCandidate>> => !!value)
.map((value) => ({
source: 'yandex_market' as const,
categoryUrl,
productUrl: value.productUrl,
sourceProductId: value.sourceProductId,
titlePreview: value.titlePreview,
discoveredAt,
}));
const byUrl = new Map<string, DiscoveredProduct>();
for (const product of products) {
byUrl.set(product.productUrl, product);
}
return [...byUrl.values()];
}
export function parseYandexMarketProduct(html: string, sourceUrl: string): ParsedMarketProduct {
const product = extractJsonLdProduct(html);
const title = extractTitle(html, product);
if (!title) {
throw new Error('Yandex Market product title was not found');
}
const canonicalUrl = normalizeProductUrlCandidate(getObjectString(product?.url) ?? sourceUrl)?.productUrl ?? sourceUrl;
const sourceProductId =
sourceProductIdFromUrl(canonicalUrl) ??
(typeof product?.sku === 'string' ? product.sku : undefined) ??
createFallbackProductId(canonicalUrl, title);
const brand = getObjectString(product?.brand);
const images = uniqueStable([...extractJsonLdImages(product), ...extractMetaImages(html)]).filter((url) =>
url.includes('/get-mpic/'),
);
const rawSpecs = extractTechnicalSpecsFromDom(html);
return {
source: 'yandex_market',
sourceProductId,
title: normalizeWhitespace(title),
...(brand ? { brand } : {}),
images,
specs: specsRecord(rawSpecs),
rawSpecs,
sourceUrl: canonicalUrl,
aggregateRating: extractRating(product?.aggregateRating),
};
}
export function canonicalizeYandexMarketProductUrl(rawUrl: string) {
return canonicalizeMarketUrl(rawUrl);
}
+35
View File
@@ -0,0 +1,35 @@
import { createHash } from 'node:crypto';
import type { NormalizedMarketProduct } from './types.ts';
function stableJson(value: unknown): string {
if (Array.isArray(value)) {
return `[${value.map(stableJson).join(',')}]`;
}
if (value && typeof value === 'object') {
const object = value as Record<string, unknown>;
return `{${Object.keys(object)
.sort()
.map((key) => `${JSON.stringify(key)}:${stableJson(object[key])}`)
.join(',')}}`;
}
return JSON.stringify(value);
}
export function createStableContentHash(product: NormalizedMarketProduct) {
const semanticPayload = {
sourceProductId: product.sourceProductId,
canonicalProductUrl: product.sourceUrl,
title: product.title,
brand: product.brand,
model: product.model,
canonicalName: product.canonicalName,
specs: product.specs,
rawSpecs: product.rawSpecs,
images: product.images,
};
return createHash('sha256').update(stableJson(semanticPayload)).digest('hex');
}
+126
View File
@@ -0,0 +1,126 @@
import type { ExportedSpecValue } from '../../web-front/src/entities/dashcam-catalog/model/schema.ts';
export type MarketSource = 'yandex_market' | 'mvideo';
export type CrawlJobType = 'category' | 'product' | 'sitemap';
export type CrawlJobState =
| 'pending'
| 'fetching'
| 'fetched'
| 'parsed'
| 'normalized'
| 'deduplicated'
| 'needs_review'
| 'approved'
| 'exported'
| CrawlFailureReason;
export type CrawlFailureReason =
| 'rate_limited'
| 'blocked'
| 'captcha'
| 'not_found'
| 'parse_failed'
| 'needs_manual_review';
export type DiscoveredProduct = {
source: MarketSource;
categoryUrl?: string;
productUrl: string;
sourceProductId?: string;
titlePreview?: string;
discoveredAt: string;
};
export type ParsedTechnicalSpec = {
group?: string;
label: string;
value: string | number | boolean;
};
export type NormalizedTechnicalSpec = ParsedTechnicalSpec & {
normalizedKey?: string;
normalizedValue?: ExportedSpecValue;
};
export type ParsedMarketProduct = {
source: MarketSource;
sourceProductId?: string;
title: string;
brand?: string;
model?: string;
aggregateRating?: {
value: number;
scale?: number;
ratingCount?: number;
};
images: string[];
specs: Record<string, string | number | boolean>;
rawSpecs: ParsedTechnicalSpec[];
sourceUrl: string;
};
export type NormalizedMarketProduct = {
source: MarketSource;
sourceProductId?: string;
sourceUrl: string;
title: string;
brand: string;
model: string;
canonicalModel: string;
canonicalName: string;
specs: Record<string, ExportedSpecValue>;
rawSpecs: NormalizedTechnicalSpec[];
images: string[];
aggregateRating?: ParsedMarketProduct['aggregateRating'];
};
export type ProductSourceMetadata = {
source: MarketSource;
sourceProductId?: string;
url: string;
fetchedAt?: string;
stableContentHash?: string;
title: string;
aggregateRating?: ParsedMarketProduct['aggregateRating'];
};
export type SpecConflict = {
key: string;
values: Array<{
source: MarketSource;
url: string;
value: ExportedSpecValue;
}>;
};
export type ProductSourceSnapshot = {
source: MarketSource;
url: string;
fetchedAt: string;
httpStatus: number;
rawSnapshotHash: string;
stableContentHash?: string;
rawHtmlPath?: string;
parsedJson: Record<string, unknown>;
};
export type FetchClassification =
| { ok: true }
| {
ok: false;
reason: CrawlFailureReason;
retryable: boolean;
};
export type CrawlJobRecord = {
id: string;
source: MarketSource;
url: string;
jobType: CrawlJobType;
state: CrawlJobState;
attempts: number;
nextRunAt: string;
lastError?: string;
};
+281
View File
@@ -0,0 +1,281 @@
import { ingestionConfig } from './config.ts';
import { createPgPool } from './db.ts';
import { fetchRawSnapshot, getFetchFailureReason, randomDelayMs, type RawFetchSnapshot } from './fetcher.ts';
import { normalizeMarketProduct } from './normalizer.ts';
import { CrawlQueue, dailyBudgetForJob } from './queue.ts';
import { IngestionRepository } from './repository.ts';
import {
buildCookieHeader,
buildMvideoBffUrl,
buildMvideoProductUrl,
discoverMvideoProductIds,
extractMvideoBffBasePath,
extractMvideoSearchQuery,
parseMvideoProduct,
} from './sources/mvideo.ts';
import { discoverYandexMarketProductsFromCategory, parseYandexMarketProduct } from './sources/yandexMarket.ts';
import { createStableContentHash } from './stableHash.ts';
import type { CrawlFailureReason, CrawlJobRecord, DiscoveredProduct } from './types.ts';
function sleep(ms: number) {
return new Promise((resolve) => setTimeout(resolve, ms));
}
function isStopReason(reason: string | undefined) {
return reason === 'captcha' || reason === 'blocked' || reason === 'rate_limited';
}
function sourceError(reason: CrawlFailureReason, message: string) {
const error = new Error(message) as Error & { reason?: CrawlFailureReason };
error.reason = reason;
return error;
}
function parseJsonSnapshot(snapshot: RawFetchSnapshot, label: string) {
try {
return JSON.parse(snapshot.html) as unknown;
} catch (error) {
throw sourceError('parse_failed', `Expected JSON from ${label}: ${(error as Error).message}`);
}
}
async function recordSnapshot(
repository: IngestionRepository,
snapshot: RawFetchSnapshot,
parsedJson: Record<string, unknown>,
stableContentHash?: string,
) {
await repository.recordSnapshot({
source: snapshot.source,
url: snapshot.url,
fetchedAt: snapshot.fetchedAt,
httpStatus: snapshot.httpStatus,
rawSnapshotHash: snapshot.rawSnapshotHash,
rawHtmlPath: snapshot.rawHtmlPath,
parsedJson,
...(stableContentHash ? { stableContentHash } : {}),
});
}
async function processYandexCategory(
job: CrawlJobRecord,
repository: IngestionRepository,
queue: CrawlQueue,
) {
const snapshot = await fetchRawSnapshot({
source: job.source,
url: job.url,
snapshotDir: ingestionConfig.snapshotDir,
});
const products = discoverYandexMarketProductsFromCategory(snapshot.html, job.url, snapshot.fetchedAt);
await recordSnapshot(repository, snapshot, { discoveredProducts: products });
await repository.recordDiscoveredProducts(products);
for (const product of products) {
await queue.enqueue({
source: product.source,
url: product.productUrl,
jobType: 'product',
});
}
await queue.markState(job.id, 'parsed');
console.log(`yandex_market category parsed: ${products.length} product candidates`);
}
async function processMvideoCategory(
job: CrawlJobRecord,
repository: IngestionRepository,
queue: CrawlQueue,
) {
const categorySnapshot = await fetchRawSnapshot({
source: job.source,
url: job.url,
snapshotDir: ingestionConfig.snapshotDir,
});
const bffBasePath = extractMvideoBffBasePath(categorySnapshot.html);
const searchQuery = extractMvideoSearchQuery(categorySnapshot.url);
const searchUrl = buildMvideoBffUrl(bffBasePath, 'products/v2/search');
const cookieHeader = buildCookieHeader(categorySnapshot.setCookieHeaders);
searchUrl.searchParams.set('query', searchQuery);
searchUrl.searchParams.set('offset', '0');
searchUrl.searchParams.set('limit', String(ingestionConfig.productDailyBudget));
const searchSnapshot = await fetchRawSnapshot({
source: job.source,
url: searchUrl.toString(),
snapshotDir: ingestionConfig.snapshotDir,
accept: 'application/json,text/plain,*/*',
...(cookieHeader ? { headers: { cookie: cookieHeader } } : {}),
});
const productIds = discoverMvideoProductIds(parseJsonSnapshot(searchSnapshot, 'MVideo search'));
const products: DiscoveredProduct[] = productIds.map((sourceProductId) => ({
source: 'mvideo',
sourceProductId,
productUrl: buildMvideoProductUrl(sourceProductId),
categoryUrl: categorySnapshot.url,
discoveredAt: categorySnapshot.fetchedAt,
}));
await recordSnapshot(repository, categorySnapshot, {
bffBasePath,
searchQuery,
discoveredProducts: products,
});
await recordSnapshot(repository, searchSnapshot, { productIds });
await repository.recordDiscoveredProducts(products);
for (const product of products) {
await queue.enqueue({
source: product.source,
url: product.productUrl,
jobType: 'product',
});
}
await queue.markState(job.id, 'parsed');
console.log(`mvideo category parsed: ${products.length} product candidates`);
}
async function processYandexProduct(job: CrawlJobRecord, repository: IngestionRepository, queue: CrawlQueue) {
const snapshot = await fetchRawSnapshot({
source: job.source,
url: job.url,
snapshotDir: ingestionConfig.snapshotDir,
});
const parsed = parseYandexMarketProduct(snapshot.html, job.url);
const normalized = normalizeMarketProduct(parsed);
const stableContentHash = createStableContentHash(normalized);
await recordSnapshot(repository, snapshot, { parsed, normalized }, stableContentHash);
await repository.upsertNormalizedProduct(normalized, stableContentHash, snapshot.fetchedAt);
await queue.markState(job.id, 'needs_review');
console.log(`yandex_market product parsed: ${normalized.brand} ${normalized.canonicalModel}`);
}
async function processMvideoProduct(job: CrawlJobRecord, repository: IngestionRepository, queue: CrawlQueue) {
const productPageSnapshot = await fetchRawSnapshot({
source: job.source,
url: job.url,
snapshotDir: ingestionConfig.snapshotDir,
});
const bffBasePath = extractMvideoBffBasePath(productPageSnapshot.html);
const cookieHeader = buildCookieHeader(productPageSnapshot.setCookieHeaders);
const sourceProductId = productPageSnapshot.url.match(/(\d+)(?:\/?(?:\?.*)?)$/)?.[1];
if (!sourceProductId) {
throw sourceError('parse_failed', `MVideo product id was not found in ${productPageSnapshot.url}`);
}
const detailsUrl = buildMvideoBffUrl(bffBasePath, 'product-details');
detailsUrl.searchParams.set('productId', sourceProductId);
const detailsSnapshot = await fetchRawSnapshot({
source: job.source,
url: detailsUrl.toString(),
snapshotDir: ingestionConfig.snapshotDir,
accept: 'application/json,text/plain,*/*',
...(cookieHeader ? { headers: { cookie: cookieHeader } } : {}),
});
const parsed = parseMvideoProduct(parseJsonSnapshot(detailsSnapshot, 'MVideo product details'));
const normalized = normalizeMarketProduct(parsed);
const stableContentHash = createStableContentHash(normalized);
await recordSnapshot(
repository,
productPageSnapshot,
{ productApiUrl: detailsUrl.toString(), parsed, normalized },
stableContentHash,
);
await recordSnapshot(repository, detailsSnapshot, { parsed, normalized }, stableContentHash);
await repository.upsertNormalizedProduct(normalized, stableContentHash, productPageSnapshot.fetchedAt);
await queue.markState(job.id, 'needs_review');
console.log(`mvideo product parsed: ${normalized.brand} ${normalized.canonicalModel}`);
}
async function processJob(job: CrawlJobRecord, queue: CrawlQueue, repository: IngestionRepository) {
const reserved = await queue.reserveDailyBudget(job.source, job.jobType, dailyBudgetForJob(job.jobType));
if (!reserved) {
await queue.markFailure(job, 'rate_limited', `Daily ${job.jobType} budget is exhausted`);
return;
}
try {
if (job.source === 'yandex_market' && job.jobType === 'category') {
await processYandexCategory(job, repository, queue);
return;
}
if (job.source === 'mvideo' && job.jobType === 'category') {
await processMvideoCategory(job, repository, queue);
return;
}
if (job.source === 'yandex_market' && job.jobType === 'product') {
await processYandexProduct(job, repository, queue);
return;
}
if (job.source === 'mvideo' && job.jobType === 'product') {
await processMvideoProduct(job, repository, queue);
return;
}
await queue.markFailure(job, 'needs_manual_review', `Unsupported job: ${job.source}/${job.jobType}`);
} catch (error) {
const reason = getFetchFailureReason(error) ?? 'parse_failed';
await queue.markFailure(job, reason, (error as Error).message);
if (isStopReason(reason)) {
await queue.pauseSource(job.source, reason, new Date(Date.now() + 24 * 60 * 60 * 1000));
}
}
}
export async function runWorker() {
const pool = createPgPool();
const queue = new CrawlQueue(pool);
const repository = new IngestionRepository(pool);
for (const seed of ingestionConfig.sourceSeeds) {
if (!seed.implemented) {
console.warn(`source seed configured but crawler is not implemented yet: ${seed.source} ${seed.categoryUrl}`);
continue;
}
await queue.enqueue({
source: seed.source,
url: seed.categoryUrl,
jobType: 'category',
});
}
try {
do {
let processed = 0;
for (let index = 0; index < ingestionConfig.workerBatchSize; index += 1) {
const job = await queue.leaseNext();
if (!job) break;
await processJob(job, queue, repository);
processed += 1;
await sleep(randomDelayMs());
}
if (process.env.INGESTION_ONCE === '1') break;
if (processed === 0) await sleep(ingestionConfig.workerIdleMs);
} while (true);
} finally {
await pool.end();
}
}
if (import.meta.url === `file://${process.argv[1]}`) {
await runWorker();
}