merge: integrate QA-003 transcript queue cleanup fix

This commit is contained in:
程序员阿江(Relakkes)
2026-09-10 18:56:29 +08:00
3 changed files with 368 additions and 123 deletions
+55 -19
View File
@@ -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
View File
@@ -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.