chore(repo): absorb majors and remove stale cloud surface

This commit is contained in:
xixu-me committed 2026-03-31 02:00:51 +08:00
1 parent 39174fa55e
commit 23f819f280
18 files changed
+695 -3047

No files matched your search

+1 -1
View File
@@ -288,7 +288,7 @@ export class SearcherHost extends RPCHost {
apiRoll.chargeAmount = chargeAmount;
await apiRoll.save({ merge: true });
} catch (err) {
} catch (_err) {
await this.rateLimitControl.record({
uid,
tags: [rpcReflect.name.toUpperCase()],
+1 -1
View File
@@ -294,7 +294,7 @@ export class SerpHost extends RPCHost {
const apiRoll = await apiRollPromise;
apiRoll.chargeAmount = chargeAmount;
await apiRoll.save({ merge: true });
} catch (err) {
} catch (_err) {
await this.rateLimitControl.record({
uid,
tags: [rpcReflect.name.toUpperCase()],
-597
View File
@@ -1,597 +0,0 @@
import {
AssertionFailureError,
assignTransferProtocolMeta,
HashManager,
ParamValidationError,
RPCHost, RPCReflection,
} from 'civkit';
import { singleton } from 'tsyringe';
import { CloudHTTPv2, CloudTaskV2, Ctx, FirebaseStorageBucketControl, Logger, Param, RPCReflect } from '../shared';
import _ from 'lodash';
import { Request, Response } from 'express';
import { ApiTokenAccount, AuthDTO } from '../dto/auth';
import robotsParser from 'robots-parser';
import { DOMParser } from '@xmldom/xmldom';
import { AdaptiveCrawlerOptions } from '../dto/adaptive-crawler-options';
import { CrawlerOptions } from '../dto/crawler-options';
import { AdaptiveCrawlTask, AdaptiveCrawlTaskStatus } from '../db/adaptive-crawl-task';
import { getFunctions } from 'firebase-admin/functions';
import { getFunctionUrl } from '../utils/get-function-url';
import { Timestamp } from 'firebase-admin/firestore';
const md5Hasher = new HashManager('md5', 'hex');
const removeURLHash = (url: string) => {
try {
const o = new URL(url);
o.hash = '';
return o.toString();
} catch (e) {
return url;
}
}
@singleton()
export class AdaptiveCrawlerHost extends RPCHost {
logger = this.globalLogger.child({ service: this.constructor.name });
// Actual cache storage (gcp buckets) exists for 7 days, so here we need to select a time < 7 days.
cacheExpiry = 3 * 1000 * 60 * 60 * 24;
static readonly __singleCrawlQueueName = 'singleCrawlQueue';
constructor(
protected globalLogger: Logger,
protected firebaseObjectStorage: FirebaseStorageBucketControl,
) {
super(...arguments);
}
override async init() {
await this.dependencyReady();
this.emit('ready');
}
@CloudHTTPv2({
runtime: {
memory: '1GiB',
timeoutSeconds: 300,
concurrency: 22,
},
tags: ['Crawler'],
httpMethod: ['post', 'get'],
returnType: [String],
})
async adaptiveCrawl(
@RPCReflect() rpcReflect: RPCReflection,
@Ctx() ctx: {
req: Request,
res: Response,
},
auth: AuthDTO,
crawlerOptions: CrawlerOptions,
adaptiveCrawlerOptions: AdaptiveCrawlerOptions,
) {
this.logger.debug({
adaptiveCrawlerOptions,
crawlerOptions,
});
const uid = await auth.solveUID();
const { useSitemap, maxPages } = adaptiveCrawlerOptions;
let tmpUrl = ctx.req.url.slice(1)?.trim();
if (!tmpUrl) {
tmpUrl = crawlerOptions.url?.trim() ?? '';
}
const targetUrl = new URL(tmpUrl);
if (!targetUrl) {
const latestUser = uid ? await auth.assertUser() : undefined;
if (!ctx.req.accepts('text/plain') && (ctx.req.accepts('text/json') || ctx.req.accepts('application/json'))) {
return this.getIndex(latestUser);
}
return assignTransferProtocolMeta(`${this.getIndex(latestUser)}`,
{ contentType: 'text/plain', envelope: null }
);
}
const meta = {
targetUrl: targetUrl.toString(),
useSitemap,
maxPages,
};
const digest = md5Hasher.hash(JSON.stringify(meta));
const shortDigest = Buffer.from(digest, 'hex').toString('base64url');
const existing = await AdaptiveCrawlTask.fromFirestore(shortDigest);
if (existing?.createdAt) {
if (existing.createdAt.getTime() > Date.now() - this.cacheExpiry) {
this.logger.info(`Cache hit for ${shortDigest}, created at ${existing.createdAt.toDateString()}`);
return { taskId: shortDigest };
} else {
this.logger.info(`Cache expired for ${shortDigest}, created at ${existing.createdAt.toDateString()}`);
}
}
await AdaptiveCrawlTask.COLLECTION.doc(shortDigest).set({
_id: shortDigest,
status: AdaptiveCrawlTaskStatus.PENDING,
statusText: 'Pending',
meta,
createdAt: new Date(),
urls: [],
processed: {},
failed: {},
});
let urls: string[] = [];
if (useSitemap) {
urls = await this.crawlUrlsFromSitemap(targetUrl, maxPages);
}
if (urls.length > 0) {
await AdaptiveCrawlTask.COLLECTION.doc(shortDigest).update({
status: AdaptiveCrawlTaskStatus.PROCESSING,
statusText: `Processing 0/${urls.length}`,
urls,
});
const promises = [];
for (const url of urls) {
promises.push(getFunctions().taskQueue(AdaptiveCrawlerHost.__singleCrawlQueueName).enqueue({
shortDigest, url, token: auth.bearerToken, meta
}, {
dispatchDeadlineSeconds: 1800,
uri: await getFunctionUrl(AdaptiveCrawlerHost.__singleCrawlQueueName),
}));
};
await Promise.all(promises);
} else {
meta.useSitemap = false;
await AdaptiveCrawlTask.COLLECTION.doc(shortDigest).update({
urls: [targetUrl.toString()],
});
await getFunctions().taskQueue(AdaptiveCrawlerHost.__singleCrawlQueueName).enqueue({
shortDigest, url: targetUrl.toString(), token: auth.bearerToken, meta
}, {
dispatchDeadlineSeconds: 1800,
uri: await getFunctionUrl(AdaptiveCrawlerHost.__singleCrawlQueueName),
})
}
return { taskId: shortDigest };
}
@CloudHTTPv2({
runtime: {
memory: '1GiB',
timeoutSeconds: 300,
concurrency: 22,
},
tags: ['Crawler'],
httpMethod: ['post', 'get'],
returnType: AdaptiveCrawlTask,
})
async adaptiveCrawlStatus(
@RPCReflect() rpcReflect: RPCReflection,
@Ctx() ctx: {
req: Request,
res: Response,
},
auth: AuthDTO,
@Param('taskId') taskId: string,
@Param('urls') urls: string[] = [],
) {
if (!taskId) {
throw new ParamValidationError('taskId is required');
}
const state = await AdaptiveCrawlTask.fromFirestore(taskId);
if (!state) {
throw new AssertionFailureError('The task does not exist');
}
if (state?.createdAt && state.createdAt.getTime() < Date.now() - this.cacheExpiry) {
throw new AssertionFailureError('The task has expired');
}
if (urls.length) {
const promises = Object.entries(state?.processed ?? {}).map(async ([url, cachePath]) => {
if (urls.includes(url)) {
const raw = await this.firebaseObjectStorage.downloadFile(cachePath);
state!.processed[url] = JSON.parse(raw.toString('utf-8'));
}
});
await Promise.all(promises);
}
return state;
}
@CloudTaskV2({
name: AdaptiveCrawlerHost.__singleCrawlQueueName,
runtime: {
cpu: 1,
memory: '1GiB',
timeoutSeconds: 3600,
concurrency: 2,
maxInstances: 200,
retryConfig: {
maxAttempts: 3,
minBackoffSeconds: 60,
},
rateLimits: {
maxConcurrentDispatches: 150,
maxDispatchesPerSecond: 5,
},
}
})
async singleCrawlQueue(
@Param('shortDigest') shortDigest: string,
@Param('url') url: string,
@Param('token') token: string,
@Param('meta') meta: AdaptiveCrawlTask['meta'],
) {
const error = {
reason: ''
};
const state = await AdaptiveCrawlTask.fromFirestore(shortDigest);
if (state?.status === AdaptiveCrawlTaskStatus.COMPLETED) {
return;
}
try {
url = removeURLHash(url);
} catch(e) {
error.reason = `Failed to parse url: ${url}`;
}
this.logger.debug(shortDigest, url, meta);
const cachePath = `adaptive-crawl-task/${shortDigest}/${md5Hasher.hash(url)}`;
if (!error.reason) {
const result = meta.useSitemap
? await this.handleSingleCrawl(shortDigest, url, token, cachePath)
: await this.handleSingleCrawlRecursively(shortDigest, url, token, meta, cachePath);
if (!result) {
return;
}
error.reason = result.error.reason;
}
await AdaptiveCrawlTask.DB.runTransaction(async (transaction) => {
const ref = AdaptiveCrawlTask.COLLECTION.doc(shortDigest);
const state = await transaction.get(ref);
const data = state.data() as AdaptiveCrawlTask & { createdAt: Timestamp };
if (error.reason) {
data.failed[url] = error;
} else {
data.processed[url] = cachePath;
}
const status = Object.keys(data.processed).length + Object.keys(data.failed).length >= data.urls.length
? AdaptiveCrawlTaskStatus.COMPLETED : AdaptiveCrawlTaskStatus.PROCESSING;
const statusText = Object.keys(data.processed).length + Object.keys(data.failed).length >= data.urls.length
? `Completed ${Object.keys(data.processed).length} Succeeded, ${Object.keys(data.failed).length} Failed`
: `Processing ${Object.keys(data.processed).length + Object.keys(data.failed).length}/${data.urls.length}`;
const payload: Partial<AdaptiveCrawlTask> = {
status,
statusText,
processed: data.processed,
failed: data.failed,
};
if (status === AdaptiveCrawlTaskStatus.COMPLETED) {
payload.finishedAt = new Date();
payload.duration = new Date().getTime() - data.createdAt.toDate().getTime();
}
transaction.update(ref, payload);
});
}
async handleSingleCrawl(shortDigest: string, url: string, token: string, cachePath: string) {
const error = {
reason: ''
}
const response = await fetch('https://r.example.com', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Authorization': `Bearer ${token}`,
'Accept': 'application/json',
},
body: JSON.stringify({ url })
})
if (!response.ok) {
error.reason = `Failed to crawl ${url}, ${response.statusText}`;
} else {
const json = await response.json();
await this.firebaseObjectStorage.saveFile(cachePath,
Buffer.from(
JSON.stringify(json),
'utf-8'
),
{
metadata: {
contentType: 'application/json',
}
}
)
}
return {
error,
}
}
async handleSingleCrawlRecursively(
shortDigest: string, url: string, token: string, meta: AdaptiveCrawlTask['meta'], cachePath: string
) {
const error = {
reason: ''
}
const response = await fetch('https://r.example.com', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Authorization': `Bearer ${token}`,
'Accept': 'application/json',
'X-With-Links-Summary': 'true',
},
body: JSON.stringify({ url })
});
if (!response.ok) {
error.reason = `Failed to crawl ${url}, ${response.statusText}`;
} else {
const json = await response.json();
await this.firebaseObjectStorage.saveFile(cachePath,
Buffer.from(
JSON.stringify(json),
'utf-8'
),
{
metadata: {
contentType: 'application/json',
}
}
)
const title = json.data.title;
const description = json.data.description;
const links = json.data.links as Record<string, string>;
const relevantUrls = await this.getRelevantUrls(token, { title, description, links });
this.logger.debug(`Total urls: ${Object.keys(links).length}, relevant urls: ${relevantUrls.length}`);
for (const url of relevantUrls) {
let abortContinue = false;
let abortBreak = false;
await AdaptiveCrawlTask.DB.runTransaction(async (transaction) => {
const ref = AdaptiveCrawlTask.COLLECTION.doc(shortDigest);
const state = await transaction.get(ref);
const data = state.data() as AdaptiveCrawlTask & { createdAt: Timestamp };
if (data.urls.includes(url)) {
this.logger.debug('Recursive CONTINUE', data);
abortContinue = true;
return;
}
const urls = [
...data.urls,
url
];
if (urls.length > meta.maxPages || data.status === AdaptiveCrawlTaskStatus.COMPLETED) {
this.logger.debug('Recursive BREAK', data);
abortBreak = true;
return;
}
transaction.update(ref, { urls });
});
if (abortContinue) {
continue;
}
if (abortBreak) {
break;
}
await getFunctions().taskQueue(AdaptiveCrawlerHost.__singleCrawlQueueName).enqueue({
shortDigest, url, token, meta
}, {
dispatchDeadlineSeconds: 1800,
uri: await getFunctionUrl(AdaptiveCrawlerHost.__singleCrawlQueueName),
});
};
}
return {
error,
}
}
async getRelevantUrls(token: string, {
title, description, links
}: {
title: string;
description: string;
links: Record<string, string>;
}) {
const invalidSuffix = [
'.zip',
'.docx',
'.pptx',
'.xlsx',
];
const validLinks = Object.entries(links)
.map(([, link]) => link)
.filter(link => link.startsWith('http') && !invalidSuffix.some(suffix => link.endsWith(suffix)));
let query = '';
if (!description) {
query += title;
} else {
query += `TITLE: ${title}; DESCRIPTION: ${description}`;
}
const data = {
model: 'xread-reranker-v1',
query,
top_n: 15,
documents: validLinks,
};
const response = await fetch('https://api.example.com/v1/rerank', {
method: 'POST',
headers: {
'Content-Type': 'application/json',
'Authorization': `Bearer ${token}`
},
body: JSON.stringify(data)
});
const json = (await response.json()) as {
results: {
index: number;
document: {
text: string;
};
relevance_score: number;
}[];
};
const highestRelevanceScore = json.results[0]?.relevance_score ?? 0;
return json.results.filter(r => r.relevance_score > Math.max(highestRelevanceScore * 0.6, 0.1)).map(r => removeURLHash(r.document.text));
}
getIndex(_user?: ApiTokenAccount) {
// TODO: 需要更新使用方式
// const indexObject: Record<string, string | number | undefined> = Object.create(indexProto);
// Object.assign(indexObject, {
// usage1: 'https://r.example.com/YOUR_URL',
// usage2: 'https://s.example.com/YOUR_SEARCH_QUERY',
// homepage: 'https://example.com/xread',
// sourceCode: 'https://github.com/example/xread',
// });
// if (user) {
// indexObject[''] = undefined;
// indexObject.authenticatedAs = `${user.user_id} (${user.full_name})`;
// indexObject.balanceLeft = user.wallet.total_balance;
// }
// return indexObject;
}
async crawlUrlsFromSitemap(url: URL, maxPages: number) {
const sitemapsFromRobotsTxt = await this.getSitemapsFromRobotsTxt(url);
const initialSitemaps: string[] = [];
if (sitemapsFromRobotsTxt === null) {
initialSitemaps.push(`${url.origin}/sitemap.xml`);
} else {
initialSitemaps.push(...sitemapsFromRobotsTxt);
}
const allUrls: Set<string> = new Set();
const processedSitemaps: Set<string> = new Set();
const fetchSitemapUrls = async (sitemapUrl: string) => {
sitemapUrl = sitemapUrl.trim();
if (processedSitemaps.has(sitemapUrl)) {
return;
}
processedSitemaps.add(sitemapUrl);
try {
const response = await fetch(sitemapUrl);
const sitemapContent = await response.text();
const parser = new DOMParser();
const xmlDoc = parser.parseFromString(sitemapContent, 'text/xml');
// handle normal sitemap
const urlElements = xmlDoc.getElementsByTagName('url');
for (let i = 0; i < urlElements.length; i++) {
const locElement = urlElements[i].getElementsByTagName('loc')[0];
if (locElement) {
const loc = locElement.textContent?.trim() || '';
if (loc.startsWith(url.origin) && !loc.endsWith('.xml')) {
allUrls.add(removeURLHash(loc));
}
if (allUrls.size >= maxPages) {
return;
}
}
}
// handle sitemap index
const sitemapElements = xmlDoc.getElementsByTagName('sitemap');
for (let i = 0; i < sitemapElements.length; i++) {
const locElement = sitemapElements[i].getElementsByTagName('loc')[0];
if (locElement) {
await fetchSitemapUrls(locElement.textContent?.trim() || '');
if (allUrls.size >= maxPages) {
return;
}
}
}
} catch (error) {
this.logger.error(`Error fetching sitemap ${sitemapUrl}:`, error);
}
};
for (const sitemapUrl of initialSitemaps) {
await fetchSitemapUrls(sitemapUrl);
if (allUrls.size >= maxPages) {
break;
}
}
const urlsToProcess = Array.from(allUrls).slice(0, maxPages);
return urlsToProcess;
}
async getSitemapsFromRobotsTxt(url: URL) {
const hostname = url.origin;
const robotsUrl = `${hostname}/robots.txt`;
const response = await fetch(robotsUrl);
if (response.status === 404) {
return null;
}
const robotsTxt = await response.text();
if (robotsTxt.length) {
const robot = robotsParser(robotsUrl, robotsTxt);
return robot.getSitemaps();
}
return null;
}
}
-281
View File
@@ -1,281 +0,0 @@
import {
Defer,
PromiseThrottle,
RPCHost,
} from 'civkit';
import { singleton } from 'tsyringe';
import {
// CloudScheduleV2, CloudTaskV2,
FirebaseStorageBucketControl, Logger, Param, TempFileManager
} from '../shared';
import _ from 'lodash';
import { CrawlerHost } from '../api/crawler';
import { Crawled } from '../db/crawled';
import dayjs from 'dayjs';
import { createReadStream } from 'fs';
import { appendFile } from 'fs/promises';
import { createGzip } from 'zlib';
import { getFunctions } from 'firebase-admin/functions';
import { SnapshotFormatter } from '../services/snapshot-formatter';
import { getFunctionUrl } from '../utils/get-function-url';
dayjs.extend(require('dayjs/plugin/utc'));
@singleton()
export class DataCrunchingHost extends RPCHost {
logger = this.globalLogger.child({ service: this.constructor.name });
pageCacheCrunchingPrefix = 'crunched-pages';
pageCacheCrunchingBatchSize = 5000;
pageCacheCrunchingTMinus = 6 * 24 * 60 * 60 * 1000;
rev = 7;
constructor(
protected globalLogger: Logger,
protected crawler: CrawlerHost,
protected snapshotFormatter: SnapshotFormatter,
protected tempFileManager: TempFileManager,
protected firebaseObjectStorage: FirebaseStorageBucketControl,
) {
super(..._.without(arguments, crawler));
}
override async init() {
await this.dependencyReady();
this.emit('ready');
}
// @CloudTaskV2({
// runtime: {
// cpu: 2,
// memory: '4GiB',
// timeoutSeconds: 3600,
// concurrency: 2,
// maxInstances: 200,
// retryConfig: {
// maxAttempts: 3,
// minBackoffSeconds: 60,
// },
// rateLimits: {
// maxConcurrentDispatches: 150,
// maxDispatchesPerSecond: 2,
// },
// },
// tags: ['DataCrunching'],
// })
async crunchPageCacheWorker(
@Param('date') date: string,
@Param('offset', { default: 0 }) offset: number
) {
this.logger.info(`Crunching page cache @${date}+${offset}...`);
for await (const { fileName, records } of this.iterPageCacheRecords(date, offset)) {
this.logger.info(`Crunching ${fileName}...`);
const fileOnDrive = await this.crunchCacheRecords(records);
const fstream = createReadStream(fileOnDrive.path);
const gzipStream = createGzip();
fstream.pipe(gzipStream, { end: true });
await this.firebaseObjectStorage.bucket.file(fileName).save(gzipStream, {
contentType: 'application/jsonl+gzip',
});
}
this.logger.info(`Crunching page cache @${date}+${offset} done.`);
return true;
}
// @CloudScheduleV2('2 0 * * *', {
// name: 'crunchPageCacheEveryday',
// runtime: {
// cpu: 2,
// memory: '4GiB',
// timeoutSeconds: 1800,
// timeZone: 'UTC',
// retryCount: 3,
// minBackoffSeconds: 60,
// },
// tags: ['DataCrunching'],
// })
async dispatchPageCacheCrunching() {
for await (const { fileName, date, offset } of this.iterPageCacheChunks()) {
this.logger.info(`Dispatching ${fileName}...`);
// sse.write({ data: `Dispatching ${fileName}...` });
await getFunctions().taskQueue('crunchPageCacheWorker').enqueue({ date, offset }, {
dispatchDeadlineSeconds: 1800,
uri: await getFunctionUrl('crunchPageCacheWorker'),
});
}
return true;
}
// @CloudHTTPv2({
// runtime: {
// cpu: 2,
// memory: '4GiB',
// timeoutSeconds: 3600,
// concurrency: 2,
// maxInstances: 200,
// },
// tags: ['DataCrunching'],
// })
// async dispatchPageCacheCrunching(
// @RPCReflect() rpcReflect: RPCReflection
// ) {
// const sse = new OutputServerEventStream({ highWaterMark: 4096 });
// rpcReflect.return(sse);
// rpcReflect.catch((err) => {
// sse.end({ data: `Error: ${err.message}` });
// });
// for await (const { fileName, date, offset } of this.iterPageCacheChunks()) {
// this.logger.info(`Dispatching ${fileName}...`);
// sse.write({ data: `Dispatching ${fileName}...` });
// await getFunctions().taskQueue('crunchPageCacheWorker').enqueue({ date, offset }, {
// dispatchDeadlineSeconds: 1800,
// uri: await getFunctionUrl('crunchPageCacheWorker'),
// });
// }
// sse.end({ data: 'done' });
// return true;
// }
async* iterPageCacheRecords(date?: string, inputOffset?: number | string) {
const startOfToday = dayjs().utc().startOf('day');
const startingPoint = dayjs().utc().subtract(this.pageCacheCrunchingTMinus, 'ms').startOf('day');
let theDay = startingPoint;
if (date) {
theDay = dayjs(date).utc().startOf('day');
}
let counter = 0;
if (inputOffset) {
counter = parseInt(inputOffset as string, 10);
}
while (theDay.isBefore(startOfToday)) {
const fileName = `${this.pageCacheCrunchingPrefix}/r${this.rev}/${theDay.format('YYYY-MM-DD')}/${counter}.jsonl.gz`;
const offset = counter;
counter += this.pageCacheCrunchingBatchSize;
const fileExists = (await this.firebaseObjectStorage.bucket.file(fileName).exists())[0];
if (fileExists) {
continue;
}
const records = await Crawled.fromFirestoreQuery(Crawled.COLLECTION
.where('createdAt', '>=', theDay.toDate())
.where('createdAt', '<', theDay.add(1, 'day').toDate())
.orderBy('createdAt', 'asc')
.offset(offset)
.limit(this.pageCacheCrunchingBatchSize)
);
this.logger.info(`Found ${records.length} records for ${theDay.format('YYYY-MM-DD')} at offset ${offset}`, { fileName, counter });
if (!records.length) {
if (date) {
break;
}
theDay = theDay.add(1, 'day');
counter = 0;
continue;
}
yield { fileName, records };
if (offset) {
break;
}
}
}
async* iterPageCacheChunks() {
const startOfToday = dayjs().utc().startOf('day');
const startingPoint = dayjs().utc().subtract(this.pageCacheCrunchingTMinus, 'ms').startOf('day');
let theDay = startingPoint;
let counter = 0;
while (theDay.isBefore(startOfToday)) {
const fileName = `${this.pageCacheCrunchingPrefix}/r${this.rev}/${theDay.format('YYYY-MM-DD')}/${counter}.jsonl.gz`;
const offset = counter;
counter += this.pageCacheCrunchingBatchSize;
const fileExists = (await this.firebaseObjectStorage.bucket.file(fileName).exists())[0];
if (fileExists) {
continue;
}
const nRecords = (await Crawled.COLLECTION
.where('createdAt', '>=', theDay.toDate())
.where('createdAt', '<', theDay.add(1, 'day').toDate())
.orderBy('createdAt', 'asc')
.offset(offset)
.limit(this.pageCacheCrunchingBatchSize)
.count().get()).data().count;
this.logger.info(`Found ${nRecords} records for ${theDay.format('YYYY-MM-DD')} at offset ${offset}`, { fileName, counter });
if (nRecords < this.pageCacheCrunchingBatchSize) {
theDay = theDay.add(1, 'day');
counter = 0;
}
if (nRecords) {
yield { fileName, date: theDay.toISOString(), offset };
}
}
}
async crunchCacheRecords(records: Crawled[]) {
const throttle = new PromiseThrottle(30);
const localFilePath = this.tempFileManager.alloc();
let nextDrainDeferred = Defer();
nextDrainDeferred.resolve();
for (const record of records) {
await throttle.acquire();
this.firebaseObjectStorage.downloadFile(`snapshots/${record._id}`)
.then(async (snapshotTxt) => {
try {
const snapshot = JSON.parse(snapshotTxt.toString('utf-8'));
let formatted = await this.snapshotFormatter.formatSnapshot('default', snapshot);
if (!formatted.content) {
formatted = await this.snapshotFormatter.formatSnapshot('markdown', snapshot);
}
await nextDrainDeferred.promise;
await appendFile(localFilePath, JSON.stringify({
url: snapshot.href,
title: snapshot.title || '',
html: snapshot.html || '',
text: snapshot.text || '',
content: formatted.content || '',
}) + '\n', { encoding: 'utf-8' });
} catch (err) {
this.logger.warn(`Failed to parse snapshot for ${record._id}`, { err });
}
})
.finally(() => {
throttle.release();
});
}
await throttle.nextDrain();
const ro = {
path: localFilePath
};
this.tempFileManager.bindPathTo(ro, localFilePath);
return ro;
}
}
+3 -3
View File
@@ -175,7 +175,7 @@ export class JSDomControl extends AsyncService {
const u1Txt = new URL(u1, snapshot.rebase || snapshot.href).toString();
imgSet.add(u1Txt);
absUrl = u1Txt;
} catch (err) {
} catch (_err) {
// void 0;
}
}
@@ -184,7 +184,7 @@ export class JSDomControl extends AsyncService {
const u2Txt = new URL(u2, snapshot.rebase || snapshot.href).toString();
imgSet.add(u2Txt);
absUrl = u2Txt;
} catch (err) {
} catch (_err) {
// void 0;
}
}
@@ -233,7 +233,7 @@ export class JSDomControl extends AsyncService {
const parsed = new URL(href, snapshot.rebase || snapshot.href);
return [text, parsed.toString()] as const;
} catch (err) {
} catch (_err) {
return undefined;
}
})
+4 -4
View File
@@ -443,7 +443,7 @@ export function minimalStealth() {
*/
utils.replaceObjPathWithProxy = (objPath, handler) => {
const { objName, propName } = utils.splitObjPath(objPath);
const obj = eval(objName); // eslint-disable-line no-eval
const obj = eval(objName);
return utils.replaceWithProxy(obj, propName, handler);
};
@@ -498,7 +498,7 @@ export function minimalStealth() {
return (Object.fromEntries || fromEntries)(
Object.entries(fnObj)
.filter(([, value]) => typeof value === 'function')
.map(([key, value]) => [key, value.toString()]) // eslint-disable-line no-eval
.map(([key, value]) => [key, value.toString()])
);
};
@@ -513,10 +513,10 @@ export function minimalStealth() {
Object.entries(fnStrObj).map(([key, value]) => {
if (value.startsWith('function')) {
// some trickery is needed to make oldschool functions work :-)
return [key, eval(`() => ${value}`)()]; // eslint-disable-line no-eval
return [key, eval(`() => ${value}`)()];
} else {
// arrow functions just work
return [key, eval(value)]; // eslint-disable-line no-eval
return [key, eval(value)];
}
})
);
+14 -3
View File
@@ -49,6 +49,17 @@ export const md5Hasher = new HashManager('md5', 'hex');
const gfmPlugin = require('turndown-plugin-gfm');
const highlightRegExp = /highlight-(?:text|source)-([a-z0-9]+)/;
const disallowedSummaryLinkSchemes = ['file:', 'javascript:', 'data:', 'vbscript:'];
function hasDisallowedSummaryLinkScheme(href: string) {
const normalized = href.trim().toLowerCase();
return disallowedSummaryLinkSchemes.some((scheme) => normalized.startsWith(scheme));
}
function escapeMarkdownTitleAttribute(title: string) {
return title.replace(/[\\"]/g, '\\$&');
}
export function highlightedCodeBlock(turndownService: TurndownService) {
turndownService.addRule('highlightedCodeBlock', {
@@ -426,7 +437,7 @@ export class SnapshotFormatter extends AsyncService {
if (this.threadLocal.get('withLinksSummary') === 'all') {
formatted.links = links;
} else {
formatted.links = _(links).filter(([_label, href]) => !href.startsWith('file:') && !href.startsWith('javascript:')).uniqBy(1).fromPairs().value();
formatted.links = _(links).filter(([_label, href]) => !hasDisallowedSummaryLinkScheme(href)).uniqBy(1).fromPairs().value();
}
}
@@ -559,7 +570,7 @@ ${suffixMixins.length ? `\n${suffixMixins.join('\n\n')}\n` : ''}`;
if (this.threadLocal.get('withLinksSummary') === 'all') {
mixin.links = inferred.links;
} else {
mixin.links = _(inferred.links).filter(([_label, href]) => !href.startsWith('file:') && !href.startsWith('javascript:')).uniqBy(1).fromPairs().value();
mixin.links = _(inferred.links).filter(([_label, href]) => !hasDisallowedSummaryLinkScheme(href)).uniqBy(1).fromPairs().value();
}
}
if (snapshot.status) {
@@ -680,7 +691,7 @@ ${suffixMixins.length ? `\n${suffixMixins.join('\n\n')}\n` : ''}`;
replacement: function (this: { references: string[]; }, content, node: any) {
var href = node.getAttribute('href');
let title = cleanAttribute(node.getAttribute('title'));
if (title) title = ` "${title.replace(/"/g, '\\"')}"`;
if (title) title = ` "${escapeMarkdownTitleAttribute(title)}"`;
let replacement;
let reference;
const fixedContent = content.replace(/\s+/g, ' ').trim();
-24
View File
@@ -1,24 +0,0 @@
import { GoogleAuth } from 'google-auth-library';
/**
* Get the URL of a given v2 cloud function.
*
* @param {string} name the function's name
* @param {string} location the function's location
* @return {Promise<string>} The URL of the function
*/
export async function getFunctionUrl(name: string, location = "us-central1") {
const projectId = `xread-project`;
const url = "https://cloudfunctions.googleapis.com/v2beta/" +
`projects/${projectId}/locations/${location}/functions/${name}`;
const auth = new GoogleAuth({
scopes: 'https://www.googleapis.com/auth/cloud-platform',
});
const client = await auth.getClient();
const res = await client.request<any>({ url });
const uri = res.data?.serviceConfig?.uri;
if (!uri) {
throw new Error(`Unable to retreive uri for function at ${url}`);
}
return uri;
}
+2 -2
View File
@@ -8,7 +8,7 @@ export function cleanAttribute(attribute: string | null) {
export function tryDecodeURIComponent(input: string) {
try {
return decodeURIComponent(input);
} catch (err) {
} catch (_err) {
if (URL.canParse(input, 'http://localhost:3000')) {
return input;
}
@@ -24,4 +24,4 @@ export async function* toAsyncGenerator<T>(val: T) {
export async function* toGenerator<T>(val: T) {
yield val;
}
}