fix(session-title): harden async provider lifecycle

This commit is contained in:
Tianyi Cui
2026-07-21 12:08:00 +08:00
parent bb7db52941
commit a9d518a38e
25 changed files with 491 additions and 102 deletions

View File

@@ -8,9 +8,9 @@ Only text blocks from human `user/message` events are eligible. The first eligib
- `get(session)` folds the latest accepted title from a live or replayed log.
- `refresh(session, signal?)` materializes the fallback when needed, then explicitly runs the registered provider over the current eligible messages. Provider errors and caller cancellation reject.
- `register(provider)` installs the sole optional provider and returns its Cordis effect disposer. A second registration throws immediately; disposal aborts pending and active calls before another provider can register.
- `register(provider)` installs the sole optional provider and returns its awaitable Cordis effect disposer. A second registration throws immediately; disposal aborts pending and active calls, waits for their settlement, and only then permits another provider to register.
Automatic work never delays the main agent response. A provider starts after the matching `request/header` records the main request's exact route; its late completion joins an open turn or uses a flushed zero-step `session-title` turn through `ctx.sessions.appendOutOfBand()`. Automatic failures warn and retain the latest title. New all-message revisions, provider disposal, session disposal, and explicit refresh abort older work, and a stale completion cannot append.
Automatic work never delays the main agent response. A provider starts after the matching `request/header` records the main request's exact route; its late completion joins an open turn or uses a flushed zero-step `session-title` turn through `ctx.sessions.appendOutOfBand()`. Automatic failures warn and retain the latest title. New all-message revisions, provider disposal, session disposal, and explicit refresh abort older work, and a stale completion cannot append. Concurrent explicit refreshes reserve their order before fallback durability waits, so only the newest call may reach the provider. Service teardown cancels queued work and drains calls that ignore cancellation before unloading completes.
Forks inherit title events in their seed unchanged. The first-message cadence does not automatically retitle a child; the all-messages cadence may append a new revision after the child receives a later human prompt.

View File

@@ -3,7 +3,7 @@
* @module @deepseek-ai/dsh-session-title
*/
import { Context, Service } from 'cordis'
import { Context, FiberState, Service, type Fiber } from 'cordis'
import z from 'schemastery'
import type { Branded } from '@deepseek-ai/dsh-brand'
import { deepFreeze } from '@deepseek-ai/dsh-llm'
@@ -200,6 +200,8 @@ interface ResolvedConfig {
/** One exact provider registration generation. */
interface ProviderRegistration {
readonly provider: SessionTitleProvider
readonly active: Set<Promise<unknown>>
closing: boolean
}
/** Automatic work waiting for the matching main-request header. */
@@ -239,11 +241,15 @@ export class SessionTitleService extends Service {
})
private readonly config: ResolvedConfig
private readonly ownerFiber: Fiber
private registration: ProviderRegistration | undefined
private readonly work = new Map<Session, SessionTitleWorkState>()
private readonly lifetime = new AbortController()
private readonly inFlight = new Set<Promise<unknown>>()
constructor(ctx: Context, config: Config) {
super(ctx, 'sessionTitle')
this.ownerFiber = ctx.fiber
const candidate: unknown = config
if (candidate === null || typeof candidate !== 'object') {
throw new Error('session-title: configuration is required')
@@ -257,6 +263,18 @@ export class SessionTitleService extends Service {
}
this.config = deepFreeze({ ...value })
ctx.effect(() => async () => {
this.lifetime.abort(new Error('session-title service disposed'))
if (this.registration !== undefined) this.registration.closing = true
this.registration = undefined
for (const state of this.work.values()) {
delete state.pending
state.active?.controller.abort(new Error('session-title service disposed'))
}
await this.drain(this.inFlight)
this.work.clear()
}, 'sessionTitle lifecycle')
ctx.on('session/event', (session, event) => {
switch (event.type) {
case 'user/message':
@@ -295,15 +313,16 @@ export class SessionTitleService extends Service {
*/
async refresh(session: Session, signal?: AbortSignal): Promise<SessionTitleSnapshot | undefined> {
signal?.throwIfAborted()
this.assertServiceActive()
if (this.ctx.sessions.get(session.id) !== session) {
throw new Error(`session "${session.id}" is not live in this store`)
}
const fallback = await this.ensureFallback(session)
const registration = this.registration
if (registration === undefined) return fallback
const messages = collectSessionTitleMessages(session.events)
const latest = messages.at(-1)
if (latest === undefined) return fallback
if (registration === undefined || registration.closing || latest === undefined) {
return this.ensureFallback(session)
}
const state = this.stateFor(session)
const revision = this.supersede(state, 'explicit title refresh superseded older generation')
const work = this.activate({
@@ -313,42 +332,48 @@ export class SessionTitleService extends Service {
}, state, signal)
const config = session.requestHeader()?.config
const route = config === undefined ? undefined : { provider: config.provider, model: config.model }
return this.runProvider(session, work, route)
return this.startProvider(session, work, route)
}
/**
* Register the sole optional title provider. Disposal aborts its pending and
* active work before another provider may register.
* @param provider - provider identity, cadence, and generation function.
* @returns exact Cordis effect disposer for HMR-safe unregistration.
* @returns exact Cordis effect disposer, which settles after active calls quiesce.
*/
register(provider: SessionTitleProvider): () => void {
register(provider: SessionTitleProvider): () => Promise<void> {
this.validateProvider(provider)
if (this.registration !== undefined) {
throw new Error(`session-title provider "${this.registration.provider.id}" is already registered`)
}
const registration: ProviderRegistration = {
provider,
active: new Set(),
closing: false,
}
const dispose = this.ctx.effect(function* (this: SessionTitleService) {
this.registration = registration
yield () => {
this.registration = undefined
yield async () => {
registration.closing = true
for (const state of this.work.values()) {
delete state.pending
state.active?.controller.abort(new Error(`session-title provider "${provider.id}" was disposed`))
if (state.pending?.registration === registration) delete state.pending
if (state.active?.registration === registration) {
state.active.controller.abort(new Error(`session-title provider "${provider.id}" was disposed`))
}
}
await this.drain(registration.active)
if (this.registration === registration) this.registration = undefined
}
}.bind(this), 'sessionTitle.register()')
// eslint-disable-next-line @typescript-eslint/no-misused-promises -- exact effect disposer preserves owner teardown ordering
return dispose
}
/** Schedule fallback creation and any provider cadence for one eligible event. */
private onUserMessage(session: Session, event: Extract<SessionEvent, { type: 'user/message' }>): void {
if (!this.serviceActive()) return
if (event.data.source.kind !== 'user' || collectSessionTitleMessages([event]).length === 0) return
const registration = this.registration
if (registration !== undefined) {
if (registration !== undefined && !registration.closing) {
const messages = collectSessionTitleMessages(session.events, event.seq)
const shouldSchedule = registration.provider.automatic === 'all-user-messages'
|| (session.header.parentSession === undefined && messages.length === 1 && this.get(session) === undefined)
@@ -358,15 +383,19 @@ export class SessionTitleService extends Service {
state.pending = { registration, revision, throughSeq: event.seq }
}
}
queueMicrotask(() => {
void this.ensureFallback(session).catch((error: unknown) => {
this.defer(async () => {
try {
await this.ensureFallback(session)
} catch (error: unknown) {
if (!this.serviceActive()) return
this.ctx.logger.warn(`session "${session.id}": fallback title update failed: ${String(error)}`)
})
}
})
}
/** Start pending automatic work only after its exact main-request route is logged. */
private onRequestHeader(session: Session, event: Extract<SessionEvent, { type: 'request/header' }>): void {
if (!this.serviceActive()) return
const state = this.work.get(session)
const pending = state?.pending
if (state === undefined || pending === undefined || pending.throughSeq >= event.seq) return
@@ -375,16 +404,31 @@ export class SessionTitleService extends Service {
provider: event.data.header.config.provider,
model: event.data.header.config.model,
}
queueMicrotask(() => {
if (this.registration !== pending.registration || state.revision !== pending.revision) return
this.defer(async () => {
if (this.registration !== pending.registration
|| pending.registration.closing
|| this.work.get(session) !== state
|| state.revision !== pending.revision) return
const work = this.activate(pending, state)
void this.runProvider(session, work, route).catch((error: unknown) => {
if (work.signal.aborted) return
try {
await this.startProvider(session, work, route)
} catch (error: unknown) {
if (work.signal.aborted || !this.serviceActive()) return
this.ctx.logger.warn(`session "${session.id}": automatic title generation failed: ${String(error)}`)
})
}
})
}
/** Start one tracked provider call after publishing its active revision. */
private startProvider(
session: Session,
work: ActiveProviderWork,
route?: SessionTitleModelProvenance,
): Promise<SessionTitleSnapshot | undefined> {
const run = Promise.resolve().then(() => this.runProvider(session, work, route))
return this.track(run, work.registration)
}
/** Execute and durably accept one current provider revision. */
private async runProvider(
session: Session,
@@ -392,6 +436,7 @@ export class SessionTitleService extends Service {
route?: SessionTitleModelProvenance,
): Promise<SessionTitleSnapshot | undefined> {
try {
this.assertCurrent(session, work)
await this.ensureFallback(session)
this.assertCurrent(session, work)
const messages = collectSessionTitleMessages(session.events, work.throughSeq)
@@ -470,6 +515,7 @@ export class SessionTitleService extends Service {
/** Fail a completion whose provider, revision, session, or signal is stale. */
private assertCurrent(session: Session, work: ActiveProviderWork): void {
this.assertServiceActive()
work.signal.throwIfAborted()
const state = this.work.get(session)
/* v8 ignore next -- every supported supersession, provider disposal, and session disposal aborts
@@ -490,8 +536,8 @@ export class SessionTitleService extends Service {
): ActiveProviderWork {
const controller = new AbortController()
const signal = upstream === undefined
? controller.signal
: AbortSignal.any([controller.signal, upstream])
? AbortSignal.any([controller.signal, this.lifetime.signal])
: AbortSignal.any([controller.signal, this.lifetime.signal, upstream])
const work: ActiveProviderWork = { ...pending, controller, signal }
state.active = work
return work
@@ -515,6 +561,44 @@ export class SessionTitleService extends Service {
return state
}
/** Queue detached service work and retain it through service disposal. */
private defer(task: () => Promise<void>): void {
const run = Promise.resolve().then(async () => {
if (!this.serviceActive()) return
await task()
})
void this.track(run)
}
/** Retain one promise until settlement for service and optional provider teardown. */
private track<T>(run: Promise<T>, registration?: ProviderRegistration): Promise<T> {
this.inFlight.add(run)
registration?.active.add(run)
const settled = (): void => {
this.inFlight.delete(run)
registration?.active.delete(run)
}
void run.then(settled, settled)
return run
}
/** Await every current and settling promise in one lifecycle registry. */
private async drain(active: Set<Promise<unknown>>): Promise<void> {
while (active.size > 0) await Promise.allSettled([...active])
}
/** Whether the owning plugin fiber can still start or commit title work. */
private serviceActive(): boolean {
return !this.lifetime.signal.aborted
&& this.ownerFiber.uid !== null
&& this.ownerFiber.state === FiberState.ACTIVE
}
/** Reject work once the owning plugin fiber has begun unloading. */
private assertServiceActive(): void {
if (!this.serviceActive()) throw new Error('session-title service disposed')
}
/** Reject malformed provider registrations before publishing an effect. */
private validateProvider(provider: unknown): asserts provider is SessionTitleProvider {
if (provider === null || typeof provider !== 'object') {
@@ -534,6 +618,7 @@ export class SessionTitleService extends Service {
/** Create the first deterministic fallback if the session still lacks a title. */
private async ensureFallback(session: Session): Promise<SessionTitleSnapshot | undefined> {
this.assertServiceActive()
const current = this.get(session)
if (current !== undefined) return current
const [first] = collectSessionTitleMessages(session.events)

View File

@@ -84,7 +84,7 @@ describe('SessionTitleService provider lifecycle', () => {
await settle()
child.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
expect(firstGenerate).not.toHaveBeenCalled()
disposeFirst()
await disposeFirst()
const allGenerate = vi.fn(async (request: SessionTitleProviderRequest) => ({
title: 'Fork all prompts',
@@ -170,7 +170,7 @@ describe('SessionTitleService provider lifecycle', () => {
expect(requests[1]?.messages.map(message => message.seq)).toEqual([first.seq, second.seq])
})
it('rejects a second provider and aborts stale work when the winner is disposed', async () => {
it('rejects a second provider and drains stale work when the winner is disposed', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
await ctx.plugin(SessionTitleService, CONFIG)
@@ -202,10 +202,15 @@ describe('SessionTitleService provider lifecycle', () => {
await settle()
expect(observedSignal?.aborted).toBe(false)
dispose()
const disposal = dispose()
expect(observedSignal?.aborted).toBe(true)
pending.resolve({ title: 'stale provider result', messageSeqs: [message.seq] })
let disposed = false
void disposal.then(() => { disposed = true })
await settle()
expect(disposed).toBe(false)
pending.resolve({ title: 'stale provider result', messageSeqs: [message.seq] })
await disposal
expect(disposed).toBe(true)
expect(ctx.sessionTitle.get(session)?.source.kind).toBe('fallback')
const replacement: SessionTitleProvider = {
@@ -214,7 +219,7 @@ describe('SessionTitleService provider lifecycle', () => {
generate: async () => ({ title: 'replacement', messageSeqs: [message.seq] }),
}
const disposeReplacement = ctx.sessionTitle.register(replacement)
disposeReplacement()
await disposeReplacement()
})
it('supersedes an older all-messages revision and cannot commit an ignored abort', async () => {

View File

@@ -1,4 +1,4 @@
import { Context } from 'cordis'
import { Context, type Fiber } from 'cordis'
import { describe, expect, it, vi } from 'vitest'
import SessionStore, { Session, SessionId } from '@deepseek-ai/dsh-session'
import SessionTitleService, {
@@ -162,6 +162,157 @@ describe('SessionTitleService configuration and refresh boundaries', () => {
expect(disposeSignal?.aborted).toBe(true)
})
it('reserves overlapping refresh order before fallback durability settles', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
await ctx.plugin(SessionTitleService, CONFIG)
const seed = new Session(SessionId('refresh-order-seed'))
seed.append('turn/start', {
turn: 1,
trigger: { kind: 'message', source: { kind: 'user' } },
})
const source = appendPrompt(seed, 'Keep the newest explicit refresh')
seed.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
const session = ctx.sessions.create(SessionId('refresh-order'), { seed: seed.events })
const flushStarted = deferred<undefined>()
const releaseFlush = deferred<undefined>()
let flushCount = 0
ctx.on('session/flush', async (subject) => {
if (subject !== session || ++flushCount !== 1) return
flushStarted.resolve(undefined)
await releaseFlush.promise
})
const result = deferred<SessionTitleProviderResult>()
const requests: SessionTitleProviderRequest[] = []
ctx.sessionTitle.register({
id: SessionTitleProviderId('refresh-order'),
automatic: 'first-message',
generate(request) {
requests.push(request)
return result.promise
},
})
const older = ctx.sessionTitle.refresh(session)
const olderOutcome = older.then(
() => undefined,
(error: unknown) => error,
)
await flushStarted.promise
const newer = ctx.sessionTitle.refresh(session)
await settle()
expect(requests).toHaveLength(1)
expect(requests[0]?.signal.aborted).toBe(false)
releaseFlush.resolve(undefined)
await settle()
expect(requests).toHaveLength(1)
expect(requests[0]?.signal.aborted).toBe(false)
result.resolve({ title: 'Newest explicit title', messageSeqs: [source.seq] })
await expect(newer).resolves.toMatchObject({ title: 'Newest explicit title' })
const olderError = await olderOutcome
expect(olderError).toBeInstanceOf(Error)
if (!(olderError instanceof Error)) throw new Error('expected older refresh to reject')
expect(olderError.message).toMatch(/superseded/)
})
it('cancels a queued fallback when the session-title service unloads', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const lifecycle: { fiber?: Fiber; session?: Session; inactiveRefresh?: Promise<unknown> } = {}
ctx.on('internal/plugin', (subject) => {
if (subject !== lifecycle.fiber || subject.uid !== null || lifecycle.session === undefined) return
appendPrompt(lifecycle.session, 'Ignore reentrant disposal prompt')
lifecycle.session.append('request/header', {
header: { config: { provider: 'main', model: 'main' } },
reason: 'initial',
})
lifecycle.inactiveRefresh = ctx.sessionTitle.refresh(lifecycle.session).then(
() => undefined,
(error: unknown) => error,
)
})
const fiber = await ctx.plugin(SessionTitleService, CONFIG)
lifecycle.fiber = fiber
const session = startSession(ctx, 'service-dispose-fallback')
lifecycle.session = session
appendPrompt(session, 'Do not publish after service disposal')
await fiber.dispose()
await settle()
expect(session.events.some(event => event.type === 'session/title')).toBe(false)
const inactiveError = await lifecycle.inactiveRefresh
expect(inactiveError).toBeInstanceOf(Error)
if (!(inactiveError instanceof Error)) throw new Error('expected inactive refresh to reject')
expect(inactiveError.message).toBe('session-title service disposed')
})
it('aborts pending and active provider work and drains ignored cancellation during service unload', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(SessionTitleService, CONFIG)
const result = deferred<SessionTitleProviderResult>()
const requests: SessionTitleProviderRequest[] = []
ctx.sessionTitle.register({
id: SessionTitleProviderId('service-unload'),
automatic: 'all-user-messages',
generate(request) {
requests.push(request)
return result.promise
},
})
const active = startSession(ctx, 'service-unload-active')
const activeMessage = appendPrompt(active, 'Active provider work')
await settle()
const refresh = ctx.sessionTitle.refresh(active)
const refreshOutcome = refresh.then(
() => undefined,
(error: unknown) => error,
)
await settle()
expect(requests).toHaveLength(1)
const pending = startSession(ctx, 'service-unload-pending')
appendPrompt(pending, 'Pending provider work')
const disposal = fiber.dispose()
let disposed = false
void disposal.then(() => { disposed = true })
await settle()
expect(requests[0]?.signal.aborted).toBe(true)
expect(disposed).toBe(false)
result.resolve({ title: 'Ignored service abort', messageSeqs: [activeMessage.seq] })
await disposal
expect(disposed).toBe(true)
await expect(refreshOutcome).resolves.toEqual(expect.objectContaining({ message: 'session-title service disposed' }))
})
it('suppresses a queued fallback failure after service unload begins', async () => {
const ctx = new Context()
await ctx.plugin(SessionStore)
const fiber = await ctx.plugin(SessionTitleService, CONFIG)
const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
const session = startSession(ctx, 'service-unload-flush')
appendPrompt(session, 'Fallback whose flush outlives the service')
session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
const flushStarted = deferred<undefined>()
const releaseFlush = deferred<undefined>()
ctx.on('session/flush', async (subject) => {
if (subject !== session) return
flushStarted.resolve(undefined)
await releaseFlush.promise
throw new Error('flush failed during service unload')
})
await flushStarted.promise
const disposal = fiber.dispose()
releaseFlush.resolve(undefined)
await disposal
expect(warn).not.toHaveBeenCalled()
})
it('warns when a detached session prevents queued fallback publication', async () => {
const ctx = await setup()
const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => undefined)
@@ -238,10 +389,13 @@ describe('SessionTitleService provider validation and stale scheduling', () => {
header: { config: { provider: 'main', model: 'main' } },
reason: 'initial',
})
dispose()
const pending = startSession(ctx, 'pending-provider-dispose')
appendPrompt(pending, 'Drop pending provider work')
await dispose()
await settle()
expect(generate).not.toHaveBeenCalled()
expect(ctx.sessionTitle.get(session)?.source.kind).toBe('fallback')
expect(ctx.sessionTitle.get(pending)?.source.kind).toBe('fallback')
})
it('rejects malformed provider results without replacing the fallback', async () => {