migrate article embeddings to support multi-model architecture and enhance data integrity
This commit is contained in:
+288
-104
@@ -1,14 +1,37 @@
|
||||
const { extractFromHtml } = require('@extractus/article-extractor');
|
||||
const sharp = require('sharp');
|
||||
const db = require('./db');
|
||||
const config = require('./config');
|
||||
const { generateAndStoreEmbedding } = require('./embeddings');
|
||||
const { fetchWithPolicy } = require('./http');
|
||||
const { getSharedBrowserSession } = require('./sources/browserCrawler');
|
||||
const { extractFromHtml } = require("@extractus/article-extractor");
|
||||
const sharp = require("sharp");
|
||||
const db = require("./db");
|
||||
const config = require("./config");
|
||||
const { generateAndStoreEmbedding } = require("./embeddings");
|
||||
const { fetchWithPolicy } = require("./http");
|
||||
const { getSharedBrowserSession } = require("./sources/browserCrawler");
|
||||
const { validateExtractedArticle } = require("./contentValidation");
|
||||
const {
|
||||
getEffectivePolicy,
|
||||
recordPlainSuccess,
|
||||
recordPlainFailure,
|
||||
recordBrowserSuccess,
|
||||
recordBrowserFailure,
|
||||
} = require("./domainPolicy");
|
||||
|
||||
|
||||
const MAX_PLAIN_HTML_LENGTH = 1_500_000;
|
||||
const PLAIN_FETCH_TIMEOUT = 12000;
|
||||
const BROWSER_FETCH_TIMEOUT = 20000;
|
||||
|
||||
// retry windows for failures that look transient (validation rejected the
|
||||
// page, fetch timed out). genuinely terminal failures (404, dead url) get
|
||||
// a hard cap on attempt count instead
|
||||
const VALIDATION_RETRY_AFTER_MS = 24 * 60 * 60 * 1000;
|
||||
const TRANSIENT_RETRY_AFTER_MS = 6 * 60 * 60 * 1000;
|
||||
const MAX_TERMINAL_ATTEMPTS = 3;
|
||||
|
||||
|
||||
const updateArticleAssets = db.prepare(`
|
||||
UPDATE articles
|
||||
SET content = ?, image = ?, content_status = 'ready', content_error = NULL, content_attempted_at = ?
|
||||
SET content = ?, image = ?, content_status = 'ready', content_error = NULL,
|
||||
content_attempted_at = ?, content_attempt_count = content_attempt_count + 1,
|
||||
content_retry_after = NULL
|
||||
WHERE id = ?
|
||||
`);
|
||||
const updateArticleTitleDescription = db.prepare(`
|
||||
@@ -18,19 +41,25 @@ const updateArticleTitleDescription = db.prepare(`
|
||||
`);
|
||||
const markContentSkipped = db.prepare(`
|
||||
UPDATE articles
|
||||
SET content_status = 'skipped', content_error = ?, content_attempted_at = ?
|
||||
SET content_status = 'skipped', content_error = ?, content_attempted_at = ?,
|
||||
content_attempt_count = content_attempt_count + 1, content_retry_after = NULL
|
||||
WHERE id = ?
|
||||
`);
|
||||
const markContentFailed = db.prepare(`
|
||||
UPDATE articles
|
||||
SET content_status = 'failed', content_error = ?, content_attempted_at = ?
|
||||
SET content_status = 'failed', content_error = ?, content_attempted_at = ?,
|
||||
content_attempt_count = content_attempt_count + 1, content_retry_after = NULL
|
||||
WHERE id = ?
|
||||
`);
|
||||
const markContentPending = db.prepare(`
|
||||
UPDATE articles
|
||||
SET content_status = NULL, content_error = NULL, content_attempted_at = ?
|
||||
SET content_status = 'pending', content_error = ?, content_attempted_at = ?,
|
||||
content_attempt_count = content_attempt_count + 1, content_retry_after = ?
|
||||
WHERE id = ?
|
||||
`);
|
||||
|
||||
// round-robin pull of articles needing content. respects content_retry_after so
|
||||
// a freshly-rejected article doesnt get retried in the next loop iteration
|
||||
const selectRoundRobinArticlesMissingContent = db.prepare(`
|
||||
SELECT id, url, title, description
|
||||
FROM (
|
||||
@@ -39,21 +68,18 @@ const selectRoundRobinArticlesMissingContent = db.prepare(`
|
||||
FROM articles
|
||||
WHERE (content IS NULL OR TRIM(content) = '')
|
||||
AND (content_status IS NULL OR content_status = 'pending')
|
||||
AND (content_retry_after IS NULL OR content_retry_after <= datetime('now'))
|
||||
)
|
||||
WHERE rn <= ?
|
||||
ORDER BY rn, source
|
||||
`);
|
||||
|
||||
const loggedBlockedDomains = new Set();
|
||||
let contentBackfillRunning = false;
|
||||
const selectAttemptCount = db.prepare(`
|
||||
SELECT content_attempt_count AS attempts FROM articles WHERE id = ?
|
||||
`);
|
||||
|
||||
function getHostname(url) {
|
||||
try {
|
||||
return new URL(url).hostname.toLowerCase();
|
||||
} catch {
|
||||
return '';
|
||||
}
|
||||
}
|
||||
|
||||
let contentBackfillRunning = false;
|
||||
|
||||
|
||||
function getErrorStatus(error) {
|
||||
@@ -61,38 +87,28 @@ function getErrorStatus(error) {
|
||||
return error.status;
|
||||
}
|
||||
|
||||
const match = String(error && error.message || '').match(/\b(401|403|404|408|429|5\d\d)\b/);
|
||||
const match = String((error && error.message) || "").match(/\b(401|403|404|408|429|5\d\d)\b/);
|
||||
return match ? Number(match[1]) : null;
|
||||
}
|
||||
|
||||
function getErrorMessage(error, fallback) {
|
||||
const message = String(error && error.message || fallback || '').trim();
|
||||
const message = String((error && error.message) || fallback || "").trim();
|
||||
return message ? message.slice(0, 500) : null;
|
||||
}
|
||||
|
||||
function markArticleStatus(statement, id, message) {
|
||||
const attemptedAt = new Date().toISOString();
|
||||
const parameterCount = statement.source.split('?').length - 1;
|
||||
|
||||
if (parameterCount === 3) {
|
||||
statement.run(message, attemptedAt, id);
|
||||
return;
|
||||
}
|
||||
|
||||
if (parameterCount === 2) {
|
||||
statement.run(attemptedAt, id);
|
||||
return;
|
||||
}
|
||||
|
||||
throw new Error(`Unexpected content status statement parameter count: ${parameterCount}`);
|
||||
function nowIso() {
|
||||
return new Date().toISOString();
|
||||
}
|
||||
|
||||
function futureIso(ms) {
|
||||
return new Date(Date.now() + ms).toISOString();
|
||||
}
|
||||
|
||||
|
||||
async function fetchCompressedImage(url) {
|
||||
const response = await fetchWithPolicy(url, {
|
||||
retries: 1,
|
||||
headers: {
|
||||
Accept: 'image/*',
|
||||
},
|
||||
headers: { Accept: "image/*" },
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
@@ -101,90 +117,255 @@ async function fetchCompressedImage(url) {
|
||||
throw error;
|
||||
}
|
||||
|
||||
const contentType = String(response.headers.get('content-type') || '').toLowerCase();
|
||||
if (!contentType.startsWith('image/')) {
|
||||
throw new Error(`image request returned ${contentType || 'unknown content-type'}`);
|
||||
const contentType = String(response.headers.get("content-type") || "").toLowerCase();
|
||||
if (!contentType.startsWith("image/")) {
|
||||
throw new Error(`image request returned ${contentType || "unknown content-type"}`);
|
||||
}
|
||||
|
||||
const input = Buffer.from(await response.arrayBuffer());
|
||||
if (input.length === 0) {
|
||||
throw new Error('image request returned an empty body');
|
||||
throw new Error("image request returned an empty body");
|
||||
}
|
||||
|
||||
const output = await sharp(input)
|
||||
.rotate()
|
||||
.resize({ width: 320, height: 320, fit: 'inside', withoutEnlargement: true })
|
||||
.resize({ width: 320, height: 320, fit: "inside", withoutEnlargement: true })
|
||||
.webp({ quality: 25 })
|
||||
.toBuffer();
|
||||
|
||||
return output.toString('base64');
|
||||
return output.toString("base64");
|
||||
}
|
||||
|
||||
|
||||
// plain http fetch — no js execution. fast, low memory, but fails on
|
||||
// js-rendered sites and gets blocked by cloudflare more often
|
||||
async function fetchPlainHtml(url) {
|
||||
const response = await fetchWithPolicy(url, {
|
||||
timeout: PLAIN_FETCH_TIMEOUT,
|
||||
retries: 1,
|
||||
});
|
||||
|
||||
if (!response.ok) {
|
||||
const error = new Error(`plain fetch returned ${response.status}`);
|
||||
error.status = response.status;
|
||||
throw error;
|
||||
}
|
||||
|
||||
const contentType = String(response.headers.get("content-type") || "").toLowerCase();
|
||||
if (contentType && !contentType.includes("html") && !contentType.includes("xml")) {
|
||||
throw new Error(`plain fetch returned non-html content-type: ${contentType}`);
|
||||
}
|
||||
|
||||
const text = await response.text();
|
||||
return {
|
||||
html: text.slice(0, MAX_PLAIN_HTML_LENGTH),
|
||||
finalUrl: response.url || url,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
async function fetchBrowserHtml(url) {
|
||||
const maxConcurrentPages = Number(config.browser?.maxConcurrentPages) || 25;
|
||||
const session = await getSharedBrowserSession({
|
||||
requestTimeout: BROWSER_FETCH_TIMEOUT,
|
||||
maxConcurrentPages,
|
||||
});
|
||||
|
||||
const html = await session.fetchRenderedHtml(url, { timeout: BROWSER_FETCH_TIMEOUT });
|
||||
return { html, finalUrl: url };
|
||||
}
|
||||
|
||||
|
||||
function stripHtmlContent(value) {
|
||||
if (typeof value !== "string") return null;
|
||||
const stripped = value.replace(/<[^>]+>/g, " ").replace(/\s+/g, " ").trim();
|
||||
return stripped || null;
|
||||
}
|
||||
|
||||
|
||||
// runs fetch → extract → validate. returns { ok, article, html, finalUrl, reason }
|
||||
// where article has been post-processed (content stripped of html). on failure,
|
||||
// reason explains what tripped — used both for logging and for the per-domain
|
||||
// policy update
|
||||
async function attemptFetch(url, fetcher) {
|
||||
let html;
|
||||
let finalUrl;
|
||||
try {
|
||||
const result = await fetcher(url);
|
||||
html = result.html;
|
||||
finalUrl = result.finalUrl;
|
||||
} catch (error) {
|
||||
return { ok: false, reason: `fetch-error:${error.message || "unknown"}`, error };
|
||||
}
|
||||
|
||||
if (!html) {
|
||||
return { ok: false, reason: "empty-html" };
|
||||
}
|
||||
|
||||
let extracted;
|
||||
try {
|
||||
extracted = await extractFromHtml(html, finalUrl || url);
|
||||
} catch (error) {
|
||||
return { ok: false, reason: `extractor-error:${error.message || "unknown"}` };
|
||||
}
|
||||
|
||||
if (extracted) {
|
||||
extracted = {
|
||||
...extracted,
|
||||
content: stripHtmlContent(extracted.content),
|
||||
};
|
||||
}
|
||||
|
||||
const validation = validateExtractedArticle({ article: extracted, html, finalUrl });
|
||||
if (!validation.ok) {
|
||||
return { ok: false, reason: validation.reason, retryable: validation.retryable, html, finalUrl };
|
||||
}
|
||||
|
||||
return { ok: true, article: extracted, html, finalUrl };
|
||||
}
|
||||
|
||||
|
||||
function getAttemptCount(id) {
|
||||
const row = selectAttemptCount.get(id);
|
||||
return row ? row.attempts || 0 : 0;
|
||||
}
|
||||
|
||||
|
||||
async function fetchAndStoreContent(id, url, storedTitle, storedDescription) {
|
||||
const policy = getEffectivePolicy(url);
|
||||
|
||||
// domains we know are blocked — skip the fetch entirely until ttl expires.
|
||||
// the row stays pending so it'll get picked up after the policy resets
|
||||
if (policy.policy === "blocked") {
|
||||
markContentPending.run(
|
||||
`domain blocked by policy`,
|
||||
nowIso(),
|
||||
futureIso(TRANSIENT_RETRY_AFTER_MS),
|
||||
id
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
const tryPlainFirst = policy.policy === "auto" || policy.policy === "plain_only";
|
||||
let plainResult = null;
|
||||
let browserResult = null;
|
||||
|
||||
|
||||
if (tryPlainFirst) {
|
||||
plainResult = await attemptFetch(url, fetchPlainHtml);
|
||||
|
||||
if (plainResult.ok) {
|
||||
recordPlainSuccess(url);
|
||||
await commitArticle(id, url, plainResult, storedTitle, storedDescription);
|
||||
return;
|
||||
}
|
||||
|
||||
recordPlainFailure(url);
|
||||
|
||||
// hard 4xx (other than 408/429) on plain — domain might serve the same to
|
||||
// browser, but try anyway since it's cheap once the policy hasnt flipped yet.
|
||||
// 408/429/5xx defer for retry
|
||||
const status = plainResult.error && getErrorStatus(plainResult.error);
|
||||
if (status === 408 || status === 429 || (status && status >= 500)) {
|
||||
markContentPending.run(
|
||||
`plain ${status}`,
|
||||
nowIso(),
|
||||
futureIso(TRANSIENT_RETRY_AFTER_MS),
|
||||
id
|
||||
);
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
// policy.policy === "plain_only" means we just tried plain and failed —
|
||||
// dont escalate to browser, the operator (or earlier domain memory) said no
|
||||
if (policy.policy === "plain_only") {
|
||||
recordValidationFailure(id, plainResult);
|
||||
return;
|
||||
}
|
||||
|
||||
|
||||
browserResult = await attemptFetch(url, fetchBrowserHtml);
|
||||
|
||||
if (browserResult.ok) {
|
||||
recordBrowserSuccess(url);
|
||||
await commitArticle(id, url, browserResult, storedTitle, storedDescription);
|
||||
return;
|
||||
}
|
||||
|
||||
recordBrowserFailure(url);
|
||||
|
||||
const browserStatus = browserResult.error && getErrorStatus(browserResult.error);
|
||||
if (browserStatus === 408 || browserStatus === 429 || (browserStatus && browserStatus >= 500)) {
|
||||
markContentPending.run(
|
||||
`browser ${browserStatus}`,
|
||||
nowIso(),
|
||||
futureIso(TRANSIENT_RETRY_AFTER_MS),
|
||||
id
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
// both paths exhausted (or browser-only path failed). decide between
|
||||
// pending-with-retry and terminal failed based on attempt count and
|
||||
// whether the validator thought it was retryable
|
||||
recordValidationFailure(id, browserResult);
|
||||
}
|
||||
|
||||
|
||||
function recordValidationFailure(id, result) {
|
||||
const reason = result?.reason || "unknown";
|
||||
const retryable = result?.retryable !== false;
|
||||
const attempts = getAttemptCount(id);
|
||||
|
||||
// hard fetch errors with no retryable signal — terminal after a few tries
|
||||
if (!retryable || attempts + 1 >= MAX_TERMINAL_ATTEMPTS) {
|
||||
markContentFailed.run(reason, nowIso(), id);
|
||||
return;
|
||||
}
|
||||
|
||||
markContentPending.run(reason, nowIso(), futureIso(VALIDATION_RETRY_AFTER_MS), id);
|
||||
}
|
||||
|
||||
|
||||
async function commitArticle(id, url, result, storedTitle, storedDescription) {
|
||||
const { article, finalUrl } = result;
|
||||
const content = article.content || null;
|
||||
|
||||
// if stored title looks like a raw url, replace with extracted one
|
||||
const titleLooksLikeUrl = storedTitle && /^https?:\/\//i.test(storedTitle.trim());
|
||||
if (titleLooksLikeUrl) {
|
||||
const scrapedTitle = typeof article.title === "string" ? article.title.trim() : null;
|
||||
const scrapedDescription = typeof article.description === "string" ? article.description.trim() : null;
|
||||
if (scrapedTitle) {
|
||||
updateArticleTitleDescription.run(scrapedTitle, scrapedDescription || storedDescription || null, id);
|
||||
}
|
||||
}
|
||||
|
||||
let image = null;
|
||||
if (article.image) {
|
||||
try {
|
||||
image = await fetchCompressedImage(article.image);
|
||||
} catch (error) {
|
||||
const status = getErrorStatus(error);
|
||||
if (status === 401 || status === 403 || status === 404 || status === 429) {
|
||||
console.warn(`image fetch skipped for ${url}: upstream returned ${status}`);
|
||||
} else {
|
||||
console.error(`image fetch failed for ${url}:`, error.message || error);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
updateArticleAssets.run(content, image, nowIso(), id);
|
||||
|
||||
try {
|
||||
const maxConcurrentPages = Number(config.browser?.maxConcurrentPages) || 25;
|
||||
const browserSession = await getSharedBrowserSession({ requestTimeout: 20000, maxConcurrentPages });
|
||||
const html = await browserSession.fetchRenderedHtml(url, { timeout: 20000 });
|
||||
const article = await extractFromHtml(html, url);
|
||||
if (!article) {
|
||||
markArticleStatus(markContentSkipped, id, 'extractor returned no article');
|
||||
return;
|
||||
}
|
||||
|
||||
const content = typeof article.content === 'string'
|
||||
? article.content.replace(/<[^>]+>/g, ' ').replace(/\s+/g, ' ').trim() || null
|
||||
: null;
|
||||
|
||||
// if stored title looks like a raw URL, try to replace with scraped title
|
||||
const titleLooksLikeUrl = storedTitle && /^https?:\/\//i.test(storedTitle.trim());
|
||||
if (titleLooksLikeUrl) {
|
||||
const scrapedTitle = typeof article.title === 'string' ? article.title.trim() : null;
|
||||
const scrapedDescription = typeof article.description === 'string' ? article.description.trim() : null;
|
||||
if (scrapedTitle) {
|
||||
updateArticleTitleDescription.run(scrapedTitle, scrapedDescription || storedDescription || null, id);
|
||||
}
|
||||
}
|
||||
|
||||
let image = null;
|
||||
if (article.image) {
|
||||
try {
|
||||
image = await fetchCompressedImage(article.image);
|
||||
} catch (error) {
|
||||
const status = getErrorStatus(error);
|
||||
if (status === 401 || status === 403 || status === 404 || status === 429) {
|
||||
console.warn(`image fetch skipped for ${url}: upstream returned ${status}`);
|
||||
} else {
|
||||
console.error(`image fetch failed for ${url}:`, error);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if (!content && !image) {
|
||||
markArticleStatus(markContentSkipped, id, 'article had no extractable content or image');
|
||||
return;
|
||||
}
|
||||
|
||||
updateArticleAssets.run(content, image, new Date().toISOString(), id);
|
||||
await generateAndStoreEmbedding(id);
|
||||
} catch (error) {
|
||||
const status = getErrorStatus(error);
|
||||
if (status === 401 || status === 403 || status === 404) {
|
||||
console.warn(`content fetch skipped for ${url}: upstream returned ${status}`);
|
||||
markArticleStatus(markContentSkipped, id, `upstream returned ${status}`);
|
||||
return;
|
||||
}
|
||||
|
||||
if (status === 408 || status === 429 || (status && status >= 500)) {
|
||||
console.warn(`content fetch deferred for ${url}: upstream returned ${status}`);
|
||||
markArticleStatus(markContentPending, id, null);
|
||||
return;
|
||||
}
|
||||
|
||||
markArticleStatus(markContentFailed, id, getErrorMessage(error, 'content fetch failed'));
|
||||
console.error(`content fetch failed for ${url}:`, error);
|
||||
console.error(`embedding failed for article ${id}:`, error.message || error);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
async function backfillMissingContent(perSource = 50, concurrency = 5) {
|
||||
if (contentBackfillRunning) {
|
||||
return;
|
||||
@@ -204,15 +385,18 @@ async function backfillMissingContent(perSource = 50, concurrency = 5) {
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
function hasPendingContent() {
|
||||
return Boolean(db.prepare(`
|
||||
SELECT 1 FROM articles
|
||||
WHERE (content IS NULL OR TRIM(content) = '')
|
||||
AND (content_status IS NULL OR content_status = 'pending')
|
||||
AND (content_retry_after IS NULL OR content_retry_after <= datetime('now'))
|
||||
LIMIT 1
|
||||
`).get());
|
||||
}
|
||||
|
||||
|
||||
module.exports = {
|
||||
fetchAndStoreContent,
|
||||
backfillMissingContent,
|
||||
|
||||
Reference in New Issue
Block a user