diff --git a/CHANGELOG.md b/CHANGELOG.md index 034725b..5e47905 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,21 @@ 本项目遵循 [Keep a Changelog](https://keepachangelog.com/zh-CN/1.1.0/) 的基本格式,并计划采用[语义化版本](https://semver.org/lang/zh-CN/)。 +## [Unreleased] + +### Changed + +- 超长会话改为流式读取,不再把整个 Codex 会话文件一次性载入内存;同时支持已归档会话和失效路径回退。 +- 通知状态与正文会优先提取明确的回复、确认和继续操作要求,并清理内部状态标记、本地路径与空泛摘要。 +- 优先使用 ChatGPT 或 Codex App 自带的兼容 Node.js,减少包管理器升级后通知失效的概率。 +- 已发送去重记录只保留 90 天,审计日志达到 10 MB 后轮换,安装备份最多保留最近 5 份。 + +### Fixed + +- 修复超大 JSONL 会话导致主任务识别超时、漏发通知或误判为内部任务的问题。 +- 修复“告诉我结果”“确认后继续”等明确请求被误判为“本轮结束”的问题。 +- 修复通知摘要为空、只显示“请回复”、保留内部隐藏标记或在完整确认事项中间截断的问题。 + ## [0.1.0] - 2026-08-01 ### Added diff --git a/scripts/install.mjs b/scripts/install.mjs index 55d5600..e996a12 100644 --- a/scripts/install.mjs +++ b/scripts/install.mjs @@ -13,6 +13,7 @@ import { assertSupportedRuntime, arraysEqual, atomicWrite, + capBackupRecords, createBackups, createRuntimeSnapshot, captureTextFileState, @@ -27,6 +28,7 @@ import { pathExists, permissionHook, promptHiddenDeviceKey, + pruneBackupDirectories, prepareRuntimeStage, readDeviceKeyFromFile, readManifest, @@ -63,7 +65,7 @@ Never pass a Bark key as a command-line value. Without --key-file, a hidden interactive prompt is used.`; function printPlan(plan) { - console.log(`Node.js: ${process.execPath} (${process.versions.node})`); + console.log(`Node.js: ${plan.managedNotify[0]} (${process.versions.node})`); console.log(`Codex home: ${plan.paths.codexHome}`); console.log(`Install root: ${plan.paths.installRoot}`); console.log(`notify mode: ${plan.notifyMode}`); @@ -216,8 +218,18 @@ export async function install({ input = process.stdin, output = process.stderr, operations = {}, + environment = process.env, } = {}) { assertSupportedRuntime(runtime); + const selectedNodePath = String( + environment._CODEX_BARK_SELECTED_NODE ?? process.execPath, + ).trim(); + if ( + !selectedNodePath.startsWith("/") || + /[\u0000-\u001f\u007f]/u.test(selectedNodePath) + ) { + throw new Error("The selected Node.js runtime path is invalid."); + } const options = parseArguments(argv, "install"); if (options.help) { console.log(HELP); @@ -233,7 +245,11 @@ export async function install({ await assertNotSymlink(paths.hooksJson, { requireFile: true }); const runtimeConfigState = await inspectRuntimeConfig(paths); const existingManifest = await readManifest(paths); - const plan = await buildInstallPlan({ paths, existingManifest }); + const plan = await buildInstallPlan({ + paths, + existingManifest, + nodePath: selectedNodePath, + }); printPlan(plan); if (options.dryRun) { console.log("Dry run complete; no files were changed and no key was read."); @@ -254,7 +270,7 @@ export async function install({ const stage = await prepareRuntimeStage(paths, { notifyMode: plan.notifyMode, previousNotify: plan.previousNotify, - nodePath: process.execPath, + nodePath: selectedNodePath, }); let backup; let snapshot; @@ -318,7 +334,7 @@ export async function install({ version: "0.1.0", status: "installed", installRoot: paths.installRoot, - nodePath: process.execPath, + nodePath: selectedNodePath, installedAt: existingManifest?.status === "installed" ? existingManifest.installedAt @@ -342,10 +358,10 @@ export async function install({ managedEntry: plan.managedHook, before: plan.hooksBeforeState, }, - backups: [ + backups: capBackupRecords([ ...(existingManifest?.backups ?? []), backup, - ], + ]), files: {}, }; manifest.files = await runtimeFileHashes(paths); @@ -353,6 +369,7 @@ export async function install({ parseJsonObject(manifestSource, "installed.json"); await writeAtomic(paths.manifest, manifestSource, 0o600); await removeRuntimeTemporary(stage, snapshot); + await pruneBackupDirectories(paths).catch(() => {}); console.log("Installed Codex Bark Notifier."); console.log(`Backup: ${backup.directory}`); console.log("Ready now: turn completion, reply-needed, and blocked/error notifications."); diff --git a/scripts/install.sh b/scripts/install.sh index 8dbc1d3..fb2c6ad 100755 --- a/scripts/install.sh +++ b/scripts/install.sh @@ -60,12 +60,6 @@ select_node() { return 0 fi - path_node=$(command -v node 2>/dev/null) || path_node= - if [ -n "$path_node" ] && is_supported_node "$path_node"; then - printf '%s\n' "$path_node" - return 0 - fi - # The override below is intentionally private to the test suite. Production # callers use the literal /Applications root. system_applications=${_CODEX_BARK_TEST_APPLICATIONS_ROOT:-/Applications} @@ -92,6 +86,15 @@ select_node() { done fi + # Keep PATH as the final fallback. Homebrew's public bin/node symlink is more + # stable than the versioned Cellar path reported by process.execPath, so the + # wrapper passes this selected path through to the installer below. + path_node=$(command -v node 2>/dev/null) || path_node= + if [ -n "$path_node" ] && is_supported_node "$path_node"; then + printf '%s\n' "$path_node" + return 0 + fi + print_error "No compatible Node.js runtime found. Node.js ${minimum_node_major}+ is required." print_error "Install/update the Codex desktop app, install Node.js ${minimum_node_major}+, or set CODEX_BARK_NODE." return 1 @@ -115,6 +118,12 @@ esac script_directory=${script_path%/*} project_directory=${script_directory%/*} node_executable=$(select_node) || exit 1 +case $node_executable in + /*) ;; + *) node_executable=$PWD/$node_executable ;; +esac +_CODEX_BARK_SELECTED_NODE=$node_executable +export _CODEX_BARK_SELECTED_NODE case ${1-} in --verify) diff --git a/scripts/lib/installer-core.mjs b/scripts/lib/installer-core.mjs index 6ce5e0f..d856b53 100644 --- a/scripts/lib/installer-core.mjs +++ b/scripts/lib/installer-core.mjs @@ -31,6 +31,9 @@ import { fileURLToPath } from "node:url"; export const PRODUCT = "codex-bark-notifier"; export const SCHEMA_VERSION = 1; +export const MAX_INSTALL_BACKUPS = 5; +const MANAGED_BACKUP_DIRECTORY_PATTERN = + /^\d{8}T\d{6}\.\d{3}Z$/u; export class UnsafeManifestError extends Error { constructor(message, options = {}) { super(message, options); @@ -152,6 +155,8 @@ export function installationPaths({ previousNotify: join(installRoot, "previous-notify.json"), manifest: join(installRoot, "installed.json"), auditLog: join(installRoot, "bark-notify.log"), + auditArchive: join(installRoot, "bark-notify.log.1"), + auditRotationLock: join(installRoot, "bark-notify.log.rotate.lock"), state: join(installRoot, "state"), jobs: join(installRoot, "jobs"), configToml: join(codexHome, "config.toml"), @@ -2020,6 +2025,64 @@ export async function createBackups(paths, files, stamp = timestamp()) { return { directory, files: records }; } +export function capBackupRecords( + backups, + maximumBackups = MAX_INSTALL_BACKUPS, +) { + if (!Array.isArray(backups)) { + throw new Error("Backup records must be an array."); + } + if (!Number.isSafeInteger(maximumBackups) || maximumBackups < 1) { + throw new Error("Backup retention must be a positive integer."); + } + return backups.slice(-maximumBackups); +} + +export async function pruneBackupDirectories( + paths, + maximumBackups = MAX_INSTALL_BACKUPS, +) { + if (!Number.isSafeInteger(maximumBackups) || maximumBackups < 1) { + throw new Error("Backup retention must be a positive integer."); + } + const rootMetadata = await assertNotSymlink(paths.backupRoot); + if (!rootMetadata) { + return 0; + } + if (!rootMetadata.isDirectory()) { + throw new Error(`Backup root is not a directory: ${paths.backupRoot}`); + } + + const managedEntries = (await readdir(paths.backupRoot, { + withFileTypes: true, + })) + .filter( + (entry) => + entry.isDirectory() && + MANAGED_BACKUP_DIRECTORY_PATTERN.test(entry.name), + ) + .sort((left, right) => left.name.localeCompare(right.name)); + const removableNames = new Set( + managedEntries + .slice(0, Math.max(0, managedEntries.length - maximumBackups)) + .map((entry) => entry.name), + ); + let removed = 0; + for (const entry of managedEntries) { + const directory = join(paths.backupRoot, entry.name); + if (!removableNames.has(entry.name)) { + continue; + } + const metadata = await lstat(directory); + if (!metadata.isDirectory() || metadata.isSymbolicLink()) { + continue; + } + await rm(directory, { recursive: true, force: true }); + removed += 1; + } + return removed; +} + export function sha256(content) { return createHash("sha256").update(content).digest("hex"); } @@ -2592,6 +2655,48 @@ export async function readManifest(paths) { "installed.json managed files are invalid.", ); } + if (!Array.isArray(manifest.backups)) { + throw new UnsafeManifestError( + "installed.json backup records must be an array.", + ); + } + for (const record of manifest.backups) { + const directory = record?.directory; + if ( + !record || + Array.isArray(record) || + typeof record !== "object" || + typeof directory !== "string" || + resolve(directory) !== directory || + dirname(directory) !== resolve(paths.backupRoot) || + !MANAGED_BACKUP_DIRECTORY_PATTERN.test(basename(directory)) || + !Array.isArray(record.files) + ) { + throw new UnsafeManifestError( + "installed.json contains invalid backup records.", + ); + } + for (const file of record.files) { + const allowedPath = [paths.configToml, paths.hooksJson].includes( + file?.path, + ); + const expectedBackup = file?.existed + ? join(directory, basename(file.path)) + : ""; + if ( + !file || + Array.isArray(file) || + typeof file !== "object" || + !allowedPath || + typeof file.existed !== "boolean" || + file.backup !== expectedBackup + ) { + throw new UnsafeManifestError( + "installed.json contains invalid backup file records.", + ); + } + } + } for (const [file, hash] of Object.entries(manifest.files)) { const resolvedFile = resolve(file); const relativeToLibrary = relative(paths.library, resolvedFile); diff --git a/src/lib/paths.mjs b/src/lib/paths.mjs index 60d9c1f..d577f4c 100644 --- a/src/lib/paths.mjs +++ b/src/lib/paths.mjs @@ -22,8 +22,12 @@ export function runtimePaths(entryUrl = import.meta.url, environment = process.e configFile: join(runtimeRoot, "config.json"), keyFile: join(runtimeRoot, "bark-device-key"), stateDirectory: join(runtimeRoot, "state"), + stateCleanupStamp: join(runtimeRoot, "state", ".sent-cleanup"), + stateCleanupLock: join(runtimeRoot, "state", ".sent-cleanup.lock"), jobsDirectory: join(runtimeRoot, "jobs"), auditLog: join(runtimeRoot, "bark-notify.log"), + auditArchive: join(runtimeRoot, "bark-notify.log.1"), + auditRotationLock: join(runtimeRoot, "bark-notify.log.rotate.lock"), sessionRoot: join(codexHome, "sessions"), sessionIndex: join(codexHome, "session_index.jsonl"), }; diff --git a/src/lib/sessions.mjs b/src/lib/sessions.mjs index e9c630e..ff09517 100644 --- a/src/lib/sessions.mjs +++ b/src/lib/sessions.mjs @@ -1,5 +1,7 @@ -import { readFile, readdir } from "node:fs/promises"; -import { basename, join } from "node:path"; +import { createReadStream } from "node:fs"; +import { open, readFile, readdir } from "node:fs/promises"; +import { basename, dirname, join } from "node:path"; +import { createInterface } from "node:readline"; import { normalizeTaskName, @@ -18,6 +20,127 @@ export const THREAD_KIND_RETRY_DELAYS = Object.freeze([ 4_000, ]); +const REVERSE_READ_CHUNK_BYTES = 64 * 1024; +const MAX_TRANSCRIPT_LINE_BYTES = 16 * 1024 * 1024; + +function transcriptRoots(sessionRoot) { + return [sessionRoot, join(dirname(sessionRoot), "archived_sessions")].filter( + (root, index, roots) => root && roots.indexOf(root) === index, + ); +} + +async function readableTranscriptPath(transcriptPath) { + if (!transcriptPath) { + return ""; + } + + let handle; + try { + handle = await open(transcriptPath, "r"); + return (await handle.stat()).isFile() ? transcriptPath : ""; + } catch { + return ""; + } finally { + await handle?.close().catch(() => {}); + } +} + +async function transcriptPathFromPayload(payload, threadId, sessionRoot) { + const preferredPath = await readableTranscriptPath( + typeof payload?.transcript_path === "string" + ? payload.transcript_path + : "", + ); + return ( + preferredPath || + (await sessionPathForThreadId(threadId, sessionRoot)) + ); +} + +async function* transcriptLinesReverse( + transcriptPath, + { + chunkBytes = REVERSE_READ_CHUNK_BYTES, + maxLineBytes = MAX_TRANSCRIPT_LINE_BYTES, + } = {}, +) { + const handle = await open(transcriptPath, "r"); + try { + let position = (await handle.stat()).size; + let lineParts = []; + let lineBytes = 0; + let skipLine = false; + + function addEarlierPart(part) { + if (skipLine || part.length === 0) { + return; + } + lineBytes += part.length; + if (lineBytes > maxLineBytes) { + lineParts = []; + lineBytes = 0; + skipLine = true; + return; + } + lineParts.push(Buffer.from(part)); + } + + function finishLine() { + if (skipLine) { + lineParts = []; + lineBytes = 0; + skipLine = false; + return null; + } + + const line = + lineParts.length === 0 + ? Buffer.alloc(0) + : Buffer.concat(lineParts.reverse(), lineBytes); + lineParts = []; + lineBytes = 0; + const end = line.at(-1) === 13 ? line.length - 1 : line.length; + return line.toString("utf8", 0, end); + } + + while (position > 0) { + const bytesToRead = Math.min(chunkBytes, position); + position -= bytesToRead; + const chunk = Buffer.allocUnsafe(bytesToRead); + const { bytesRead } = await handle.read( + chunk, + 0, + bytesToRead, + position, + ); + const data = bytesRead === chunk.length + ? chunk + : chunk.subarray(0, bytesRead); + let segmentEnd = data.length; + + for (let index = data.length - 1; index >= 0; index -= 1) { + if (data[index] !== 10) { + continue; + } + addEarlierPart(data.subarray(index + 1, segmentEnd)); + const line = finishLine(); + if (line !== null) { + yield line; + } + segmentEnd = index; + } + addEarlierPart(data.subarray(0, segmentEnd)); + } + + const firstLine = finishLine(); + if (firstLine !== null) { + yield firstLine; + } + } finally { + await handle.close(); + } +} + export function classifySessionMetaPayload(metaPayload, threadId = "") { if (!metaPayload || (threadId && metaPayload.id !== threadId)) { return "unknown"; @@ -46,18 +169,23 @@ export async function sessionPathForThreadId(threadId, sessionRoot) { return ""; } - try { - const sessionFiles = await readdir(sessionRoot, { recursive: true }); - const suffix = `-${threadId}.jsonl`; - const sessionFile = sessionFiles.find( - (entry) => - typeof entry === "string" && - basename(entry).endsWith(suffix), - ); - return sessionFile ? join(sessionRoot, sessionFile) : ""; - } catch { - return ""; + const suffix = `-${threadId}.jsonl`; + for (const root of transcriptRoots(sessionRoot)) { + try { + const sessionFiles = await readdir(root, { recursive: true }); + const sessionFile = sessionFiles.find( + (entry) => + typeof entry === "string" && + basename(entry).endsWith(suffix), + ); + if (sessionFile) { + return join(root, sessionFile); + } + } catch { + // Try the active/archive fallback root below. + } } + return ""; } export async function sessionMetaFromTranscript( @@ -69,10 +197,10 @@ export async function sessionMetaFromTranscript( return null; } + const input = createReadStream(transcriptPath, { encoding: "utf8" }); + const rows = createInterface({ input, crlfDelay: Infinity }); try { - const rows = (await readFile(transcriptPath, "utf8")).split(/\r?\n/u); - let fallbackMeta = null; - for (const row of rows) { + for await (const row of rows) { if (!row.trim()) { continue; } @@ -80,24 +208,53 @@ export async function sessionMetaFromTranscript( if (record?.type !== "session_meta") { continue; } - fallbackMeta ??= record.payload; if (!threadId || record?.payload?.id === threadId) { return record.payload; } + return allowAnySessionMetaFallback ? record.payload : null; } - return allowAnySessionMetaFallback ? fallbackMeta : null; + return null; } catch { return null; + } finally { + rows.close(); + input.destroy(); } } export async function threadKind(payload, paths) { const { threadId } = payloadIds(payload); - const transcriptPath = - payload?.transcript_path || - (await sessionPathForThreadId(threadId, paths.sessionRoot)); - const meta = await sessionMetaFromTranscript(transcriptPath, threadId); - return classifySessionMetaPayload(meta, threadId); + const preferredPath = await readableTranscriptPath( + typeof payload?.transcript_path === "string" + ? payload.transcript_path + : "", + ); + if (preferredPath) { + const preferredMeta = await sessionMetaFromTranscript( + preferredPath, + threadId, + ); + const preferredKind = classifySessionMetaPayload( + preferredMeta, + threadId, + ); + if (preferredKind !== "unknown") { + return preferredKind; + } + } + + const fallbackPath = await sessionPathForThreadId( + threadId, + paths.sessionRoot, + ); + if (!fallbackPath || fallbackPath === preferredPath) { + return "unknown"; + } + const fallbackMeta = await sessionMetaFromTranscript( + fallbackPath, + threadId, + ); + return classifySessionMetaPayload(fallbackMeta, threadId); } export async function resolvedThreadKind( @@ -164,9 +321,11 @@ export async function taskNameFromIndex( export async function notificationContext(payload, paths) { const { threadId } = payloadIds(payload); - const currentTranscript = - payload?.transcript_path || - (await sessionPathForThreadId(threadId, paths.sessionRoot)); + const currentTranscript = await transcriptPathFromPayload( + payload, + threadId, + paths.sessionRoot, + ); const parentThreadId = await parentThreadIdFromTranscript( currentTranscript, threadId, @@ -203,10 +362,7 @@ export async function conversationNameFromTranscript( } try { - const rows = (await readFile(transcriptPath, "utf8")) - .split(/\r?\n/u) - .reverse(); - for (const row of rows) { + for await (const row of transcriptLinesReverse(transcriptPath)) { if (!row.trim()) { continue; } @@ -249,10 +405,7 @@ export async function assistantSummaryFromTranscript( for (let attempt = 0; attempt <= retryDelays.length; attempt += 1) { try { - const rows = (await readFile(transcriptPath, "utf8")) - .split(/\r?\n/u) - .reverse(); - for (const row of rows) { + for await (const row of transcriptLinesReverse(transcriptPath)) { if (!row.trim()) { continue; } diff --git a/src/lib/state.mjs b/src/lib/state.mjs index 6018ebc..c815216 100644 --- a/src/lib/state.mjs +++ b/src/lib/state.mjs @@ -1,7 +1,6 @@ import { createHash, randomUUID } from "node:crypto"; import { constants as fileConstants } from "node:fs"; import { - appendFile, chmod, lstat, mkdir, @@ -11,6 +10,7 @@ import { rename, stat, unlink, + utimes, } from "node:fs/promises"; import { basename, dirname, join, resolve } from "node:path"; @@ -18,10 +18,14 @@ import { payloadIds, shortenConversationName } from "./text.mjs"; export const STALE_LOCK_MILLISECONDS = 2 * 60 * 1_000; export const STALE_JOB_MILLISECONDS = 60 * 60 * 1_000; +export const SENT_MARKER_RETENTION_MILLISECONDS = 90 * 24 * 60 * 60 * 1_000; +export const SENT_MARKER_CLEANUP_INTERVAL_MILLISECONDS = 24 * 60 * 60 * 1_000; +export const AUDIT_LOG_MAX_BYTES = 10 * 1024 * 1024; export const PRIVATE_JOB_SCHEMA_VERSION = 2; const LEGACY_PRIVATE_JOB_SCHEMA_VERSION = 1; const MANAGED_JOB_NAME_PATTERN = /^[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}\.json$/u; +const MANAGED_SENT_MARKER_PATTERN = /^[0-9a-f]{64}\.sent$/u; export function eventKey(payload) { const { threadId, turnId } = payloadIds(payload); @@ -93,9 +97,108 @@ export function safeErrorReason(error) { return "unknown_error"; } -export async function writeAudit(paths, payload, event, details = {}) { +async function acquireMaintenanceLock( + lockPath, + { mayRecoverStaleLock = true } = {}, +) { + try { + const handle = await open(lockPath, "wx", 0o600); + await handle.close(); + return true; + } catch (error) { + if (error?.code !== "EEXIST") { + throw error; + } + } + + if (mayRecoverStaleLock) { + try { + const lockInfo = await stat(lockPath); + if (Date.now() - lockInfo.mtimeMs > STALE_LOCK_MILLISECONDS) { + await removeIfPresent(lockPath); + return acquireMaintenanceLock(lockPath, { + mayRecoverStaleLock: false, + }); + } + } catch { + return acquireMaintenanceLock(lockPath, { + mayRecoverStaleLock: false, + }); + } + } + + return false; +} + +export async function rotateAuditLog( + paths, + { maximumBytes = AUDIT_LOG_MAX_BYTES } = {}, +) { + if (!Number.isSafeInteger(maximumBytes) || maximumBytes < 1) { + throw new Error("Audit log maximum size must be a positive integer."); + } + + let auditInfo; + try { + auditInfo = await lstat(paths.auditLog); + } catch (error) { + if (error?.code === "ENOENT") { + return false; + } + throw error; + } + if ( + !auditInfo.isFile() || + auditInfo.isSymbolicLink() || + auditInfo.size < maximumBytes + ) { + return false; + } + + const archivePath = paths.auditArchive ?? `${paths.auditLog}.1`; + const lockPath = paths.auditRotationLock ?? `${paths.auditLog}.rotate.lock`; + if (!(await acquireMaintenanceLock(lockPath))) { + return false; + } + + try { + try { + auditInfo = await lstat(paths.auditLog); + } catch (error) { + if (error?.code === "ENOENT") { + return false; + } + throw error; + } + if ( + !auditInfo.isFile() || + auditInfo.isSymbolicLink() || + auditInfo.size < maximumBytes + ) { + return false; + } + await rename(paths.auditLog, archivePath); + await chmod(archivePath, 0o600); + return true; + } finally { + await removeIfPresent(lockPath); + } +} + +export async function writeAudit( + paths, + payload, + event, + details = {}, + { maximumBytes = AUDIT_LOG_MAX_BYTES } = {}, +) { try { await ensurePrivateDirectory(paths.runtimeRoot); + try { + await rotateAuditLog(paths, { maximumBytes }); + } catch { + // A rotation problem must not suppress the current audit record. + } const { threadId, turnId } = payloadIds(payload); const record = { timestamp: new Date().toISOString(), @@ -115,15 +218,128 @@ export async function writeAudit(paths, payload, event, details = {}) { } } - await appendFile(paths.auditLog, `${JSON.stringify(record)}\n`, { - mode: 0o600, - }); + const handle = await open( + paths.auditLog, + fileConstants.O_APPEND | + fileConstants.O_CREAT | + fileConstants.O_WRONLY | + fileConstants.O_NOFOLLOW, + 0o600, + ); + try { + await handle.writeFile(`${JSON.stringify(record)}\n`, "utf8"); + } finally { + await handle.close(); + } await chmod(paths.auditLog, 0o600); } catch { // Delivery should not fail only because audit logging failed. } } +export async function cleanupExpiredSentMarkers( + paths, + { + now = Date.now(), + retentionMilliseconds = SENT_MARKER_RETENTION_MILLISECONDS, + minimumIntervalMilliseconds = SENT_MARKER_CLEANUP_INTERVAL_MILLISECONDS, + } = {}, +) { + if (!Number.isSafeInteger(retentionMilliseconds) || retentionMilliseconds < 1) { + throw new Error("Sent marker retention must be a positive integer."); + } + if ( + !Number.isSafeInteger(minimumIntervalMilliseconds) || + minimumIntervalMilliseconds < 0 + ) { + throw new Error("Sent marker cleanup interval must be a non-negative integer."); + } + + await ensurePrivateDirectory(paths.stateDirectory); + const stampPath = + paths.stateCleanupStamp ?? join(paths.stateDirectory, ".sent-cleanup"); + const lockPath = + paths.stateCleanupLock ?? join(paths.stateDirectory, ".sent-cleanup.lock"); + const cleanupIsRecent = async () => { + if (minimumIntervalMilliseconds === 0) { + return false; + } + try { + const info = await lstat(stampPath); + return ( + info.isFile() && + !info.isSymbolicLink() && + now - info.mtimeMs <= minimumIntervalMilliseconds + ); + } catch { + return false; + } + }; + if (await cleanupIsRecent()) { + return 0; + } + if (!(await acquireMaintenanceLock(lockPath))) { + return 0; + } + + try { + let entries = []; + if (await cleanupIsRecent()) { + return 0; + } + try { + entries = await readdir(paths.stateDirectory, { withFileTypes: true }); + } catch { + return 0; + } + + let removed = 0; + await Promise.all( + entries.map(async (entry) => { + if (!entry.isFile() || !MANAGED_SENT_MARKER_PATTERN.test(entry.name)) { + return; + } + const candidatePath = join(paths.stateDirectory, entry.name); + try { + const info = await lstat(candidatePath); + if ( + !info.isFile() || + info.isSymbolicLink() || + now - info.mtimeMs <= retentionMilliseconds + ) { + return; + } + await unlink(candidatePath); + removed += 1; + } catch { + // Cleanup is best effort and must not prevent deduplication. + } + }), + ); + if (!entries.some( + (entry) => + entry.isFile() && MANAGED_SENT_MARKER_PATTERN.test(entry.name), + )) { + return removed; + } + const stampHandle = await open( + stampPath, + fileConstants.O_CREAT | + fileConstants.O_TRUNC | + fileConstants.O_WRONLY | + fileConstants.O_NOFOLLOW, + 0o600, + ); + await stampHandle.close(); + await chmod(stampPath, 0o600); + const stampTime = new Date(now); + await utimes(stampPath, stampTime, stampTime); + return removed; + } finally { + await removeIfPresent(lockPath); + } +} + export async function acquireEventLock( payload, paths, @@ -138,6 +354,7 @@ export async function acquireEventLock( } await ensurePrivateDirectory(paths.stateDirectory); + await cleanupExpiredSentMarkers(paths).catch(() => {}); const lockPath = join(paths.stateDirectory, `${key}.lock`); const sentPath = join(paths.stateDirectory, `${key}.sent`); diff --git a/src/lib/text.mjs b/src/lib/text.mjs index b056f0f..983f64f 100644 --- a/src/lib/text.mjs +++ b/src/lib/text.mjs @@ -40,9 +40,14 @@ const LEGACY_CREDENTIAL_ONLY_PATTERN = const STANDALONE_CREDENTIAL_PATTERN = /^\s*[A-Za-z0-9_-]{28,}\s*[。.!!]?\s*$/u; const SHORT_STANDALONE_BARK_KEY_PATTERN = - /^(?=\S{8,27}$)(?=\S*[A-Za-z])(?=\S*\d)\S+$/u; + /^(?=[A-Za-z0-9_-]{8,27}[。.!!]?$)(?=[A-Za-z0-9_-]*[A-Za-z])(?=[A-Za-z0-9_-]*\d)[A-Za-z0-9_-]+[。.!!]?$/u; const PUNCTUATED_STANDALONE_BARK_KEY_PATTERN = - /^(?=\S{8,256}$)(?=\S*[A-Za-z])(?=\S*\d)(?=\S*[._+\/-])\S+$/u; + /^(?=[A-Za-z0-9._+\/-]{8,256}[。!!]?$)(?=[A-Za-z0-9._+\/-]*[A-Za-z])(?=[A-Za-z0-9._+\/-]*\d)(?=[A-Za-z0-9._+\/-]*[._+\/-])[A-Za-z0-9._+\/-]+[。!!]?$/u; +const INTERNAL_STATUS_MARKER_PATTERN = + /^(?:\[(?:COMPLETE(?:D)?|SUCCESS(?:FUL)?|SUCCEEDED|BLOCKED|ERROR|FAILED|NEEDS[_ -]?(?:INPUT|ACTION|REPLY)|WAITING[_ -]?(?:FOR[_ -]?)?USER)\]\s*)+/iu; +const INTERNAL_ERROR_STATUS_MARKER_PATTERN = + /(?:^|\n)\s*\[(?:BLOCKED|ERROR|FAILED)\](?:\s|$)/iu; +const LOCAL_PATH_PLACEHOLDER = "__CODEX_LOCAL_PATH__"; function stripBidiControls(value) { return String(value ?? "").replace(BIDI_CONTROL_PATTERN, ""); @@ -338,6 +343,10 @@ function stripListMarker(value) { ); } +function stripInternalStatusMarker(value) { + return String(value ?? "").replace(INTERNAL_STATUS_MARKER_PATTERN, ""); +} + function truncateSummaryText(rawText, fallback) { const safeFallback = stripBidiControls(fallback); const candidate = @@ -412,7 +421,7 @@ export function shortenAssistantSummary(rawText, fallback = "未生成摘要") { const candidates = answerSummarySegments(cleanedText) .map((segment) => redactSensitiveSummaryValues( - stripListMarker(segment.trim()) + stripInternalStatusMarker(stripListMarker(segment.trim())) .replace(/^#{1,6}\s*/u, "") .replace( /^(?:好的|可以|明白了?|理解(?:了)?|没问题|收到|好哒?|当然|行)[,,。!!\s]+/u, @@ -611,7 +620,9 @@ const STATUS_DIRECT_REQUEST_CLAUSE_PATTERNS = [ ]; const STATUS_CONTINUATION_REQUEST_PATTERNS = [ /(?:^|[,,;;::\s])(?:有[^。!?!?;;\n]{0,32})?需要你确认(?:[^。!?!?;;\n]{0,48})/u, - /(?:^|[,,;;::\s])你确认[^。!?!?;;\n]{0,32}(?:的话|后)[^。!?!?;;\n]{0,48}我(?:就|会)/u, + /(?:^|[,,;;::\s])你确认[^。!?!?;;\n]{0,32}(?:的话|后)[^。!?!?;;\n]{0,48}我(?:就|会|再|将|继续)/u, + /^(?:等待|等)你[^。!?!?;;\n]{0,40}确认(?:后|以后)/u, + /^(?:(?:现在|接下来|目前|这边)\s*)?需要先确认[^。!?!?;;\n]{1,180}/u, /(?:^|[,,;;::\s])(?:下一步|接下来)(?![^。!?!?;;\n]{0,20}(?:我会|我将|系统会|自动))[^。!?!?;;\n]{0,64}(?:发给我|提供给我|上传给我|截图给我|回复|确认|提交给我)/u, /(?:^|[,,;;::\s])(?:下一张(?:先|优先)?|最优先|优先)(?:请)?发(?!送)/u, /(?:完成(?:以后|后)|改完(?:保存)?后)[^。!?!?;;\n]{0,48}(?:回复(?:我)?|告诉我|发给我)/u, @@ -625,7 +636,7 @@ function statusSentences(value) { line.match(/[^。!?!?]+(?:[。!?!?]+|$)/gu) ?? [], ) .map((sentence) => - stripListMarker(sentence.trim()) + stripInternalStatusMarker(stripListMarker(sentence.trim())) .replace(/^#{1,6}\s*/u, "") .trim(), ) @@ -697,7 +708,7 @@ function requiredActionCandidates(value) { }, ]; const clauses = sentence - .split(/(? clause.trim()) .filter(Boolean); if (clauses.length > 1) { @@ -726,12 +737,22 @@ function requiredActionPriority(value, inheritedOptional = false) { priority += 70; } if ( - /你确认[^。!?!?;;\n]{0,32}(?:的话|后)[^。!?!?;;\n]{0,48}我(?:就|会)/u.test( + /你确认[^。!?!?;;\n]{0,32}(?:的话|后)[^。!?!?;;\n]{0,48}我(?:就|会|再|将|继续)/u.test( text, ) ) { priority += 90; } + if (/(?:等待|等)你[^。!?!?;;\n]{0,40}确认(?:后|以后)/u.test(text)) { + priority += 90; + } + if ( + /需要先确认[^。!?!?;;\n]{1,180}/u.test( + text, + ) + ) { + priority += 80; + } if (/(?:下一步|下一张|接下来|最优先|优先)/u.test(text)) { priority += 60; } @@ -765,7 +786,12 @@ function requiredActionPriority(value, inheritedOptional = false) { } function compactRequiredAction(value) { - const text = String(value ?? "").trim(); + const text = String(value ?? "") + .trim() + .replace( + /^(?:(?:现在|接下来|目前|这边)\s*)?需要先确认(?:[一二两三四五六七八九十\d]+)(?:件|项)(?:事)?[::]/u, + "请确认:", + ); if ( !/(?:回复我|告诉我|发给我|回复(?:\s|[“「『"'`::]|$))/u.test(text) ) { @@ -788,6 +814,23 @@ function punctuateRequiredAction(value) { return `${text}。`; } +function isVagueRequiredAction(value) { + return /^(?:请)?回复[。!?!?]*$/u.test( + String(value ?? "").trim(), + ); +} + +function safeRequiredActionFallback(value, { sourceContainsLocalPath = false } = {}) { + if (sourceContainsLocalPath) { + return /确认/u.test(String(value ?? "")) + ? "请确认是否继续处理本地文件。" + : "请按提示处理本地文件后继续。"; + } + return /确认/u.test(String(value ?? "")) + ? "需要你确认后才能继续。" + : "需要你回复后才能继续。"; +} + export function notificationBodyFromAssistantReply( rawText, status, @@ -806,22 +849,43 @@ export function notificationBodyFromAssistantReply( .replace(/~~~[\s\S]*?~~~/gu, " ") .replace(//giu, " ") .replace(/!\[[^\]]*\]\([^)]*\)/gu, " ") - .replace(/\[([^\]]+)\]\([^)]*\)/gu, "$1") + .replace( + /([::])\s*\r?\n+\s*`([^`\r\n]+)`/gu, + (match, delimiter, codeValue) => + containsAbsoluteLocalPath(codeValue) + ? `${delimiter}${LOCAL_PATH_PLACEHOLDER}` + : match, + ) + .replace(/\[([^\]]+)\]\(([^)]*)\)/gu, (_match, labelText, target) => + containsAbsoluteLocalPath(target) + ? `${labelText}${LOCAL_PATH_PLACEHOLDER}` + : labelText, + ) .replace(/\bhttps?:\/\/[^\s<>"'`,;,。!?;、]+/giu, " ") - .replace(/::(?:code-comment|created-thread)\{[^}\n]*\}/gu, " "); + .replace(/::(?:code-comment|created-thread)\{[^}\n]*\}/gu, " ") + .replace(/((?:请)?回复\s*[::])\s*\r?\n+\s*/gu, "$1"); + const candidates = requiredActionCandidates(cleanedText) - .map(({ text, inheritedOptional, wholeSentence }) => ({ - inheritedOptional, - wholeSentence, - text: redactSensitiveSummaryValues( - stripListMarker(text.trim()) - .replace(/^#{1,6}\s*/u, "") - .replace(/[*_`~]/gu, "") - .replace(/\s+/gu, " ") - .trim(), - ), - })) + .map(({ text, inheritedOptional, wholeSentence }) => { + const containsLocalPathPlaceholder = text.includes( + LOCAL_PATH_PLACEHOLDER, + ); + return { + inheritedOptional, + wholeSentence, + containsLocalPathPlaceholder, + text: redactSensitiveSummaryValues( + stripInternalStatusMarker(stripListMarker(text.trim())) + .replaceAll(LOCAL_PATH_PLACEHOLDER, "") + .replace(/^#{1,6}\s*/u, "") + .replace(/[*_`~]/gu, "") + .replace(/\s+/gu, " ") + .trim(), + ), + }; + }) .filter(({ text }) => Boolean(text)) + .filter(({ text }) => !isVagueRequiredAction(text)) .filter(({ text }) => !containsAbsoluteLocalPath(text)) .filter(({ text }) => { const classificationText = textForStatusClassification(text); @@ -836,8 +900,18 @@ export function notificationBodyFromAssistantReply( candidate.text, candidate.inheritedOptional, ) + + (candidate.wholeSentence && + /(?:等待|等)你[^。!?!?;;\n]{0,40}确认(?:后|以后)[^。!?!?;;\n]{0,48}我(?:再|将|继续)/u.test( + candidate.text, + ) + ? 5 + : 0) + (candidate.wholeSentence && /[??]\s*$/u.test(candidate.text) ? 10 + : 0) + + (candidate.wholeSentence && + /(?:回复(?:我)?|告诉我|发给我)\s*[::]/u.test(candidate.text) + ? 15 : 0), })); const candidate = candidates.reduce((best, current) => { @@ -851,16 +925,32 @@ export function notificationBodyFromAssistantReply( return current; } return best; - }, null)?.text; - return candidate - ? truncateSummaryText( - sanitizeNotificationText( - punctuateRequiredAction(compactRequiredAction(candidate)), - safeFallback, - ), - safeFallback, - ) - : shortenAssistantSummary(rawText, fallback); + }, null); + if (candidate) { + if (candidate.containsLocalPathPlaceholder) { + return safeRequiredActionFallback(candidate.text, { + sourceContainsLocalPath: true, + }); + } + const compactCandidate = sanitizeNotificationText( + punctuateRequiredAction(compactRequiredAction(candidate.text)), + safeFallback, + ); + return isVagueRequiredAction(compactCandidate) + ? safeRequiredActionFallback(cleanedText) + : truncateSummaryText(compactCandidate, safeFallback); + } + + const summaryFallback = shortenAssistantSummary(rawText, fallback); + const actionTailContainsLocalPath = statusSentences(cleanedText) + .slice(-3) + .some((sentence) => containsAbsoluteLocalPath(sentence)); + return actionTailContainsLocalPath || + isVagueRequiredAction(summaryFallback) + ? safeRequiredActionFallback(cleanedText, { + sourceContainsLocalPath: actionTailContainsLocalPath, + }) + : summaryFallback; } export function classifyLastReply(lastReply) { @@ -884,6 +974,7 @@ export function classifyLastReply(lastReply) { const unresolvedEnding = withoutResolvedStatusPhrases(decisiveEnding); if ( + INTERNAL_ERROR_STATUS_MARKER_PATTERN.test(cleanedReply) || STATUS_ERROR_PATTERN.test(unresolvedSummary) || STATUS_DECISIVE_ERROR_PATTERN.test(unresolvedEnding) || STATUS_FINAL_ERROR_PATTERN.test( diff --git a/test/bootstrap.test.mjs b/test/bootstrap.test.mjs index 7566c7f..9eea19c 100644 --- a/test/bootstrap.test.mjs +++ b/test/bootstrap.test.mjs @@ -169,7 +169,7 @@ test("CODEX_BARK_NODE is an explicit validated override and preserves arguments" ); }); -test("PATH node wins bundled runtimes", async (t) => { +test("bundled runtime wins PATH to avoid versioned package-manager paths", async (t) => { const context = await temporaryBootstrapEnvironment(t); await createFakeNode(join(context.pathDirectory, "node"), { label: "path", @@ -181,8 +181,8 @@ test("PATH node wins bundled runtimes", async (t) => { const result = await runBootstrap(context, ["--dry-run"]); assert.equal(result.code, 0, result.stderr); - assert.match(await invocationLog(context), /^BEGIN:path$/mu); - assert.doesNotMatch(await invocationLog(context), /system-chatgpt/u); + assert.match(await invocationLog(context), /^BEGIN:system-chatgpt$/mu); + assert.doesNotMatch(await invocationLog(context), /^BEGIN:path$/mu); }); test("system ChatGPT runtime wins later bundled candidates", async (t) => { diff --git a/test/helpers.mjs b/test/helpers.mjs index dbcf16d..397028d 100644 --- a/test/helpers.mjs +++ b/test/helpers.mjs @@ -18,8 +18,12 @@ export async function temporaryPaths() { configFile: join(runtimeRoot, "config.json"), keyFile: join(runtimeRoot, "bark-device-key"), stateDirectory: join(runtimeRoot, "state"), + stateCleanupStamp: join(runtimeRoot, "state", ".sent-cleanup"), + stateCleanupLock: join(runtimeRoot, "state", ".sent-cleanup.lock"), jobsDirectory: join(runtimeRoot, "jobs"), auditLog: join(runtimeRoot, "bark-notify.log"), + auditArchive: join(runtimeRoot, "bark-notify.log.1"), + auditRotationLock: join(runtimeRoot, "bark-notify.log.rotate.lock"), sessionRoot, sessionIndex: join(codexHome, "session_index.jsonl"), }; diff --git a/test/installer.test.mjs b/test/installer.test.mjs index e151eed..16bcdbe 100644 --- a/test/installer.test.mjs +++ b/test/installer.test.mjs @@ -22,6 +22,7 @@ import { install } from "../scripts/install.mjs"; import { atomicWrite, acquireLifecycleLock, + capBackupRecords, createRuntimeSnapshot, dispatcherSource, enableHooksFeature, @@ -33,6 +34,7 @@ import { pathExists, prepareRuntimeStage, promptHiddenDeviceKey, + pruneBackupDirectories, readDeviceKeyFromFile, readManifest, releaseLifecycleLock, @@ -84,6 +86,63 @@ function fileMode(metadata) { return metadata.mode & 0o777; } +test("backup retention keeps five recent managed snapshots only", async (t) => { + const context = await fixture(t); + const names = Array.from( + { length: 7 }, + (_, index) => `2026080${index + 1}T120000.000Z`, + ); + await mkdir(context.paths.backupRoot, { recursive: true, mode: 0o700 }); + for (const name of names) { + await mkdir(join(context.paths.backupRoot, name), { mode: 0o700 }); + } + const unrelated = join(context.paths.backupRoot, "keep-me"); + await mkdir(unrelated); + + const allRecords = names.map((name) => ({ + directory: join(context.paths.backupRoot, name), + files: [], + })); + const kept = capBackupRecords(allRecords); + assert.deepEqual( + kept.map((record) => record.directory), + allRecords.slice(-5).map((record) => record.directory), + ); + assert.equal(await pruneBackupDirectories(context.paths), 2); + assert.deepEqual( + (await readdir(context.paths.backupRoot)).sort(), + [...names.slice(-5), "keep-me"].sort(), + ); +}); + +test("invalid backup metadata stops upgrades without deleting snapshots", async (t) => { + const context = await fixture(t); + await installFixture(context); + const extraBackup = join( + context.paths.backupRoot, + "20260809T120000.000Z", + ); + await mkdir(extraBackup, { mode: 0o700 }); + + const manifest = JSON.parse(await readFile(context.paths.manifest, "utf8")); + manifest.backups = {}; + await writeFile( + context.paths.manifest, + `${JSON.stringify(manifest, null, 2)}\n`, + { mode: 0o600 }, + ); + const before = (await readdir(context.paths.backupRoot)).sort(); + + await assert.rejects( + installFixture(context), + /backup records must be an array/iu, + ); + assert.deepEqual( + (await readdir(context.paths.backupRoot)).sort(), + before, + ); +}); + test("key inputs reject argv secrets and unsafe source files", async (t) => { const sentinel = `DO_NOT_PERSIST_${process.hrtime.bigint()}`; for (const argv of [ diff --git a/test/sessions.test.mjs b/test/sessions.test.mjs index 33294d9..921f86f 100644 --- a/test/sessions.test.mjs +++ b/test/sessions.test.mjs @@ -1,14 +1,17 @@ import assert from "node:assert/strict"; import { mkdir, writeFile } from "node:fs/promises"; -import { join } from "node:path"; +import { dirname, join } from "node:path"; import test from "node:test"; import { buildPermissionNotification } from "../src/bark-notify.mjs"; import { assistantSummaryFromTranscript, classifySessionMetaPayload, + conversationNameFromTranscript, notificationContext, resolvedThreadKind, + sessionMetaFromTranscript, + sessionPathForThreadId, taskNameFromIndex, THREAD_KIND_RETRY_DELAYS, threadKind, @@ -71,6 +74,73 @@ test("thread lookup resolves real main/subagent files and unknown safely", async ); }); +test("thread lookup falls back to archived sessions", async (t) => { + const paths = await temporaryPaths(); + t.after(() => removeTemporaryPaths(paths)); + const archivedRoot = join(dirname(paths.sessionRoot), "archived_sessions"); + await mkdir(archivedRoot, { recursive: true }); + const archivedTranscript = join( + archivedRoot, + "rollout-archived-thread.jsonl", + ); + await writeFile( + archivedTranscript, + jsonl({ type: "session_meta", payload: { id: "archived-thread" } }), + ); + + assert.equal( + await sessionPathForThreadId("archived-thread", paths.sessionRoot), + archivedTranscript, + ); + assert.equal( + await threadKind({ "thread-id": "archived-thread" }, paths), + "main", + ); +}); + +test("thread lookup prefers a supplied transcript and recovers from a stale path", async (t) => { + const paths = await temporaryPaths(); + t.after(() => removeTemporaryPaths(paths)); + const threadId = "preferred-thread"; + const discoveredTranscript = join( + paths.sessionRoot, + `rollout-discovered-${threadId}.jsonl`, + ); + const preferredTranscript = join( + paths.root, + `rollout-preferred-${threadId}.jsonl`, + ); + await writeFile( + discoveredTranscript, + jsonl({ + type: "session_meta", + payload: { id: threadId, parent_thread_id: "parent-thread" }, + }), + ); + await writeFile( + preferredTranscript, + jsonl({ type: "session_meta", payload: { id: threadId } }), + ); + + assert.equal( + await threadKind( + { "thread-id": threadId, transcript_path: preferredTranscript }, + paths, + ), + "main", + ); + assert.equal( + await threadKind( + { + "thread-id": threadId, + transcript_path: join(paths.root, "moved-away.jsonl"), + }, + paths, + ), + "subagent", + ); +}); + test("resolvedThreadKind retries until session metadata appears", async (t) => { const paths = await temporaryPaths(); t.after(() => removeTemporaryPaths(paths)); @@ -355,3 +425,68 @@ test("assistant summary falls back when the completed turn is absent or has no f ); } }); + +test("large transcripts stay readable across reverse chunk boundaries", async (t) => { + const paths = await temporaryPaths(); + t.after(() => removeTemporaryPaths(paths)); + const threadId = "large-thread"; + const transcript = join( + paths.sessionRoot, + `rollout-large-${threadId}.jsonl`, + ); + const filler = "边界填充".repeat(24 * 1024); + const records = [ + { type: "session_meta", payload: { id: threadId } }, + { + type: "event_msg", + payload: { type: "user_message", message: "旧的对话问题" }, + }, + ...Array.from({ length: 12 }, (_, index) => ({ + type: "event_msg", + payload: { type: "token_count", index, filler }, + })), + { + type: "event_msg", + payload: { type: "user_message", message: "核验超长会话通知" }, + }, + { + type: "event_msg", + payload: { + type: "task_complete", + turn_id: "target-turn", + last_agent_message: "超长会话现在可以正确发送通知。", + }, + }, + ...Array.from({ length: 12 }, (_, index) => ({ + type: "event_msg", + payload: { type: "token_count", index: index + 12, filler }, + })), + { + type: "event_msg", + payload: { + type: "task_complete", + turn_id: "newer-turn", + last_agent_message: "另一轮的结果不能覆盖目标轮。", + }, + }, + ]; + await writeFile(transcript, jsonl(...records)); + + assert.deepEqual( + await sessionMetaFromTranscript(transcript, threadId), + { id: threadId }, + ); + assert.equal( + await conversationNameFromTranscript(transcript), + "核验超长会话通知", + ); + assert.equal( + await assistantSummaryFromTranscript( + transcript, + "target-turn", + "摘要回退", + { retryDelays: [] }, + ), + "超长会话现在可以正确发送通知。", + ); +}); diff --git a/test/state.test.mjs b/test/state.test.mjs index db1e6d8..0921a85 100644 --- a/test/state.test.mjs +++ b/test/state.test.mjs @@ -13,6 +13,7 @@ import test from "node:test"; import { acquireEventLock, + cleanupExpiredSentMarkers, cleanupStaleJobs, consumePrivateJob, createPrivateJob, @@ -20,6 +21,7 @@ import { markEventSent, readPrivateJob, releaseEventLock, + rotateAuditLog, writeAudit, } from "../src/lib/state.mjs"; import { @@ -71,6 +73,34 @@ test("a stale lock is recovered exactly once", async (t) => { assert.equal((await lstat(recovered.lockPath)).isFile(), true); }); +test("sent markers expire after 90 days without weakening recent deduplication", async (t) => { + const paths = await temporaryPaths(); + t.after(() => removeTemporaryPaths(paths)); + const oldPayload = { "thread-id": "thread", "turn-id": "old" }; + const recentPayload = { "thread-id": "thread", "turn-id": "recent" }; + const oldLock = await acquireEventLock(oldPayload, paths); + const recentLock = await acquireEventLock(recentPayload, paths); + await markEventSent(oldLock); + await markEventSent(recentLock); + + const now = Date.now(); + const old = new Date(now - 91 * 24 * 60 * 60 * 1_000); + await utimes(oldLock.sentPath, old, old); + assert.equal( + await cleanupExpiredSentMarkers(paths, { + now, + minimumIntervalMilliseconds: 0, + }), + 1, + ); + await assert.rejects(lstat(oldLock.sentPath), { code: "ENOENT" }); + assert.equal((await lstat(recentLock.sentPath)).isFile(), true); + + const oldAgain = await acquireEventLock(oldPayload, paths); + assert.equal(oldAgain.duplicate, false); + assert.equal((await acquireEventLock(recentPayload, paths)).duplicate, true); +}); + test("private jobs use 0700 directory and 0600 regular files", async (t) => { const paths = await temporaryPaths(); t.after(() => removeTemporaryPaths(paths)); @@ -301,3 +331,24 @@ test("audit log excludes keys, conversation text, and arbitrary details", async ); assert.equal(permissions((await lstat(paths.auditLog)).mode), 0o600); }); + +test("audit log rotates to one private archive before appending", async (t) => { + const paths = await temporaryPaths(); + t.after(() => removeTemporaryPaths(paths)); + await writeFile(paths.auditLog, "old audit records\n", { mode: 0o644 }); + + await writeAudit( + paths, + { "thread-id": "thread", "turn-id": "turn" }, + "sent", + {}, + { maximumBytes: 8 }, + ); + + assert.equal(await readFile(paths.auditArchive, "utf8"), "old audit records\n"); + assert.equal(permissions((await lstat(paths.auditArchive)).mode), 0o600); + assert.equal(permissions((await lstat(paths.auditLog)).mode), 0o600); + const current = JSON.parse((await readFile(paths.auditLog, "utf8")).trim()); + assert.equal(current.event, "sent"); + assert.equal(await rotateAuditLog(paths, { maximumBytes: 10_000 }), false); +}); diff --git a/test/text.test.mjs b/test/text.test.mjs index b502416..c9818c1 100644 --- a/test/text.test.mjs +++ b/test/text.test.mjs @@ -198,6 +198,13 @@ test("reply classification covers real continuation requests without optional fa "有一个隐私变化需要你确认。你确认允许保存后,我就可以实现。", "不用一次发齐,可以逐个平台发。最优先发未来30天内最早到期的那笔。", "两笔账单已经并入。下一张优先发京东金条的还款计划。", + "方案已经整理。你确认后我再修改配置。", + "草稿已经准备好。你确认后我将提交。", + "剩余步骤已列出。你确认后我继续发布。", + "草稿已经准备好。等你审阅并确认后,我再发布。", + "现在需要先确认三件事,才能继续处理。", + "需要先确认页面状态,才能得出结论。", + "现在需要先确认四件事:第18间房是谁、201押金到底是多少、40元是什么、退押金300走现金还是微信。确认后才能得出准确利润和剩余现金。", ]; for (const reply of replies) { assert.deepEqual(classifyLastReply(reply), { @@ -220,6 +227,8 @@ test("reply classification covers real continuation requests without optional fa "下一张截图将自动发送。", "请先做一次验收。后来我已代你完成,不需要你操作。", "请先点击确认;刚刚已经完成,不需要你操作。", + "建议你确认后再继续使用,本轮无需回复。", + "这里说明需要先确认参数才能继续,规则已经通过。", ]) { assert.deepEqual(classifyLastReply(reply), { icon: "✅", @@ -323,6 +332,102 @@ test("reply notification body prioritizes the concrete user action", () => { ), "下一张优先发京东金条64,552元的还款计划。", ); + assert.equal( + notificationBodyFromAssistantReply( + "方案已整理。请回复:继续安装、稍后处理,或取消。", + status, + "选择安装方案", + ), + "请回复:继续安装、稍后处理,或取消。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "[NEEDS_INPUT] 请回复:继续安装或停止。", + status, + "选择安装方案", + ), + "请回复:继续安装或停止。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "请回复。", + status, + "选择安装方案", + ), + "需要你回复后才能继续。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "请回复:", + status, + "选择安装方案", + ), + "需要你回复后才能继续。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "请回复:https://example.com", + status, + "选择安装方案", + ), + "需要你回复后才能继续。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "若同意以上规则,请回复:\n\n`封面默认不插入视频,按这些规则修改 Skill。`", + status, + "确认视频规则", + ), + "若同意以上规则,请回复:封面默认不插入视频,按这些规则修改 Skill。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "正式文件是 [SKILL.md](/Users/example/private/SKILL.md)。\n\n若同意以上规则,请回复:\n\n`按这些规则修改 Skill。`", + status, + "确认 Skill 规则", + ), + "若同意以上规则,请回复:按这些规则修改 Skill。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "草稿已完成。等你审阅并确认后,我再修改正式版本。", + status, + "审阅草稿", + ), + "等你审阅并确认后,我再修改正式版本。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "现在需要先确认四件事:第18间房是谁、201押金到底是多少、40元是什么、退押金300走现金还是微信。确认后才能得出准确结果。", + status, + "核对账目", + ), + "请确认:第18间房是谁、201押金到底是多少、40元是什么、退押金300走现金还是微信。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "你确认后,我再修改 /Users/example/private/config.json 里的设置。", + status, + "更新通知配置", + ), + "请确认是否继续处理本地文件。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "你确认后,我就点击 Relink File,选择 [本地素材](/Users/example/private/video.mp4)。", + status, + "重新关联素材", + ), + "请确认是否继续处理本地文件。", + ); + assert.equal( + notificationBodyFromAssistantReply( + "你确认后,我就点击 Relink File,选择:\n\n`/Users/example/private/video.mp4`", + status, + "重新关联素材", + ), + "请确认是否继续处理本地文件。", + ); }); test("Unicode truncation counts code points instead of UTF-16 units", () => { @@ -661,6 +766,25 @@ test("assistant summary removes hidden markers and avoids cutting a complete cla ), "点击 Bark 通知可以直接进入 ChatGPT 的 Codex Remote 界面。", ); + for (const [marker, expected] of [ + ["COMPLETE", "通知优化完成。"], + ["SUCCESS", "通知测试通过。"], + ["BLOCKED", "当前无法继续。"], + ["NEEDS_INPUT", "请确认是否继续。"], + ]) { + assert.equal( + shortenAssistantSummary(`[${marker}] ${expected}`, "通知状态"), + expected, + ); + } + assert.deepEqual(classifyLastReply("[BLOCKED] 当前无法继续。"), { + icon: "⛔", + label: "受阻或出错", + }); + assert.deepEqual(classifyLastReply("[NEEDS_INPUT] 请确认是否继续。"), { + icon: "🔁", + label: "需要你回复", + }); assert.equal( shortenAssistantSummary( "通知结果已确认:点击 Bark 通知可以直接进入 ChatGPT 的 Codex RemoteSettingsPanel 并显示结果。", @@ -774,6 +898,25 @@ test("credential redaction covers Bark keys with natural-language separators", ( ); }); +test("natural Chinese result sentences are not mistaken for standalone Bark keys", () => { + for (const reply of [ + "原因找到了:网页端读不到本地4K音轨。", + "网页端声音已修复,4K画面保持不变。", + "建议取消1400元Pro,但不要和别人共用账号。", + ]) { + const summary = shortenAssistantSummary(reply, "通知结果"); + assert.notEqual(summary, "未生成摘要"); + assert.equal( + formatNotification( + { icon: "✅", label: "本轮结束" }, + "通知检查", + summary, + ).body, + `💬${summary}`, + ); + } +}); + test("conversation names redact credentials and safely fall back when credential-only", () => { assert.equal( shortenConversationName(