mirror of
https://github.com/NanmiCoder/claude-code-haha.git
synced 2026-10-10 11:53:10 +08:00
fix: discard stale transcript queues after retention cleanup
This commit is contained in:
@@ -1,32 +1,53 @@
|
||||
import { expect, test } from 'bun:test'
|
||||
import { mkdir, mkdtemp, readFile, rm, unlink, writeFile } from 'node:fs/promises'
|
||||
import { mkdir, mkdtemp, readFile, rename, rm, unlink, writeFile } from 'node:fs/promises'
|
||||
import { join } from 'node:path'
|
||||
import { createSandboxedTestEnvironment } from '../../../scripts/pr/test-environment.js'
|
||||
|
||||
// Exercise the registered runtime shutdown callback in another process. That
|
||||
// process caches enabled settings and never observes the parent's 365→0→365.
|
||||
for (const deleted of [true, false]) {
|
||||
test(`runtime exit ${deleted ? 'does not recreate a deleted transcript' : 'preserves metadata for an existing transcript'}`, async () => {
|
||||
// Exercise the registered shutdown callback in another process. It caches
|
||||
// enabled settings and never observes the parent's intervening 365→0→365.
|
||||
const scenarios = [
|
||||
'deleted-idle',
|
||||
'kept-idle',
|
||||
'deleted-queued',
|
||||
'first-turn-exit',
|
||||
'deleted-then-new-turn',
|
||||
'deleted-queued-then-new-turn',
|
||||
'replaced-queued-then-new-turn',
|
||||
] as const
|
||||
for (const scenario of scenarios) {
|
||||
test(`runtime exit retention lifecycle: ${scenario}`, async () => {
|
||||
const directory = await mkdtemp('/tmp/session-exit-retention-')
|
||||
const env = createSandboxedTestEnvironment(directory, { TEST_ENABLE_SESSION_PERSISTENCE: '1' })
|
||||
const configDir = env.CLAUDE_CONFIG_DIR!
|
||||
await mkdir(configDir, { recursive: true })
|
||||
await writeFile(join(configDir, 'settings.json'), JSON.stringify({ cleanupPeriodDays: 365 }))
|
||||
const source = `
|
||||
import { mkdir } from 'node:fs/promises'
|
||||
import { switchSession } from './src/bootstrap/state.ts'
|
||||
import { getSettings_DEPRECATED } from './src/utils/settings/settings.ts'
|
||||
import { runCleanupFunctions } from './src/utils/cleanupRegistry.ts'
|
||||
import { cacheSessionTitle, flushSessionStorage, getTranscriptPathForSession, recordTranscript } from './src/utils/sessionStorage.ts'
|
||||
const originalSetTimeout = globalThis.setTimeout
|
||||
// Hold the normal 100 ms transcript queue until the real exit callback
|
||||
// drains it, making the external deletion race deterministic.
|
||||
globalThis.setTimeout = (callback, delay, ...args) => originalSetTimeout(callback, delay === 100 ? 60000 : delay, ...args)
|
||||
const scenario = ${JSON.stringify(scenario)}
|
||||
const id = 'deadbeef-0000-4000-8000-000000000001'
|
||||
const old = {type:'user', uuid:'deadbeef-0000-4000-8000-000000000002', timestamp:'2026-09-10T00:00:00.000Z', message:{role:'user',content:'CACHED EXIT PRIVATE PROMPT'}}
|
||||
const queued = {type:'user', uuid:'deadbeef-0000-4000-8000-000000000003', timestamp:'2026-09-10T00:00:01.000Z', message:{role:'user',content:'QUEUED BEFORE DELETE'}}
|
||||
const fresh = {type:'user', uuid:'deadbeef-0000-4000-8000-000000000004', timestamp:'2026-09-10T00:00:02.000Z', message:{role:'user',content:'FRESH EXPLICIT TURN'}}
|
||||
switchSession(id)
|
||||
await mkdir(process.env.CLAUDE_CONFIG_DIR, { recursive: true })
|
||||
getSettings_DEPRECATED()
|
||||
cacheSessionTitle('CACHED EXIT TITLE')
|
||||
await recordTranscript([{type:'user', uuid:'deadbeef-0000-4000-8000-000000000002', timestamp:'2026-09-10T00:00:00.000Z', message:{role:'user',content:'CACHED EXIT PRIVATE PROMPT'}}])
|
||||
await flushSessionStorage()
|
||||
process.stdout.write(JSON.stringify({ path: getTranscriptPathForSession(id) }) + '\\n')
|
||||
if (scenario.endsWith('idle')) cacheSessionTitle('CACHED EXIT TITLE')
|
||||
if (scenario !== 'first-turn-exit') {
|
||||
await recordTranscript([old])
|
||||
await flushSessionStorage()
|
||||
}
|
||||
if (scenario.includes('queued')) await recordTranscript([old, queued])
|
||||
if (scenario === 'first-turn-exit') await recordTranscript([fresh])
|
||||
process.stdout.write(JSON.stringify({ path: getTranscriptPathForSession(id), cachedDays: getSettings_DEPRECATED().cleanupPeriodDays }) + '\\n')
|
||||
await new Promise(resolve => process.stdin.once('data', resolve))
|
||||
if (scenario.endsWith('new-turn')) await recordTranscript(scenario.includes('queued') ? [old, queued, fresh] : [old, fresh])
|
||||
if (getSettings_DEPRECATED().cleanupPeriodDays !== 365) throw new Error('Child observed disabled settings')
|
||||
// This is the same cleanup registry invoked by gracefulShutdown.
|
||||
await runCleanupFunctions()
|
||||
process.exit(0)
|
||||
@@ -43,26 +64,41 @@ for (const deleted of [true, false]) {
|
||||
if (result.done) throw new Error(`Runtime exited before ready: ${await errorOutput}`)
|
||||
ready += new TextDecoder().decode(result.value)
|
||||
}
|
||||
const { path } = JSON.parse(ready.trim()) as { path: string }
|
||||
const initial = await readFile(path, 'utf8')
|
||||
expect(initial).toContain('CACHED EXIT PRIVATE PROMPT')
|
||||
expect(initial).toContain('CACHED EXIT TITLE')
|
||||
const { path, cachedDays } = JSON.parse(ready.trim()) as { path: string; cachedDays: number }
|
||||
expect(cachedDays).toBe(365)
|
||||
const initial = await readFile(path, 'utf8').catch(() => '')
|
||||
expect(initial).not.toContain('QUEUED BEFORE DELETE')
|
||||
if (scenario === 'first-turn-exit') expect(initial).toBe('')
|
||||
else expect(initial).toContain('CACHED EXIT PRIVATE PROMPT')
|
||||
if (scenario.endsWith('idle')) expect(initial).toContain('CACHED EXIT TITLE')
|
||||
expect(initial).not.toContain('last-prompt')
|
||||
if (deleted) {
|
||||
if (scenario.startsWith('deleted')) {
|
||||
await writeFile(join(configDir, 'settings.json'), JSON.stringify({ cleanupPeriodDays: 0 }))
|
||||
await unlink(path)
|
||||
await writeFile(join(configDir, 'settings.json'), JSON.stringify({ cleanupPeriodDays: 365 }))
|
||||
}
|
||||
if (scenario.startsWith('replaced')) {
|
||||
// Keep the old inode allocated so replacement detection is deterministic.
|
||||
await rename(path, path + '.removed')
|
||||
await writeFile(path, '')
|
||||
}
|
||||
child.stdin.write('exit\n')
|
||||
child.stdin.end()
|
||||
expect(await child.exited).toBe(0)
|
||||
expect(await errorOutput).toBe('')
|
||||
if (deleted) {
|
||||
expect(await readFile(path, 'utf8').catch(() => null)).toBeNull()
|
||||
} else {
|
||||
if (scenario === 'deleted-idle' || scenario === 'deleted-queued') {
|
||||
await expect(readFile(path, 'utf8')).rejects.toMatchObject({ code: 'ENOENT' })
|
||||
} else if (scenario === 'kept-idle') {
|
||||
const final = await readFile(path, 'utf8')
|
||||
expect(final).toContain('"lastPrompt":"CACHED EXIT PRIVATE PROMPT"')
|
||||
expect(final.match(/CACHED EXIT TITLE/g)).toHaveLength(2)
|
||||
} else {
|
||||
const final = await readFile(path, 'utf8')
|
||||
expect(final).not.toContain('CACHED EXIT PRIVATE PROMPT')
|
||||
expect(final).not.toContain('QUEUED BEFORE DELETE')
|
||||
const messages = final.trim().split('\n').map(line => JSON.parse(line)).filter(entry => entry.type === 'user')
|
||||
expect(messages).toHaveLength(1)
|
||||
expect(messages[0]).toMatchObject({ uuid: 'deadbeef-0000-4000-8000-000000000004', parentUuid: null })
|
||||
}
|
||||
} finally {
|
||||
child.kill()
|
||||
|
||||
@@ -15,6 +15,18 @@ import {
|
||||
recordSidechainTranscript,
|
||||
getAgentTranscriptPath,
|
||||
resetProjectForTesting,
|
||||
getCurrentSessionTitle,
|
||||
saveAiGeneratedTitle,
|
||||
saveTaskSummary,
|
||||
saveCustomTitle,
|
||||
saveTag,
|
||||
saveAgentName,
|
||||
saveAgentColor,
|
||||
saveWorktreeState,
|
||||
linkSessionToPR,
|
||||
adoptResumedSessionFile,
|
||||
recordContentReplacement,
|
||||
removeTranscriptMessage,
|
||||
} from '../sessionStorage.js'
|
||||
import { resetSettingsCache, setSessionSettingsCache } from '../settings/settingsCache.js'
|
||||
|
||||
@@ -174,6 +186,104 @@ describe('session retention transitions', () => {
|
||||
expect(JSON.parse(raw.trim())).toMatchObject({ uuid: fresh.uuid, parentUuid: null })
|
||||
})
|
||||
|
||||
async function updatePrivateMetadata() {
|
||||
saveAiGeneratedTitle(sessionId, 'PRIVATE AI TITLE')
|
||||
saveTaskSummary(sessionId, 'PRIVATE TASK SUMMARY')
|
||||
await saveCustomTitle(sessionId, 'PRIVATE CUSTOM TITLE')
|
||||
await saveTag(sessionId, 'PRIVATE TAG')
|
||||
await saveAgentName(sessionId, 'PRIVATE AGENT')
|
||||
await saveAgentColor(sessionId, 'PRIVATE COLOR')
|
||||
await linkSessionToPR(sessionId, 1, 'https://example.invalid/PRIVATE', 'PRIVATE REPO')
|
||||
saveWorktreeState(null)
|
||||
}
|
||||
|
||||
it('removes an older failed message without losing later transcript content during a full rewrite', async () => {
|
||||
const failed = user('FAILED STREAM MESSAGE')
|
||||
const later = user('later retained content '.repeat(4000))
|
||||
await recordTranscript([failed, later] as never[])
|
||||
await flushSessionStorage()
|
||||
const transcriptPath = getTranscriptPathForSession(sessionId)
|
||||
// Preserve an unparseable line as well: tombstoning only removes the
|
||||
// selected UUID, even when the target is outside the tail read window.
|
||||
await fs.appendFile(transcriptPath, 'incomplete trailing entry\n')
|
||||
await removeTranscriptMessage(failed.uuid)
|
||||
const saved = await fs.readFile(transcriptPath, 'utf8')
|
||||
expect(saved).not.toContain('FAILED STREAM MESSAGE')
|
||||
const lines = saved.trim().split('\n')
|
||||
expect(JSON.parse(lines[0]!)).toMatchObject({ uuid: later.uuid, message: later.message })
|
||||
expect(lines[1]).toBe('incomplete trailing entry')
|
||||
expect(lines).toHaveLength(2)
|
||||
})
|
||||
|
||||
for (const replaceFile of [false, true]) {
|
||||
it(`binds resumed content replacement to its existing file (replaced=${replaceFile})`, async () => {
|
||||
const transcriptPath = getTranscriptPathForSession(sessionId)
|
||||
await fs.mkdir(path.dirname(transcriptPath), { recursive: true })
|
||||
await fs.writeFile(transcriptPath, '')
|
||||
adoptResumedSessionFile()
|
||||
const originalSetTimeout = globalThis.setTimeout
|
||||
globalThis.setTimeout = ((callback: any, delay?: number, ...args: any[]) =>
|
||||
originalSetTimeout(callback, delay === 100 ? 60_000 : delay, ...args)) as typeof setTimeout
|
||||
try {
|
||||
await recordContentReplacement([{ kind: 'tool-result', toolUseId: 'tool-private', replacement: 'PRIVATE REPLACEMENT' }])
|
||||
if (replaceFile) {
|
||||
await fs.rename(transcriptPath, transcriptPath + '.removed')
|
||||
await fs.writeFile(transcriptPath, '')
|
||||
}
|
||||
await flushSessionStorage()
|
||||
const saved = await fs.readFile(transcriptPath, 'utf8')
|
||||
if (replaceFile) expect(saved).toBe('')
|
||||
else expect(JSON.parse(saved)).toMatchObject({ type: 'content-replacement', replacements: [{ replacement: 'PRIVATE REPLACEMENT' }] })
|
||||
} finally {
|
||||
globalThis.setTimeout = originalSetTimeout
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
it('synchronous metadata helpers do not modify saved content while retention is zero', async () => {
|
||||
await recordTranscript([user('public before private metadata')] as never[])
|
||||
await flushSessionStorage()
|
||||
const transcriptPath = getTranscriptPathForSession(sessionId)
|
||||
const before = await fs.readFile(transcriptPath, 'utf8')
|
||||
retention(0)
|
||||
await updatePrivateMetadata()
|
||||
expect(await fs.readFile(transcriptPath, 'utf8')).toBe(before)
|
||||
expect(getCurrentSessionTitle(sessionId)).toBe('PRIVATE CUSTOM TITLE')
|
||||
})
|
||||
|
||||
it('synchronous metadata helpers do not recreate a removed transcript with enabled cached settings', async () => {
|
||||
await recordTranscript([user('public removed transcript')] as never[])
|
||||
await flushSessionStorage()
|
||||
const transcriptPath = getTranscriptPathForSession(sessionId)
|
||||
await fs.unlink(transcriptPath)
|
||||
await updatePrivateMetadata()
|
||||
await expect(fs.readFile(transcriptPath, 'utf8')).rejects.toMatchObject({ code: 'ENOENT' })
|
||||
})
|
||||
|
||||
it('does not reappend private metadata caches when new public records materialize after deletion', async () => {
|
||||
const old = user('old removed history')
|
||||
await recordTranscript([old] as never[])
|
||||
await flushSessionStorage()
|
||||
const transcriptPath = getTranscriptPathForSession(sessionId)
|
||||
retention(0)
|
||||
await updatePrivateMetadata()
|
||||
await fs.unlink(transcriptPath)
|
||||
retention(365)
|
||||
await recordTranscript([old, user('new public history')] as never[])
|
||||
await flushSessionStorage()
|
||||
reAppendSessionMetadata()
|
||||
const saved = await fs.readFile(transcriptPath, 'utf8')
|
||||
expect(saved).toContain('new public history')
|
||||
expect(saved).not.toContain('PRIVATE')
|
||||
await saveCustomTitle(sessionId, 'new public title')
|
||||
await saveTag(sessionId, 'new public tag')
|
||||
reAppendSessionMetadata()
|
||||
const updated = await fs.readFile(transcriptPath, 'utf8')
|
||||
expect(updated).not.toContain('PRIVATE')
|
||||
expect(updated.match(/new public title/g)).toHaveLength(2)
|
||||
expect(updated.match(/new public tag/g)).toHaveLength(2)
|
||||
})
|
||||
|
||||
it('does not rebuild a removed session through cached last-prompt metadata while disabled', async () => {
|
||||
await recordTranscript([user('old prompt')] as never[])
|
||||
await flushSessionStorage()
|
||||
|
||||
+203
-104
@@ -6,7 +6,6 @@ import type { Dirent } from 'fs'
|
||||
// with the async-suffixed names.
|
||||
import { appendFileSync as fsAppendFileSync, closeSync, constants, fchmodSync, fstatSync, openSync, readSync } from 'fs'
|
||||
import {
|
||||
appendFile as fsAppendFile,
|
||||
open as fsOpen,
|
||||
mkdir,
|
||||
readdir,
|
||||
@@ -533,7 +532,8 @@ export async function enqueueSessionEntryAfterPendingForTesting(
|
||||
if (delayMs > 0) {
|
||||
await new Promise(resolve => setTimeout(resolve, delayMs))
|
||||
}
|
||||
void projectForTesting.enqueueWrite(path, entry)
|
||||
await getProject().prepareTranscriptFile(path)
|
||||
await projectForTesting.enqueueWrite(path, entry)
|
||||
})
|
||||
}
|
||||
|
||||
@@ -604,6 +604,18 @@ class Project {
|
||||
// growing-history callers. Keep these separate from the on-disk dedup cache:
|
||||
// excluded messages cannot become parents, even after cache invalidation.
|
||||
private excludedMessages = new Map<string, Set<UUID>>()
|
||||
// UI metadata remains available in memory while recording is disabled,
|
||||
// but those values must not later be replayed by materialize/exit cleanup.
|
||||
private excludedMetadata = new Set<Entry['type']>()
|
||||
|
||||
markMetadataPersistence(type: Entry['type']): void {
|
||||
if (getSettings_DEPRECATED()?.cleanupPeriodDays === 0) this.excludedMetadata.add(type)
|
||||
else this.excludedMetadata.delete(type)
|
||||
}
|
||||
|
||||
resetMetadataPersistence(): void {
|
||||
this.excludedMetadata.clear()
|
||||
}
|
||||
|
||||
isMessageExcluded(uuid: UUID, sessionId = getSessionId()): boolean {
|
||||
return this.excludedMessages.get(sessionId)?.has(uuid) ?? false
|
||||
@@ -630,17 +642,20 @@ class Project {
|
||||
|
||||
async discardDeletedTranscriptHistory(messageSet: Set<UUID>, sessionId: UUID): Promise<void> {
|
||||
const filePath = this.sessionFile ?? getTranscriptPathForSession(sessionId)
|
||||
// An unflushed first turn has UUIDs in the dedup cache before its file
|
||||
// exists. Only treat a missing file as cleanup when no writes are pending.
|
||||
if (messageSet.size === 0 || this.activeDrain || this.writeQueues.get(filePath)?.length) return
|
||||
if (messageSet.size === 0 && !this.writeQueues.get(filePath)?.length) return
|
||||
try {
|
||||
await stat(filePath)
|
||||
const identity = await stat(filePath)
|
||||
const previous = this.transcriptFileIdentities.get(filePath)
|
||||
if (!previous || (previous.dev === identity.dev && previous.ino === identity.ino)) return
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== 'ENOENT') return
|
||||
this.excludeMessageUuids(messageSet, sessionId)
|
||||
messageSet.clear()
|
||||
this.currentSessionLastPrompt = undefined
|
||||
}
|
||||
this.excludeMessageUuids(messageSet, sessionId)
|
||||
messageSet.clear()
|
||||
const queue = this.writeQueues.get(filePath)
|
||||
if (queue) this.discardQueuedEntries(queue.splice(0))
|
||||
this.pendingEntries = []
|
||||
this.currentSessionLastPrompt = undefined
|
||||
}
|
||||
private remoteIngressUrl: string | null = null
|
||||
private internalEventWriter: InternalEventWriter | null = null
|
||||
@@ -648,12 +663,13 @@ class Project {
|
||||
private internalSubagentEventReader: InternalEventReader | null = null
|
||||
private pendingWriteCount: number = 0
|
||||
private flushResolvers: Array<() => void> = []
|
||||
// Per-file write queues. Each entry carries a resolve callback so
|
||||
// callers of enqueueWrite can optionally await their specific write.
|
||||
// Per-file queues capture the destination's identity before returning to
|
||||
// the caller. Actual disk writes are awaited by flush(), not enqueueWrite.
|
||||
private writeQueues = new Map<
|
||||
string,
|
||||
Array<{ entry: Entry; resolve: () => void }>
|
||||
Array<{ entry: Entry; identity: { dev: number; ino: number } }>
|
||||
>()
|
||||
private transcriptFileIdentities = new Map<string, { dev: number; ino: number }>()
|
||||
private flushTimer: ReturnType<typeof setTimeout> | null = null
|
||||
private activeDrain: Promise<void> | null = null
|
||||
private FLUSH_INTERVAL_MS = 100
|
||||
@@ -696,13 +712,33 @@ class Project {
|
||||
}
|
||||
|
||||
private enqueueWrite(filePath: string, entry: Entry): Promise<void> {
|
||||
return new Promise<void>(resolve => {
|
||||
return this.trackWrite(async () => {
|
||||
let identity = this.transcriptFileIdentities.get(filePath)
|
||||
if (!identity) {
|
||||
const noFollow = process.platform === 'win32' ? 0 : constants.O_NOFOLLOW
|
||||
let file: Awaited<ReturnType<typeof fsOpen>>
|
||||
try {
|
||||
file = await fsOpen(filePath, constants.O_WRONLY | constants.O_APPEND | noFollow)
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error
|
||||
this.discardQueuedEntries([{ entry }])
|
||||
return
|
||||
}
|
||||
try {
|
||||
const info = await file.stat()
|
||||
if (!info.isFile()) throw new Error(`Refusing non-regular transcript target: ${filePath}`)
|
||||
identity = { dev: info.dev, ino: info.ino }
|
||||
this.transcriptFileIdentities.set(filePath, identity)
|
||||
} finally {
|
||||
await file.close()
|
||||
}
|
||||
}
|
||||
let queue = this.writeQueues.get(filePath)
|
||||
if (!queue) {
|
||||
queue = []
|
||||
this.writeQueues.set(filePath, queue)
|
||||
}
|
||||
queue.push({ entry, resolve })
|
||||
queue.push({ entry, identity })
|
||||
this.scheduleDrain()
|
||||
})
|
||||
}
|
||||
@@ -723,17 +759,31 @@ class Project {
|
||||
}, this.FLUSH_INTERVAL_MS)
|
||||
}
|
||||
|
||||
private async appendToFile(filePath: string, data: string): Promise<void> {
|
||||
// Creation belongs to an explicit new message chain, never a delayed drain.
|
||||
// Keep the identity with each queued entry so a later turn recreating the
|
||||
// same path cannot receive a batch belonging to the deleted transcript.
|
||||
async prepareTranscriptFile(filePath: string): Promise<void> {
|
||||
if (this.shouldSkipPersistence()) return
|
||||
await mkdir(dirname(filePath), { recursive: true, mode: 0o700 })
|
||||
if (this.shouldSkipPersistence()) return
|
||||
const noFollow = process.platform === 'win32' ? 0 : constants.O_NOFOLLOW
|
||||
const file = await fsOpen(filePath, constants.O_WRONLY | constants.O_APPEND | constants.O_CREAT | noFollow, 0o600)
|
||||
try {
|
||||
await fsAppendFile(filePath, data, { mode: 0o600 })
|
||||
} catch {
|
||||
// Directory may not exist — some NFS-like filesystems return
|
||||
// unexpected error codes, so don't discriminate on code.
|
||||
await mkdir(dirname(filePath), { recursive: true, mode: 0o700 })
|
||||
await fsAppendFile(filePath, data, { mode: 0o600 })
|
||||
const identity = await file.stat()
|
||||
if (!identity.isFile()) throw new Error(`Refusing non-regular transcript target: ${filePath}`)
|
||||
this.transcriptFileIdentities.set(filePath, { dev: identity.dev, ino: identity.ino })
|
||||
} finally {
|
||||
await file.close()
|
||||
}
|
||||
}
|
||||
|
||||
private discardQueuedEntries(batch: Array<{ entry: Entry }>): void {
|
||||
for (const { entry } of batch) {
|
||||
if ('uuid' in entry) this.excludeMessageUuids([entry.uuid], entry.sessionId as SessionId)
|
||||
}
|
||||
this.currentSessionLastPrompt = undefined
|
||||
}
|
||||
|
||||
private async drainWriteQueue(): Promise<void> {
|
||||
for (const [filePath, queue] of this.writeQueues) {
|
||||
if (queue.length === 0) {
|
||||
@@ -743,39 +793,41 @@ class Project {
|
||||
// Cleanup can disable recording while a batch is waiting for its timer.
|
||||
// Discard it instead of recreating a transcript removed by cleanup.
|
||||
if (this.shouldSkipPersistence()) {
|
||||
for (const { entry, resolve } of batch) {
|
||||
if ('uuid' in entry) this.filterPersistableMessages([entry as TranscriptMessage], entry.sessionId as SessionId)
|
||||
resolve()
|
||||
}
|
||||
this.currentSessionLastPrompt = undefined
|
||||
this.discardQueuedEntries(batch)
|
||||
continue
|
||||
}
|
||||
|
||||
let content = ''
|
||||
const resolvers: Array<() => void> = []
|
||||
|
||||
for (const { entry, resolve } of batch) {
|
||||
const line = jsonStringify(entry) + '\n'
|
||||
|
||||
if (content.length + line.length >= this.MAX_CHUNK_BYTES) {
|
||||
// Flush chunk and resolve its entries before starting a new one
|
||||
await this.appendToFile(filePath, content)
|
||||
for (const r of resolvers) {
|
||||
r()
|
||||
}
|
||||
resolvers.length = 0
|
||||
content = ''
|
||||
}
|
||||
|
||||
content += line
|
||||
resolvers.push(resolve)
|
||||
const noFollow = process.platform === 'win32' ? 0 : constants.O_NOFOLLOW
|
||||
let file: Awaited<ReturnType<typeof fsOpen>>
|
||||
try {
|
||||
file = await fsOpen(filePath, constants.O_WRONLY | constants.O_APPEND | noFollow)
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== 'ENOENT') throw error
|
||||
this.discardQueuedEntries(batch)
|
||||
continue
|
||||
}
|
||||
|
||||
if (content.length > 0) {
|
||||
await this.appendToFile(filePath, content)
|
||||
for (const r of resolvers) {
|
||||
r()
|
||||
try {
|
||||
const identity = await file.stat()
|
||||
if (!identity.isFile()) throw new Error(`Refusing non-regular transcript target: ${filePath}`)
|
||||
let content = ''
|
||||
for (const item of batch) {
|
||||
if (this.shouldSkipPersistence() ||
|
||||
item.identity.dev !== identity.dev || item.identity.ino !== identity.ino) {
|
||||
this.discardQueuedEntries([item])
|
||||
continue
|
||||
}
|
||||
const line = jsonStringify(item.entry) + '\n'
|
||||
if (content.length + line.length >= this.MAX_CHUNK_BYTES) {
|
||||
await file.appendFile(content)
|
||||
content = ''
|
||||
}
|
||||
content += line
|
||||
}
|
||||
if (content.length > 0) {
|
||||
await file.appendFile(content)
|
||||
}
|
||||
} finally {
|
||||
await file.close()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -792,6 +844,12 @@ class Project {
|
||||
this.pendingEntries = []
|
||||
}
|
||||
|
||||
adoptSessionFile(filePath: string): void {
|
||||
this.sessionFile = filePath
|
||||
const identity = readTranscriptFileIdentitySync(filePath)
|
||||
if (identity) this.transcriptFileIdentities.set(filePath, identity)
|
||||
}
|
||||
|
||||
/**
|
||||
* Re-append cached session metadata to the end of the transcript file.
|
||||
* This ensures metadata stays within the tail window that readLiteMetadata
|
||||
@@ -839,7 +897,7 @@ class Project {
|
||||
// and not "type":"tag" appearing inside a nested tool_use input that
|
||||
// happens to be JSON-serialized into a message.
|
||||
const tailLines = tail.split('\n')
|
||||
if (!skipTitleRefresh) {
|
||||
if (!skipTitleRefresh && !this.excludedMetadata.has('custom-title')) {
|
||||
const titleLine = tailLines.findLast(l =>
|
||||
l.startsWith('{"type":"custom-title"'),
|
||||
)
|
||||
@@ -855,7 +913,7 @@ class Project {
|
||||
}
|
||||
}
|
||||
const tagLine = tailLines.findLast(l => l.startsWith('{"type":"tag"'))
|
||||
if (tagLine) {
|
||||
if (tagLine && !this.excludedMetadata.has('tag')) {
|
||||
const tailTag = extractLastJsonStringField(tagLine, 'tag')
|
||||
// Same: tagSession(id, null) writes `tag:""` to clear.
|
||||
if (tailTag !== undefined) {
|
||||
@@ -875,49 +933,49 @@ class Project {
|
||||
}
|
||||
// Unconditional: cache was refreshed from tail above; re-append keeps
|
||||
// the entry at EOF so compaction-pushed content doesn't evict it.
|
||||
if (this.currentSessionTitle) {
|
||||
if (this.currentSessionTitle && !this.excludedMetadata.has('custom-title')) {
|
||||
appendEntryToFile(this.sessionFile, {
|
||||
type: 'custom-title',
|
||||
customTitle: this.currentSessionTitle,
|
||||
sessionId,
|
||||
}, allowCreate)
|
||||
}
|
||||
if (this.currentSessionTag) {
|
||||
if (this.currentSessionTag && !this.excludedMetadata.has('tag')) {
|
||||
appendEntryToFile(this.sessionFile, {
|
||||
type: 'tag',
|
||||
tag: this.currentSessionTag,
|
||||
sessionId,
|
||||
}, allowCreate)
|
||||
}
|
||||
if (this.currentSessionAgentName) {
|
||||
if (this.currentSessionAgentName && !this.excludedMetadata.has('agent-name')) {
|
||||
appendEntryToFile(this.sessionFile, {
|
||||
type: 'agent-name',
|
||||
agentName: this.currentSessionAgentName,
|
||||
sessionId,
|
||||
}, allowCreate)
|
||||
}
|
||||
if (this.currentSessionAgentColor) {
|
||||
if (this.currentSessionAgentColor && !this.excludedMetadata.has('agent-color')) {
|
||||
appendEntryToFile(this.sessionFile, {
|
||||
type: 'agent-color',
|
||||
agentColor: this.currentSessionAgentColor,
|
||||
sessionId,
|
||||
}, allowCreate)
|
||||
}
|
||||
if (this.currentSessionAgentSetting) {
|
||||
if (this.currentSessionAgentSetting && !this.excludedMetadata.has('agent-setting')) {
|
||||
appendEntryToFile(this.sessionFile, {
|
||||
type: 'agent-setting',
|
||||
agentSetting: this.currentSessionAgentSetting,
|
||||
sessionId,
|
||||
}, allowCreate)
|
||||
}
|
||||
if (this.currentSessionMode) {
|
||||
if (this.currentSessionMode && !this.excludedMetadata.has('mode')) {
|
||||
appendEntryToFile(this.sessionFile, {
|
||||
type: 'mode',
|
||||
mode: this.currentSessionMode,
|
||||
sessionId,
|
||||
}, allowCreate)
|
||||
}
|
||||
if (this.currentSessionWorktree !== undefined) {
|
||||
if (this.currentSessionWorktree !== undefined && !this.excludedMetadata.has('worktree-state')) {
|
||||
appendEntryToFile(this.sessionFile, {
|
||||
type: 'worktree-state',
|
||||
worktreeSession: this.currentSessionWorktree,
|
||||
@@ -926,6 +984,7 @@ class Project {
|
||||
}
|
||||
if (
|
||||
this.currentSessionPrNumber !== undefined &&
|
||||
!this.excludedMetadata.has('pr-link') &&
|
||||
this.currentSessionPrUrl &&
|
||||
this.currentSessionPrRepository
|
||||
) {
|
||||
@@ -1031,32 +1090,40 @@ class Project {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// Slow path: target was not in the last 64KB. Rare - requires many
|
||||
// large entries to have landed between the write and the tombstone.
|
||||
if (fileSize > MAX_TOMBSTONE_REWRITE_BYTES) {
|
||||
logForDebugging(
|
||||
`Skipping tombstone removal: session file too large (${formatFileSize(fileSize)})`,
|
||||
{ level: 'warn' },
|
||||
)
|
||||
return
|
||||
}
|
||||
const content = await fh.readFile({ encoding: 'utf-8' })
|
||||
const lines = content.split('\n').filter((line: string) => {
|
||||
if (!line.trim()) return true
|
||||
try {
|
||||
const entry = jsonParse(line)
|
||||
return entry.uuid !== targetUuid
|
||||
} catch {
|
||||
return true // Keep malformed lines
|
||||
}
|
||||
})
|
||||
// Keep the original descriptor: a replaced path must not receive
|
||||
// content read from the previous transcript. readFile advanced this
|
||||
// descriptor's offset, so rewrite with an explicit zero position.
|
||||
const replacement = Buffer.from(lines.join('\n'), 'utf8')
|
||||
await fh.truncate(0)
|
||||
let written = 0
|
||||
while (written < replacement.length) {
|
||||
const result = await fh.write(replacement, written, replacement.length - written, written)
|
||||
if (result.bytesWritten === 0) throw new Error('Incomplete transcript rewrite')
|
||||
written += result.bytesWritten
|
||||
}
|
||||
} finally {
|
||||
await fh.close()
|
||||
}
|
||||
|
||||
// Slow path: target was not in the last 64KB. Rare - requires many
|
||||
// large entries to have landed between the write and the tombstone.
|
||||
if (fileSize > MAX_TOMBSTONE_REWRITE_BYTES) {
|
||||
logForDebugging(
|
||||
`Skipping tombstone removal: session file too large (${formatFileSize(fileSize)})`,
|
||||
{ level: 'warn' },
|
||||
)
|
||||
return
|
||||
}
|
||||
const content = await readFile(this.sessionFile, { encoding: 'utf-8' })
|
||||
const lines = content.split('\n').filter((line: string) => {
|
||||
if (!line.trim()) return true
|
||||
try {
|
||||
const entry = jsonParse(line)
|
||||
return entry.uuid !== targetUuid
|
||||
} catch {
|
||||
return true // Keep malformed lines
|
||||
}
|
||||
})
|
||||
await writeFile(this.sessionFile, lines.join('\n'), {
|
||||
encoding: 'utf8',
|
||||
})
|
||||
} catch {
|
||||
// Silently ignore errors - the file might not exist yet
|
||||
}
|
||||
@@ -1092,6 +1159,7 @@ class Project {
|
||||
// and create a metadata-only file despite --no-session-persistence.
|
||||
if (this.shouldSkipPersistence()) return
|
||||
this.ensureCurrentSessionFile()
|
||||
await this.prepareTranscriptFile(this.sessionFile!)
|
||||
// Only a new user/assistant turn may create a transcript. Exit/compaction
|
||||
// metadata refreshes must never resurrect a file removed by cleanup.
|
||||
this.reAppendSessionMetadata(false, true)
|
||||
@@ -1127,6 +1195,12 @@ class Project {
|
||||
) {
|
||||
await this.materializeSessionFile()
|
||||
}
|
||||
if (messages.some(m => m.type === 'user' || m.type === 'assistant')) {
|
||||
const targetFile = isSidechain && agentId
|
||||
? getAgentTranscriptPath(asAgentId(agentId))
|
||||
: this.sessionFile
|
||||
if (targetFile) await this.prepareTranscriptFile(targetFile)
|
||||
}
|
||||
|
||||
// Get current git branch once for this message chain
|
||||
let gitBranch: string | undefined
|
||||
@@ -1281,46 +1355,46 @@ class Project {
|
||||
// Only load current session messages if needed
|
||||
if (entry.type === 'summary') {
|
||||
// Summaries can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'custom-title') {
|
||||
// Custom titles can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'ai-title') {
|
||||
// AI titles can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'last-prompt') {
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'task-summary') {
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'tag') {
|
||||
// Tags can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'agent-name') {
|
||||
// Agent names can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'agent-color') {
|
||||
// Agent colors can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'agent-setting') {
|
||||
// Agent settings can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'pr-link') {
|
||||
// PR links can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'file-history-snapshot') {
|
||||
// File history snapshots can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'attribution-snapshot') {
|
||||
// Attribution snapshots can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'speculation-accept') {
|
||||
// Speculation accept entries can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'mode') {
|
||||
// Mode entries can always be appended
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'worktree-state') {
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'content-replacement') {
|
||||
// Content replacement records can always be appended. Subagent records
|
||||
// go to the sidechain file (for AgentTool resume); main-thread
|
||||
@@ -1328,20 +1402,20 @@ class Project {
|
||||
const targetFile = entry.agentId
|
||||
? getAgentTranscriptPath(entry.agentId)
|
||||
: sessionFile
|
||||
void this.enqueueWrite(targetFile, entry)
|
||||
await this.enqueueWrite(targetFile, entry)
|
||||
} else if (entry.type === 'marble-origami-commit') {
|
||||
// Always append. Commit order matters for restore (later commits may
|
||||
// reference earlier commits' summary messages), so these must be
|
||||
// written in the order received and read back sequentially.
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else if (entry.type === 'marble-origami-snapshot') {
|
||||
// Always append. Last-wins on restore — later entries supersede.
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else {
|
||||
const messageSet = await getSessionMessages(sessionId)
|
||||
if (entry.type === 'queue-operation') {
|
||||
// Queue operations are always appended to the session file
|
||||
void this.enqueueWrite(sessionFile, entry)
|
||||
await this.enqueueWrite(sessionFile, entry)
|
||||
} else {
|
||||
// At this point, entry must be a TranscriptMessage (user/assistant/attachment/system)
|
||||
// All other entry types have been handled above
|
||||
@@ -1365,8 +1439,8 @@ class Project {
|
||||
// exhausts retries → gracefulShutdownSync(1). See inc-4718.
|
||||
const isNewUuid = !messageSet.has(entry.uuid)
|
||||
if (isAgentSidechain || isNewUuid) {
|
||||
// Enqueue write — appendToFile handles ENOENT by creating directories
|
||||
void this.enqueueWrite(targetFile, entry)
|
||||
// The message-chain entry point has already materialized the file.
|
||||
await this.enqueueWrite(targetFile, entry)
|
||||
|
||||
if (!isAgentSidechain) {
|
||||
// messageSet is main-file-authoritative. Sidechain entries go to a
|
||||
@@ -1663,7 +1737,7 @@ export async function resetSessionFilePointer() {
|
||||
*/
|
||||
export function adoptResumedSessionFile(): void {
|
||||
const project = getProject()
|
||||
project.sessionFile = getTranscriptPath()
|
||||
project.adoptSessionFile(getTranscriptPath())
|
||||
project.reAppendSessionMetadata(true)
|
||||
}
|
||||
|
||||
@@ -2704,11 +2778,26 @@ export async function fetchLogs(limit?: number): Promise<LogOption[]> {
|
||||
* a stale process may not have observed retention-zero before it was restored.
|
||||
*/
|
||||
/* eslint-disable custom-rules/no-sync-fs -- sync callers (exit cleanup, materialize) */
|
||||
function readTranscriptFileIdentitySync(fullPath: string): { dev: number; ino: number } | undefined {
|
||||
let fd: number | undefined
|
||||
try {
|
||||
const noFollow = process.platform === 'win32' ? 0 : constants.O_NOFOLLOW
|
||||
fd = openSync(fullPath, constants.O_RDONLY | noFollow)
|
||||
const info = fstatSync(fd)
|
||||
return info.isFile() ? { dev: info.dev, ino: info.ino } : undefined
|
||||
} catch {
|
||||
return undefined
|
||||
} finally {
|
||||
if (fd !== undefined) closeSync(fd)
|
||||
}
|
||||
}
|
||||
|
||||
function appendEntryToFile(
|
||||
fullPath: string,
|
||||
entry: Record<string, unknown>,
|
||||
allowCreate = true,
|
||||
allowCreate = false,
|
||||
): void {
|
||||
if (getSettings_DEPRECATED()?.cleanupPeriodDays === 0) return
|
||||
const fs = getFsImplementation()
|
||||
const line = jsonStringify(entry) + '\n'
|
||||
if (!allowCreate) {
|
||||
@@ -2784,6 +2873,7 @@ export async function saveCustomTitle(
|
||||
// Cache for current session only (for immediate visibility)
|
||||
if (sessionId === getSessionId()) {
|
||||
getProject().currentSessionTitle = customTitle
|
||||
getProject().markMetadataPersistence('custom-title')
|
||||
}
|
||||
logEvent('tengu_session_renamed', {
|
||||
source:
|
||||
@@ -2848,6 +2938,7 @@ export async function saveTag(sessionId: UUID, tag: string, fullPath?: string) {
|
||||
// Cache for current session only (for immediate visibility)
|
||||
if (sessionId === getSessionId()) {
|
||||
getProject().currentSessionTag = tag
|
||||
getProject().markMetadataPersistence('tag')
|
||||
}
|
||||
logEvent('tengu_session_tagged', {})
|
||||
}
|
||||
@@ -2878,6 +2969,7 @@ export async function linkSessionToPR(
|
||||
project.currentSessionPrNumber = prNumber
|
||||
project.currentSessionPrUrl = prUrl
|
||||
project.currentSessionPrRepository = prRepository
|
||||
project.markMetadataPersistence('pr-link')
|
||||
}
|
||||
logEvent('tengu_session_linked_to_pr', { prNumber })
|
||||
}
|
||||
@@ -2945,6 +3037,7 @@ export function restoreSessionMetadata(meta: {
|
||||
*/
|
||||
export function clearSessionMetadata(): void {
|
||||
const project = getProject()
|
||||
project.resetMetadataPersistence()
|
||||
project.currentSessionTitle = undefined
|
||||
project.currentSessionTag = undefined
|
||||
project.currentSessionAgentName = undefined
|
||||
@@ -2981,6 +3074,7 @@ export async function saveAgentName(
|
||||
// Cache for current session only (for immediate visibility)
|
||||
if (sessionId === getSessionId()) {
|
||||
getProject().currentSessionAgentName = agentName
|
||||
getProject().markMetadataPersistence('agent-name')
|
||||
void updateSessionName(agentName)
|
||||
}
|
||||
logEvent('tengu_agent_name_set', {
|
||||
@@ -3003,6 +3097,7 @@ export async function saveAgentColor(
|
||||
// Cache for current session only (for immediate visibility)
|
||||
if (sessionId === getSessionId()) {
|
||||
getProject().currentSessionAgentColor = agentColor
|
||||
getProject().markMetadataPersistence('agent-color')
|
||||
}
|
||||
logEvent('tengu_agent_color_set', {})
|
||||
}
|
||||
@@ -3014,6 +3109,7 @@ export async function saveAgentColor(
|
||||
*/
|
||||
export function saveAgentSetting(agentSetting: string): void {
|
||||
getProject().currentSessionAgentSetting = agentSetting
|
||||
getProject().markMetadataPersistence('agent-setting')
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -3023,6 +3119,7 @@ export function saveAgentSetting(agentSetting: string): void {
|
||||
*/
|
||||
export function cacheSessionTitle(customTitle: string): void {
|
||||
getProject().currentSessionTitle = customTitle
|
||||
getProject().markMetadataPersistence('custom-title')
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -3032,6 +3129,7 @@ export function cacheSessionTitle(customTitle: string): void {
|
||||
*/
|
||||
export function saveMode(mode: 'coordinator' | 'normal'): void {
|
||||
getProject().currentSessionMode = mode
|
||||
getProject().markMetadataPersistence('mode')
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -3061,6 +3159,7 @@ export function saveWorktreeState(
|
||||
: null
|
||||
const project = getProject()
|
||||
project.currentSessionWorktree = stripped
|
||||
project.markMetadataPersistence('worktree-state')
|
||||
// Write eagerly when the file already exists (mid-session enter/exit).
|
||||
// For --worktree startup, sessionFile is null — materializeSessionFile
|
||||
// will write it on the first message via reAppendSessionMetadata.
|
||||
|
||||
Reference in New Issue
Block a user