|
|
import type { BrowserContext, Page } from "playwright";
|
|
|
import { QUOTE_SESSION_TTL_SECONDS } from "@/lib/constants/quote-session";
|
|
|
import { getRedis } from "@/lib/redis";
|
|
|
import { RpaError } from "@/modules/rpa/errors";
|
|
|
import {
|
|
|
dismissCookieConsent,
|
|
|
scrollQuoteWidgetIntoView,
|
|
|
} from "@/workers/rpa/page-prep";
|
|
|
|
|
|
type ParkedEntry = {
|
|
|
sessionId: string;
|
|
|
page: Page;
|
|
|
context: BrowserContext;
|
|
|
parkedAt: number;
|
|
|
};
|
|
|
|
|
|
const pool = new Map<string, ParkedEntry>();
|
|
|
|
|
|
const PARKED_REDIS_PREFIX = "quote_session:parked:";
|
|
|
|
|
|
function ttlMs(): number {
|
|
|
return QUOTE_SESSION_TTL_SECONDS * 1_000;
|
|
|
}
|
|
|
|
|
|
function isExpired(entry: ParkedEntry): boolean {
|
|
|
return Date.now() - entry.parkedAt > ttlMs();
|
|
|
}
|
|
|
|
|
|
/** 仅 rpa-worker 进程可驻留 page(spawn 子进程退出后内存池无效) */
|
|
|
export function canPersistParkedQuoteSession(): boolean {
|
|
|
return Boolean(process.env.RPA_WORKER_ID?.trim());
|
|
|
}
|
|
|
|
|
|
async function markParkedInRedis(sessionId: string): Promise<void> {
|
|
|
const workerId = process.env.RPA_WORKER_ID?.trim();
|
|
|
if (!workerId) {
|
|
|
return;
|
|
|
}
|
|
|
try {
|
|
|
await getRedis().setex(
|
|
|
`${PARKED_REDIS_PREFIX}${sessionId}`,
|
|
|
QUOTE_SESSION_TTL_SECONDS,
|
|
|
workerId,
|
|
|
);
|
|
|
} catch {
|
|
|
/* Redis 不可用时仍保留进程内驻留 */
|
|
|
}
|
|
|
}
|
|
|
|
|
|
async function clearParkedInRedis(sessionId: string): Promise<void> {
|
|
|
try {
|
|
|
await getRedis().del(`${PARKED_REDIS_PREFIX}${sessionId}`);
|
|
|
} catch {
|
|
|
/* ignore */
|
|
|
}
|
|
|
}
|
|
|
|
|
|
/** Redis 侧是否曾有 worker 驻留(诊断 spawn/跨进程) */
|
|
|
export async function readParkedSessionWorkerId(
|
|
|
sessionId: string,
|
|
|
): Promise<string | null> {
|
|
|
try {
|
|
|
return await getRedis().get(`${PARKED_REDIS_PREFIX}${sessionId}`);
|
|
|
} catch {
|
|
|
return null;
|
|
|
}
|
|
|
}
|
|
|
|
|
|
/** 清理过期驻留会话 */
|
|
|
export function sweepExpiredParkedSessions(): number {
|
|
|
let removed = 0;
|
|
|
for (const [id, entry] of pool) {
|
|
|
if (isExpired(entry)) {
|
|
|
void releaseParkedQuoteSession(id);
|
|
|
removed += 1;
|
|
|
}
|
|
|
}
|
|
|
return removed;
|
|
|
}
|
|
|
|
|
|
/** 候选完成后驻留 page/context,供询价 job 同页续跑 */
|
|
|
export async function parkQuoteSession(
|
|
|
sessionId: string,
|
|
|
page: Page,
|
|
|
context: BrowserContext,
|
|
|
): Promise<void> {
|
|
|
if (!canPersistParkedQuoteSession()) {
|
|
|
throw new RpaError(
|
|
|
"STRUCT_CHANGE",
|
|
|
"非 RPA Worker 进程,无法驻留浏览器会话",
|
|
|
{ retryable: false },
|
|
|
);
|
|
|
}
|
|
|
sweepExpiredParkedSessions();
|
|
|
await releaseParkedQuoteSession(sessionId);
|
|
|
pool.set(sessionId, {
|
|
|
sessionId,
|
|
|
page,
|
|
|
context,
|
|
|
parkedAt: Date.now(),
|
|
|
});
|
|
|
await markParkedInRedis(sessionId);
|
|
|
console.log(
|
|
|
`[parked-session] 已驻留 session=${sessionId.slice(0, 8)}… count=${pool.size} worker=${process.env.RPA_WORKER_ID}`,
|
|
|
);
|
|
|
}
|
|
|
|
|
|
/** 询价 job 接管驻留会话(从池中移除) */
|
|
|
export async function takeParkedQuoteSession(
|
|
|
sessionId: string,
|
|
|
options?: { waitMs?: number },
|
|
|
): Promise<{ page: Page; context: BrowserContext } | null> {
|
|
|
const waitMs = options?.waitMs ?? 0;
|
|
|
const deadline = Date.now() + waitMs;
|
|
|
|
|
|
while (Date.now() <= deadline) {
|
|
|
sweepExpiredParkedSessions();
|
|
|
const entry = pool.get(sessionId);
|
|
|
if (entry && !isExpired(entry)) {
|
|
|
pool.delete(sessionId);
|
|
|
await clearParkedInRedis(sessionId);
|
|
|
console.log(
|
|
|
`[parked-session] 已接管 session=${sessionId.slice(0, 8)}…(单次 RPA 续跑)`,
|
|
|
);
|
|
|
return { page: entry.page, context: entry.context };
|
|
|
}
|
|
|
if (waitMs <= 0) {
|
|
|
break;
|
|
|
}
|
|
|
await new Promise((resolve) => setTimeout(resolve, 200));
|
|
|
}
|
|
|
|
|
|
const redisWorker = await readParkedSessionWorkerId(sessionId);
|
|
|
if (redisWorker && !pool.has(sessionId)) {
|
|
|
console.warn(
|
|
|
`[parked-session] Redis 记录 worker=${redisWorker} 但本进程无驻留页(可能 spawn 候选或 worker 已重启)session=${sessionId.slice(0, 8)}…`,
|
|
|
);
|
|
|
}
|
|
|
return null;
|
|
|
}
|
|
|
|
|
|
/** 用户取消或超时:关闭并移除 */
|
|
|
export async function releaseParkedQuoteSession(
|
|
|
sessionId: string,
|
|
|
): Promise<boolean> {
|
|
|
const entry = pool.get(sessionId);
|
|
|
if (!entry) {
|
|
|
await clearParkedInRedis(sessionId);
|
|
|
return false;
|
|
|
}
|
|
|
pool.delete(sessionId);
|
|
|
await clearParkedInRedis(sessionId);
|
|
|
await entry.page.close().catch(() => undefined);
|
|
|
await entry.context.close().catch(() => undefined);
|
|
|
console.log(
|
|
|
`[parked-session] 已释放 session=${sessionId.slice(0, 8)}… count=${pool.size}`,
|
|
|
);
|
|
|
return true;
|
|
|
}
|
|
|
|
|
|
export async function releaseAllParkedQuoteSessions(): Promise<void> {
|
|
|
const ids = [...pool.keys()];
|
|
|
for (const id of ids) {
|
|
|
await releaseParkedQuoteSession(id);
|
|
|
}
|
|
|
}
|
|
|
|
|
|
export function hasParkedQuoteSession(sessionId: string): boolean {
|
|
|
sweepExpiredParkedSessions();
|
|
|
const entry = pool.get(sessionId);
|
|
|
return Boolean(entry && !isExpired(entry));
|
|
|
}
|
|
|
|
|
|
/** 用户确认弹窗期间:轻量维持 widget 就绪 */
|
|
|
export async function preheatParkedQuoteSession(
|
|
|
sessionId: string,
|
|
|
): Promise<boolean> {
|
|
|
sweepExpiredParkedSessions();
|
|
|
const entry = pool.get(sessionId);
|
|
|
if (!entry || isExpired(entry)) {
|
|
|
return false;
|
|
|
}
|
|
|
await dismissCookieConsent(entry.page);
|
|
|
await scrollQuoteWidgetIntoView(entry.page, { force: true });
|
|
|
console.log(
|
|
|
`[parked-session] 已预热 session=${sessionId.slice(0, 8)}…`,
|
|
|
);
|
|
|
return true;
|
|
|
}
|
|
|
|
|
|
/** 驻留页复用前轻量检查(跳过 15~45s widget hydration 轮询) */
|
|
|
export async function prepareParkedQuotePage(page: Page): Promise<void> {
|
|
|
await dismissCookieConsent(page);
|
|
|
await scrollQuoteWidgetIntoView(page, { force: true });
|
|
|
}
|
|
|
|
|
|
export function parkedQuoteSessionCount(): number {
|
|
|
sweepExpiredParkedSessions();
|
|
|
return pool.size;
|
|
|
}
|