You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

201 lines
5.5 KiB

This file contains ambiguous Unicode characters!

This file contains ambiguous Unicode characters that may be confused with others in your current locale. If your use case is intentional and legitimate, you can safely ignore this warning. Use the Escape button to highlight these characters.

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 进程可驻留 pagespawn 子进程退出后内存池无效) */
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;
}